Strumienie Kinesis, Firehose i analityka w czasie rzeczywistym
Pozyskiwać dane strumieniowe za pomocą Kinesis Data Streams, dostarczać je do S3 lub Redshift przez Firehose oraz analizować w czasie rzeczywistym za pomocą Managed Service for Apache Flink.
Strumienie Kinesis, Firehose i analityka w czasie rzeczywistym to bezpłatna lekcja AWS Solutions Architect na CoddyKit. To lekcja 4 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.
Przegląd rodziny Kinesis
Amazon Kinesis to rodzina usług do gromadzenia, przetwarzania i analizowania strumieni danych w czasie rzeczywistym. Trzy podstawowe usługi to: Kinesis Data Streams (niestandardowe przetwarzanie z małymi opóźnieniami), Kinesis Data Firehose (w pełni zarządzane dostarczanie do S3/Redshift/OpenSearch) oraz Managed Service for Apache Flink (dawniej Kinesis Data Analytics) do obsługi zapytań SQL i przetwarzania za pomocą Flink w czasie rzeczywistym. Każda usługa odpowiada innemu etapowi potoku strumieniowego.
Architektura Kinesis Data Streams
Kinesis Data Stream to trwały, uporządkowany dziennik podzielony na fragmenty. Każdy fragment zapewnia przepustowość zapisu 1 MB/s i odczytu 2 MB/s. Rekordy danych są domyślnie przechowywane przez 24 godziny (okres ten można wydłużyć do 7 lub 365 dni). Producenci zapisują rekordy w strumieniu, a odbiorcy — Lambda, aplikacje KCL, Firehose lub Flink — odczytują je równolegle z co najmniej jednego fragmentu. Po zapisaniu rekordy są niezmienne.
# Create a Kinesis Data Stream with 4 shards
aws kinesis create-stream \
--stream-name clickstream \
--shard-count 4
# Put a record into the stream
aws kinesis put-record \
--stream-name clickstream \
--partition-key 'user-123' \
--data 'eyJldmVudCI6ICJjbGljayJ9'Fragmenty, przepustowość i skalowanie
Liczba fragmentów określa całkowitą przepustowość strumienia. Można podzielić fragment, aby podwoić przepustowość, albo scalić dwa fragmenty, aby obniżyć koszty. Użyj funkcji Enhanced Fan-Out, aby zapewnić każdemu zarejestrowanemu odbiorcy własną przepustowość odczytu 2 MB/s, niezależnie od innych odbiorców. Eliminuje to ograniczanie odczytu, gdy wiele aplikacji korzysta z tego samego strumienia. Monitoruj GetRecords.IteratorAgeMilliseconds, aby wykrywać opóźnienia odbiorców.
# Split shard to increase throughput
aws kinesis split-shard \
--stream-name clickstream \
--shard-to-split shardId-000000000001 \
--new-starting-hash-key 170141183460469231731687303715884105728
# Register an enhanced fan-out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/clickstream \
--consumer-name analytics-appKinesis Data Firehose: zarządzane dostarczanie
Kinesis Data Firehose to w pełni zarządzana usługa, która przechwytuje, przekształca i dostarcza dane strumieniowe do miejsc docelowych, takich jak S3, Amazon Redshift, Amazon OpenSearch Service, Splunk i punkty końcowe HTTP. Nie ma tu fragmentów, którymi trzeba zarządzać — Firehose skaluje się automatycznie. Konfiguruje się rozmiar bufora (1–128 MB) i interwał bufora (60–900 sekund); Firehose dostarcza dane, gdy jako pierwszy zostanie osiągnięty którykolwiek z tych limitów.
# Create a Firehose delivery stream to S3
aws firehose create-delivery-stream \
--delivery-stream-name clickstream-to-s3 \
--s3-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"BucketARN": "arn:aws:s3:::my-data-lake-123",
"Prefix": "landing/clickstream/year=!{timestamp:yyyy}/month=!{timestamp:MM}/",
"BufferingHints": {"SizeInMBs": 64, "IntervalInSeconds": 300},
"CompressionFormat": "GZIP"
}'Przekształcanie danych Firehose za pomocą Lambda
Firehose może wywoływać funkcję Lambda dla każdej partii rekordów przed ich dostarczeniem, aby przekształcać, wzbogacać lub filtrować dane w locie. Typowe zastosowania obejmują konwersję JSON do Parquet (za pomocą schematu Glue), maskowanie pól PII lub odrzucanie zdarzeń o małej wartości. Rekordy, których przekształcenie się nie powiodło, można opcjonalnie wysłać do osobnego prefiksu błędów w S3 w celu ponownego przetworzenia, dzięki czemu dane nie zostaną utracone.
# Lambda transform function signature for Firehose
def lambda_handler(event, context):
output = []
for record in event['records']:
import base64, json
payload = json.loads(base64.b64decode(record['data']))
# Drop events with no user_id
if not payload.get('user_id'):
output.append({'recordId': record['recordId'], 'result': 'Dropped', 'data': record['data']})
else:
output.append({'recordId': record['recordId'], 'result': 'Ok', 'data': record['data']})
return {'records': output}Managed Service for Apache Flink
Amazon Managed Service for Apache Flink (dawniej Kinesis Data Analytics) uruchamia aplikacje Apache Flink w pełni zarządzanej infrastrukturze. Usługa służy do stanowej analityki w czasie rzeczywistym: agregowania w przesuwających się oknach, wykrywania anomalii, dopasowywania wzorców w sekwencjach zdarzeń oraz łączenia danych strumieniowych z tabelami referencyjnymi. Kod Flink można pisać w językach Java, Python lub Scala, a Flink zarządza punktami kontrolnymi i stanem z gwarancją dokładnie jednokrotnego przetworzenia.
# Flink SQL-style tumbling window (conceptual)
# Count page views per URL every 5 minutes
CREATE TABLE clickstream (
url STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH ('connector' = 'kinesis', 'stream' = 'clickstream', ...);
SELECT
url,
COUNT(*) AS views,
TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start
FROM clickstream
GROUP BY url, TUMBLE(event_time, INTERVAL '5' MINUTE);Wybór między Streams a Firehose
Na egzaminie SAA-C03 należy znać kryteria wyboru. Użyj Kinesis Data Streams, gdy potrzebujesz opóźnień poniżej sekundy, jednoczesnego odczytu przez wielu odbiorców lub niestandardowej logiki przetwarzania z pełną kontrolą nad przechowywaniem i ponownym odtwarzaniem danych. Użyj Firehose, gdy wystarczy niezawodne dostarczanie danych strumieniowych do S3, Redshift lub OpenSearch przy minimalnej ilości kodu, opcjonalnym przekształcaniu w locie i automatycznym skalowaniu, ale z większym opóźnieniem (co najmniej 60 sekund).
Kinesis a SQS: klasyczny wybór egzaminacyjny
Typowe pytanie na egzaminie SAA-C03 wymaga wyboru między Kinesis a SQS. Najważniejsze różnice: Kinesis zachowuje kolejność komunikatów w obrębie fragmentu, umożliwia wielu odbiorcom jednoczesny odczyt tych samych danych i przechowuje rekordy na potrzeby ponownego odtwarzania. SQS usuwa komunikaty po ich pobraniu (bez możliwości ponownego odtworzenia), tryb FIFO gwarantuje ścisłą kolejność, a usługa lepiej nadaje się do rozdzielania mikrousług. Jeśli scenariusz wspomina o analityce w czasie rzeczywistym lub ponownym odtwarzaniu, wybierz Kinesis.
Producenci Kinesis: SDK i KPL
Kinesis Producer Library (KPL) to wysokowydajny klient do zapisywania danych w Kinesis Data Streams z poziomu aplikacji. KPL automatycznie agreguje wiele małych rekordów w jednym wywołaniu API (do 1 MB) i obsługuje ponowienia z wykładniczym wycofaniem. Znacznie ogranicza to koszt operacji PUT dla pojedynczych rekordów i zwiększa przepustowość fragmentu. KPL należy używać w przypadku producentów generujących duże ilości danych, takich jak strumienie kliknięć internetowych, czujniki IoT lub potoki dzienników.
# Basic KPL usage (Java pseudocode, no backticks)
KinesisProducer producer = new KinesisProducer();
byte[] data = 'hello world'.getBytes();
ByteBuffer buf = ByteBuffer.wrap(data);
// addUserRecord handles aggregation and retry internally
ListenableFuture future = producer.addUserRecord('clickstream', 'partitionKey', buf);
Futures.addCallback(future, new FutureCallback() { ... });Wzorzec Firehose do Redshift
Typowa architektura wykorzystuje Firehose jako zarządzany potok ze strumieni Kinesis lub bezpośrednio od producentów do Amazon Redshift w celu budowy hurtowni danych. Firehose najpierw zapisuje dane w pośrednim zasobniku przejściowym S3, a następnie wydaje polecenie COPY, aby załadować dane do Redshift. Jest to najbardziej wydajny sposób zbiorczego ładowania danych strumieniowych do Redshift — bezpośrednie wstawianie wiersz po wierszu byłoby niezwykle powolne z powodu narzutu przypadającego na każdy wiersz.
# Firehose Redshift destination (CLI snippet)
--redshift-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"ClusterJDBCURL": "jdbc:redshift://cluster.xyz.us-east-1.redshift.amazonaws.com:5439/sales",
"CopyCommand": {
"DataTableName": "clickevents",
"CopyOptions": "JSON 'auto'"
},
"Username": "firehose_user",
"Password": "{{resolve:secretsmanager:redshift-pw}}",
"S3Configuration": {
"RoleARN": "...",
"BucketARN": "arn:aws:s3:::firehose-staging"
}
}'Monitorowanie Kinesis za pomocą CloudWatch
Najważniejsze metryki CloudWatch dla Kinesis to: IncomingBytes i IncomingRecords do pomiaru przepustowości producentów, GetRecords.IteratorAgeMilliseconds do pomiaru opóźnienia odbiorców (wysoka wartość oznacza, że odbiorcy nie nadążają) oraz WriteProvisionedThroughputExceeded do wykrywania sytuacji, w których producenci osiągają limity fragmentów. Ustaw alarmy CloudWatch dla wieku iteratora i przekroczenia przydzielonej przepustowości, aby uruchamiać automatyczne skalowanie fragmentów za pośrednictwem Application Auto Scaling.
# CloudWatch alarm on high consumer lag
aws cloudwatch put-metric-alarm \
--alarm-name kinesis-high-lag \
--metric-name GetRecords.IteratorAgeMilliseconds \
--namespace AWS/Kinesis \
--dimensions Name=StreamName,Value=clickstream \
--statistic Maximum \
--period 60 \
--threshold 60000 \
--comparison-operator GreaterThanThreshold \
--evaluation-periods 3 \
--alarm-actions arn:aws:sns:us-east-1:123456789012:ops-alertsSzybkie sprawdzenie
Sprawdź swoją znajomość zagadnień AWS Solutions Architect (SAA-C03) omówionych w tej lekcji.
Podsumowanie lekcji
W tej lekcji nauczyli się Państwo, że: Kinesis Data Streams zapewnia trwałe, uporządkowane i podzielone na fragmenty strumieniowanie z obsługą wielu odbiorców oraz ponownego odtwarzania danych, Kinesis Data Firehose oferuje w pełni zarządzane dostarczanie bez kodu do S3, Redshift i OpenSearch, a Managed Service for Apache Flink umożliwia stanową analitykę strumieni w czasie rzeczywistym. W następnej części omówimy 7 reguł strategii migracji do chmury.
Często zadawane pytania
Czy lekcja „Strumienie Kinesis, Firehose i analityka w czasie rzeczywistym” jest bezpłatna?
Tak — pełny tekst „Strumienie Kinesis, Firehose i analityka 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 „Strumienie Kinesis, Firehose i analityka w czasie rzeczywistym”?
Pozyskiwać dane strumieniowe za pomocą Kinesis Data Streams, dostarczać je do S3 lub Redshift przez Firehose oraz analizować w czasie rzeczywistym za pomocą Managed Service for Apache Flink. Ć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 4 z 4.
Ile czasu zajmuje lekcja „Strumienie Kinesis, Firehose i analityka 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
- Budowanie jeziora danych w S3
- AWS Glue: ETL i katalog danych
- Amazon Athena: bezserwerowy SQL w S3
- Strumienie Kinesis, Firehose i analityka w czasie rzeczywistym