Strumieniowanie danych

Tworzenie aplikacji na AWS

Ricardo Sueiras

Principal Technologist

Podstawy strumieniowania danych

 

streaming-data

  • Rekordy są przetwarzane na bieżąco, w miarę napływania.
  • Producenci zapisują dane do strumienia.
  • Konsumenci odczytują dane niezależnie i równolegle.
  • System absorbuje skoki ruchu bez oczekiwania na zadania wsadowe.
Tworzenie aplikacji na AWS

Z czego składa się strumieniowanie danych?

 

komponenty systemu strumieniowego

  • Producenci: aplikacje generujące rekordy.
  • Strumienie: trwałe, uporządkowane bufory oddzielające producentów od konsumentów.
  • Konsumenci: aplikacje odczytujące i przetwarzające rekordy.
  • Przechowywanie i retencja: określa, jak długo rekordy są dostępne do ponownego odtworzenia.
Tworzenie aplikacji na AWS

Fragmenty (shards)

 

fragmenty

 

  • Fragmenty (shards) to jednostki przepustowości i przetwarzania równoległego.
  • Klucze partycji rozkładają rekordy między fragmenty.
  • Kolejność jest gwarantowana tylko w obrębie jednego fragmentu.
  • Konsumenci używają punktów kontrolnych (check-pointing) do śledzenia postępu i odtwarzania po awarii.
Tworzenie aplikacji na AWS

Zarządzane usługi AWS do strumieniowania danych

  • Amazon Kinesis Data Streams.
  • Amazon Data Firehose.
  • Podstawowe elementy składowe aplikacji do strumieniowania danych.

 

usługi AWS

Tworzenie aplikacji na AWS

Kinesis Data Streams

  • Strumieniowanie z małymi opóźnieniami.
  • Konfigurowalna retencja: od 24 godzin (domyślnie) do 365 dni.
  • Trwałe, możliwe do ponownego odtworzenia strumienie.
  • Wielu konsumentów może odczytywać dane równolegle.
  • Programiści zarządzają skalowaniem za pomocą fragmentów.

 

kinesis data streams

Tworzenie aplikacji na AWS

Kiedy używać Kinesis Data Streams?

  • Wybierz tę usługę, gdy potrzebujesz niskich opóźnień w czasie rzeczywistym.
  • Wymagasz możliwości ponownego odtworzenia rekordów.
  • Masz wielu niezależnych konsumentów.
  • Potrzebujesz własnej logiki transformacji.

 

kinesis data streams

Tworzenie aplikacji na AWS

Zapisywanie rekordów

 

biblioteki klienckie

  • Klucze partycji kontrolują przypisanie do fragmentu i kolejność rekordów.
  • PutRecords grupuje do 500 rekordów lub 5 MB na żądanie.
  • Błędy są raportowane osobno dla każdego rekordu.
  • FailedRecordCount pomaga zidentyfikować rekordy, które się nie powiodły.
Tworzenie aplikacji na AWS

Odczyt rekordów: tryb klasyczny

 

klasyczny

 

  • Konsumenci odpytują Kinesis API na żądanie.
  • Limity zapisu: 1 MB/s lub 1 000 rekordów/s.
  • Limity odczytu: 2 MB/s lub 5 wywołań GetRecords/s.
Tworzenie aplikacji na AWS

Odczyt rekordów: punkt startowy strumienia

 

początek strumienia

  • Konsumenci wskazują punkt startowy odczytu w strumieniu.
  • TRIM_HORIZON: najstarszy dostępny rekord w fragmencie.
  • LATEST: kolejny nowy rekord.
  • AT_SEQUENCE_NUMBER lub AFTER_SEQUENCE_NUMBER: od określonego numeru sekwencyjnego.
  • AT_TIMESTAMP: od wskazanego znacznika czasu lub po nim.
Tworzenie aplikacji na AWS

Odczyt rekordów: rozszerzony fan-out

  • Każdy zarejestrowany konsument otrzymuje dedykowane 2 MB/s na fragment.
  • Rekordy są dostarczane przez HTTP/2 push.
  • Używaj, gdy wielu niezależnych konsumentów potrzebuje niskich opóźnień.
  • A także wysokiej przepustowości z tego samego strumienia.

 

wzorzec — rozszerzony fan-out

Tworzenie aplikacji na AWS

Mapowanie źródeł zdarzeń Lambda

  • Lambda odpytuje Kinesis i wywołuje funkcje w partiach.
  • BatchSize i MaximumBatchingWindow kontrolują kształt partii.
  • ParallelizationFactor: równoczesne wywołania na fragment.
  • BisectBatchOnFunctionError: dzieli niepowodzenie na mniejsze partie, aby wyizolować problematyczne rekordy.
  • ReportBatchItemFailures: pomija ponowne przetwarzanie rekordów, które już się powiodły.
  • OnFailure: kieruje nieprzetworzone partie do kolejki SQS lub tematu SNS (DLQ).

 

mapowanie źródeł zdarzeń Lambda

Tworzenie aplikacji na AWS

