Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym
Tworzyć i odbierać wysokoprzepustowe strumienie zdarzeń za pomocą Kinesis Data Streams, zarządzać fragmentami w celu zapewnienia przepustowości oraz używać Lambda jako konsumenta.
Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym to bezpłatna lekcja AWS Solutions Architect na CoddyKit. To lekcja 3 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AWS Solutions Architect, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AWS Solutions Architect zawiera 4 lekcji w sumie.
Podstawowe pojęcia Kinesis Data Streams
Kinesis Data Streams (KDS) to trwała, uporządkowana usługa strumieniowania danych w czasie rzeczywistym. Dane są organizowane w ramach strumienia składającego się z co najmniej jednego fragmentu (shard). Każdy fragment jest uporządkowaną sekwencją rekordów danych. Producenci umieszczają rekordy we fragmentach za pomocą klucza partycji, który określa, do którego fragmentu trafi rekord. Konsumenci odczytują rekordy z fragmentów i przetwarzają je w kolejności, w jakiej pojawiły się w danym fragmencie.
Limity pojemności i przepustowości fragmentu
Każdy fragment obsługuje przepustowość zapisu 1 MB/s lub 1000 rekordów/s oraz przepustowość odczytu 2 MB/s (współdzieloną przez wszystkich standardowych konsumentów tego fragmentu). Całkowita pojemność strumienia skaluje się liniowo wraz z liczbą fragmentów. Proszę użyć wzoru: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Jeśli producenci osiągną limit zapisu, pojawią się błędy ProvisionedThroughputExceededException — należy je rozwiązać, dzieląc fragmenty lub równomierniej rozkładając klucze partycji.
# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed
aws kinesis create-stream \
--stream-name iot-telemetry \
--shard-count 5Klucze partycji i dystrybucja danych
Klucz partycji to ciąg znaków, który Kinesis haszuje (MD5), aby określić, do którego fragmentu trafi rekord. Dobrze dobrany klucz partycji równomiernie rozdziela rekordy między fragmenty (zapobieganie przeciążeniu fragmentów). W przypadku IoT należy użyć identyfikatora urządzenia. W przypadku strumieni kliknięć należy użyć identyfikatora sesji lub użytkownika. Należy unikać kluczy o małej liczbie możliwych wartości (np. nazwy kraju, gdy występuje tylko 5 wartości), ponieważ powodują one przeciążenie fragmentów — jeden fragment otrzymuje nieproporcjonalnie duży ruch zapisu, podczas gdy pozostałe pozostają bezczynne.
import boto3, json
client = boto3.client('kinesis', region_name='us-east-1')
# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
StreamName='iot-telemetry',
Data=json.dumps(event),
PartitionKey='sensor-42' # high-cardinality -> even distribution
)Standardowi konsumenci a Enhanced Fan-Out
Standardowi konsumenci współdzielą przepustowość odczytu 2 MB/s na fragment, korzystając z GetRecords i odpytywania. Jeśli mają Państwo 3 konsumentów fragmentu, z których każdy potrzebuje 2 MB/s, będą oni podlegać ograniczeniu przepustowości, współdzieląc łącznie 2 MB/s. Enhanced Fan-Out (EFO) zapewnia każdemu zarejestrowanemu konsumentowi własny dedykowany kanał odczytu 2 MB/s za pośrednictwem trwałego połączenia push HTTP/2 (SubscribeToShard). EFO wiąże się z kosztem za godzinę konsumenta i fragmentu, ale całkowicie eliminuje rywalizację o przepustowość odczytu.
# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--consumer-name real-time-analytics
# List registered consumers
aws kinesis list-stream-consumers \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetryLambda jako konsument Kinesis
Lambda natywnie integruje się z Kinesis Data Streams za pośrednictwem Event Source Mapping. Lambda odpytuje strumień, odczytuje partie rekordów i wywołuje Państwa funkcję. Należy skonfigurować BatchSize (od 1 do 10 000 rekordów), StartingPosition (TRIM_HORIZON dla najstarszych rekordów, LATEST dla najnowszych) oraz BisectBatchOnFunctionError, aby dzielić partie, których przetwarzanie zakończyło się niepowodzeniem. Współczynnik równoległości (Parallelisation Factor) (od 1 do 10) pozwala usłudze Lambda uruchamiać wiele współbieżnych wywołań na fragment, aby nadążyć za szybko napływającymi danymi.
# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
--function-name ProcessIoTEvents \
--event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--batch-size 100 \
--starting-position LATEST \
--parallelization-factor 5 \
--bisect-batch-on-function-error true \
--destination-config '{
"OnFailure": {
"Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
}
}'Kinesis Client Library (KCL)
Kinesis Client Library (KCL) to framework aplikacyjny do tworzenia niezawodnych konsumentów Kinesis w języku Java (lub w wielu językach za pośrednictwem MultiLangDaemon). KCL obsługuje wyliczanie fragmentów, zarządzanie dzierżawami (rozdzielanie fragmentów między instancje procesów roboczych), zapisywanie punktów kontrolnych postępu w DynamoDB oraz bezpieczną obsługę podziałów i scalania fragmentów. Każdy proces roboczy KCL przetwarza co najmniej jeden fragment, a KCL automatycznie ponownie równoważy fragmenty w miarę skalowania procesów roboczych lub ich awarii. W przypadku produkcyjnych aplikacji konsumenckich KCL jest preferowane zamiast bezpośredniego odpytywania za pomocą SDK.
# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
# 'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();Przechowywanie i ponowne odtwarzanie danych
Kinesis Data Streams przechowuje rekordy przez domyślnie 24 godziny (okres ten można wydłużyć do 7 dni lub maksymalnie 365 dni dzięki Long-Term Retention, co wiąże się z dodatkowym kosztem). W przeciwieństwie do SQS skonsumowane rekordy nie są usuwane po odczycie — pozostają dostępne do upływu okresu przechowywania. Umożliwia to wielu konsumentom niezależne odczytywanie tych samych rekordów, a także ponowne odtwarzanie przez ustawienie punktu kontrolnego konsumenta na wcześniejszy numer sekwencyjny — co jest niezwykle przydatne przy usuwaniu błędów lub uzupełnianiu danych dla nowych usług.
# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
--stream-name iot-telemetry \
--retention-period-hours 168
# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
--stream-name iot-telemetry \
--shard-id shardId-000000000000 \
--shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' --output text)
aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100Tryb przepustowości On-Demand
Kinesis Data Streams obsługuje dwa tryby przepustowości. Tryb Provisioned: ręcznie zarządzają Państwo liczbą fragmentów i płacą za każdą godzinę fragmentu. Tryb On-Demand: Kinesis automatycznie skaluje pojemność fragmentów na podstawie przychodzącej przepustowości (domyślnie do 200 MB/s zapisu i 400 MB/s odczytu), a opłaty są naliczane za każdy GB zapisanych i pobranych danych. Tryb On-Demand jest idealny w przypadku zmiennego lub nieprzewidywalnego ruchu, gdy nie chcą Państwo ręcznie zarządzać skalowaniem fragmentów.
# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--stream-mode-details StreamMode=ON_DEMANDGwarancje kolejności w obrębie fragmentów
Kinesis gwarantuje kolejność w obrębie fragmentu — rekordy z tym samym kluczem partycji zawsze trafiają do tego samego fragmentu i są odczytywane w kolejności zapisu. Nie ma jednak gwarancji kolejności między fragmentami. Jeśli aplikacja wymaga globalnego uporządkowania wszystkich rekordów, należy użyć jednego fragmentu (ograniczając przepustowość do 1 MB/s) albo przeprojektować ją tak, aby kolejność była wymagana tylko w obrębie grupy kluczy partycji (np. kolejność danych dla poszczególnych urządzeń). Jest to częsta różnica egzaminacyjna w porównaniu z SQS FIFO, które zapewnia ścisłą deduplikację i uporządkowanie.
Bezpieczeństwo: szyfrowanie i VPC
Kinesis Data Streams szyfruje rekordy po stronie serwera podczas przechowywania za pomocą AWS KMS (klucza CMK lub klucza zarządzanego przez AWS), gdy szyfrowanie po stronie serwera jest włączone. Wszystkie dane przesyłane podczas transmisji są szyfrowane za pomocą TLS. W przypadku aplikacji działających w VPC, które nie powinny przesyłać danych strumienia przez publiczny internet, należy użyć VPC Interface Endpoint (PrivateLink) dla Kinesis, aby ruch pozostawał w całości w szkielecie sieci AWS — jest to istotne w przypadku obciążeń podlegających regulacjom.
# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
--stream-name iot-telemetry \
--encryption-type KMS \
--key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123Kinesis Data Streams a SQS i Kafka
Na egzaminie SAA-C03 należy porównywać Kinesis Data Streams z alternatywnymi rozwiązaniami. Kinesis a SQS: Kinesis zachowuje kolejność w obrębie fragmentu i umożliwia wielu konsumentom odczytywanie tych samych danych; SQS usuwa wiadomości po ich skonsumowaniu. Kinesis a MSK (Kafka): MSK to zarządzany Apache Kafka — należy go użyć, gdy potrzebna jest zgodność z protokołem Kafka, zaawansowana konfiguracja tematów lub migracja z lokalnego środowiska Kafka. Kinesis należy wybrać do strumieniowania natywnego dla AWS, zapewniającego ściślejszą integrację z Lambda, Firehose i Flink. Kinesis należy wybrać, chyba że pytanie wyraźnie wspomina o Kafka lub wymaganiach dotyczących zgodności z Kafka.
Szybki sprawdzian
Sprawdź swoją znajomość zagadnień AWS Solutions Architect (SAA-C03) z tej lekcji.
Podsumowanie lekcji
W tej lekcji poznali Państwo: Kinesis Data Streams zapewnia podzielone na fragmenty, uporządkowane i trwałe strumienie z domyślnym przechowywaniem przez 24 godziny oraz możliwością ponownego odtwarzania, Enhanced Fan-Out zapewnia każdemu konsumentowi dedykowaną przepustowość 2 MB/s na fragment, eliminując rywalizację o odczyt, a także tryb On-Demand automatycznie skaluje fragmenty w przypadku nieprzewidywalnego ruchu. W następnej części omówimy wzorce choreografii i orkiestracji w architekturze sterowanej zdarzeniami.
Często zadawane pytania
Czy lekcja „Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym” jest bezpłatna?
Tak — pełny tekst „Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AWS Solutions Architect, przejdź na CoddyKit PRO. Kurs AWS Solutions Architect zawiera 4 lekcji w sumie.
Co nauczysz się w „Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym”?
Tworzyć i odbierać wysokoprzepustowe strumienie zdarzeń za pomocą Kinesis Data Streams, zarządzać fragmentami w celu zapewnienia przepustowości oraz używać Lambda jako konsumenta. Ćwiczysz AWS Solutions Architect z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć AWS Solutions Architect?
Nie wymagamy żadnego doświadczenia. AWS Solutions Architect w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 3 z 4.
Ile czasu zajmuje lekcja „Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji AWS Solutions Architect?
Tak. Każda lekcja AWS Solutions Architect zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- EventBridge: magistrala zdarzeń i reguły
- Step Functions: orkiestracja bezserwerowych przepływów pracy
- Kinesis Data Streams do przetwarzania zdarzeń w czasie rzeczywistym
- Wzorce choreografii a orkiestracji