Data Firehose

  • Zarządzana usługa dostarczająca dane strumieniowe do miejsc docelowych.
  • Obsługuje Amazon S3, Redshift i punkty końcowe HTTP.
  • Wspiera buforowanie, konwersję formatów i transformacje Lambda.
  • Wyższe opóźnienia — dane buforowane od sekund do minut.
  • Brak możliwości ponownego odtworzenia ani długoterminowej retencji.
  • Idealne rozwiązanie, gdy celem jest po prostu dostarczenie danych.

 

data firehose

Tworzenie aplikacji na AWS

Wzorzec: gorący-zimny

 

wzorzec — gorący/zimny

 

  • Wzorzec gorący-zimny łączy Data Streams z Data Firehose.
  • Data Streams umożliwia przetwarzanie w czasie rzeczywistym dla producentów.
  • Firehose archiwizuje rekordy z tego samego źródła do S3.
Tworzenie aplikacji na AWS

Kinesis vs SQS

 

SQS vs Kinesis

 

  • Kinesis sprawdza się przy uporządkowanych, możliwych do ponownego odtworzenia strumieniach z wieloma konsumentami.
  • Kinesis jest odpowiedni dla danych szeregów czasowych i uporządkowanych.
  • SQS stosuj do dystrybucji zadań.
Tworzenie aplikacji na AWS

Skalowanie strumieniowania danych

 

skalowanie

  • Strumienie danych są podzielone na fragmenty.
  • Klucze partycji kierują rekordy do określonego fragmentu.
  • Ten sam klucz zawsze trafia do tego samego fragmentu, zachowując kolejność.
  • Nierównomierne klucze partycji prowadzą do powstawania gorących fragmentów.
  • Stosuj klucze o dużej kardynalności, aby temu zapobiec.
  • Ponowne fragmentowanie (resharding) służy do skalowania w górę lub w dół.
Tworzenie aplikacji na AWS

Tryby pojemności: aprowizowany

  • Dostępne są dwa tryby pojemności.
  • Tryb aprowizowany używa stałej liczby fragmentów zarządzanej ręcznie.
  • Zapewnia przewidywalny koszt i pojemność.
  • Wymaga monitorowania, aby unikać ograniczania przepustowości i nadmiernej aprowizacji.

 

tryby pojemności

Tworzenie aplikacji na AWS

Tryby pojemności: na żądanie

  • Tryb na żądanie skaluje się automatycznie w zależności od ruchu.
  • Bez ręcznego zarządzania fragmentami.
  • Upraszcza obsługę nieprzewidywalnych obciążeń.
  • Wyższy koszt przy utrzymującym się dużym obciążeniu.

 

tryby pojemności

Tworzenie aplikacji na AWS

Kinesis Producer Library (KPL)

  • Biblioteki AWS upraszczają budowanie skalowalnych producentów i konsumentów.
  • Kinesis Producer Library (KPL) grupuje i agreguje rekordy.
  • Kompresuje rekordy i ponawia próby w razie błędu.
  • Buforuje asynchronicznie, aby zmaksymalizować przepustowość producentów.

 

biblioteki klienckie

Tworzenie aplikacji na AWS

Kinesis Consumer Library (KCL)

 

biblioteki klienckie

 

  • Kinesis Consumer Library (KCL) zarządza koordynacją fragmentów.
  • Obsługuje dzierżawy wątków roboczych i punkty kontrolne.
  • Stan jest przechowywany w tabeli DynamoDB.
  • Konsumenci korzystają z tego stanu, aby wiedzieć, gdzie wznowić odczyt.
Tworzenie aplikacji na AWS

Obsługa błędów: idempotentność

 

obsługa błędów

 

  • Kinesis gwarantuje dostarczenie co najmniej raz, więc duplikaty są możliwe.
  • Projektuj konsumentów tak, aby byli idempotentni.
  • Deduplikuj za pomocą numeru sekwencyjnego lub klucza biznesowego w DynamoDB.
  • Konsumenci mogą ponownie odczytać zachowane rekordy od numeru sekwencyjnego lub znacznika czasu.
Tworzenie aplikacji na AWS

Obsługa błędów: przekroczenie przepustowości

 

obsługa błędów

 

  • ProvisionedThroughputExceededException sygnalizuje ograniczenie przepustowości.
  • Rozwiąż to, zwiększając liczbę fragmentów.
  • Popraw rozkład kluczy partycji.
  • Ponów próbę z opóźnieniem (backoff).
  • Zmień tryb pojemności.
Tworzenie aplikacji na AWS
  • Śledź IncomingBytes i IncomingRecords w CloudWatch.
  • Monitoruj WriteProvisionedThroughputExceeded pod kątem ograniczania przepustowości.
  • Opóźnienie konsumenta pojawia się, gdy odczyt jest wolniejszy niż zapis.
  • Śledź GetRecords.IteratorAgeMilliseconds: wysoka wartość wskazuje na opóźnienie.

 

monitorowanie

Tworzenie aplikacji na AWS

Bezpieczeństwo

  • Szyfrowanie po stronie serwera w spoczynku przy użyciu AWS KMS.
  • Szyfrowanie w tranzycie przy użyciu TLS.
  • Dostęp kontrolowany przez polityki IAM.

 

bezpieczeństwo

Tworzenie aplikacji na AWS

Czas na praktykę!

Tworzenie aplikacji na AWS

Preparing Video For Download...