Cloud & IT Cert Prep · Lektion

Kinesis Streams, Firehose und Echtzeitanalysen

Nehmen Sie Streaming-Daten mit Kinesis Data Streams auf, liefern Sie sie mit Firehose an S3 oder Redshift und analysieren Sie sie mit Managed Service for Apache Flink in Echtzeit.

Lektion 4 von 413 Schritte

Kinesis Streams, Firehose und Echtzeitanalysen ist eine kostenlose Cloud & IT Cert Prep-Lektion auf CoddyKit. Dies ist Lektion 4 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des Cloud & IT Cert Prep-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der Cloud & IT Cert Prep-Kurs umfasst insgesamt 4 Lektionen.

Überblick über die Kinesis-Produktfamilie

Amazon Kinesis ist eine Familie von Services zum Erfassen, Verarbeiten und Analysieren von Streamingdaten in Echtzeit. Die drei zentralen Services sind: Kinesis Data Streams (niedrige Latenz und benutzerdefinierte Verarbeitung), Kinesis Data Firehose (vollständig verwaltete Übermittlung an S3/Redshift/OpenSearch) und Managed Service for Apache Flink (früher Kinesis Data Analytics) für Echtzeit-SQL und die Verarbeitung mit Flink. Jeder Service ist auf einen anderen Abschnitt der Streamingpipeline ausgerichtet.

Architektur von Kinesis Data Streams

Ein Kinesis Data Stream ist ein dauerhaftes, geordnetes Protokoll, das in Shards partitioniert ist. Jeder Shard bietet einen Schreibdurchsatz von 1 MB/s und einen Lesedurchsatz von 2 MB/s. Datenaufzeichnungen werden standardmäßig 24 Stunden aufbewahrt (erweiterbar auf 7 oder 365 Tage). Produzenten schreiben Aufzeichnungen in einen Stream; Konsumenten – Lambda, KCL-Anwendungen, Firehose oder Flink – lesen parallel aus einem oder mehreren Shards. Nach dem Schreiben sind die Aufzeichnungen unveränderlich.

# 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'

Shards, Durchsatz und Skalierung

Die Anzahl der Shards bestimmt den Gesamtdurchsatz eines Streams. Sie können einen Shard aufteilen, um den Durchsatz zu verdoppeln, oder zwei Shards zusammenführen, um Kosten zu reduzieren. Verwenden Sie Enhanced Fan-Out, um jedem registrierten Konsumenten unabhängig von anderen Konsumenten einen eigenen Lesedurchsatz von 2 MB/s bereitzustellen. Dadurch wird eine Drosselung beim Lesen verhindert, wenn mehrere Anwendungen denselben Stream konsumieren. Überwachen Sie GetRecords.IteratorAgeMilliseconds, um Verzögerungen bei Konsumenten zu erkennen.

# 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-app

Kinesis Data Firehose: Verwaltete Übermittlung

Kinesis Data Firehose ist ein vollständig verwalteter Service, der Streamingdaten erfasst, transformiert und an Ziele wie S3, Amazon Redshift, Amazon OpenSearch Service, Splunk und HTTP-Endpunkte übermittelt. Es gibt keine Shards zu verwalten – Firehose skaliert automatisch. Sie konfigurieren eine Puffergröße (1–128 MB) und ein Pufferintervall (60–900 Sekunden); Firehose übermittelt die Daten, sobald eines der beiden Limits zuerst erreicht ist.

# 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"
  }'

Datentransformation in Firehose mit Lambda

Firehose kann vor der Übermittlung für jedes Batch von Aufzeichnungen eine Lambda-Funktion aufrufen, um die Daten während der Übertragung zu transformieren, anzureichern oder zu filtern. Häufige Anwendungsfälle sind die Konvertierung von JSON in Parquet (über ein Glue-Schema), das Maskieren von PII-Feldern oder das Verwerfen von Ereignissen mit geringem Wert. Aufzeichnungen, bei denen die Transformation fehlschlägt, können optional an ein separates S3-Fehlerpräfix zur erneuten Verarbeitung gesendet werden, sodass keine Daten verloren gehen.

# 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 (früher Kinesis Data Analytics) führt Apache-Flink-Anwendungen auf einer vollständig verwalteten Infrastruktur aus. Verwenden Sie den Service für zustandsbehaftete Echtzeitanalysen: Aggregationen über gleitende Fenster, Anomalieerkennung, Musterabgleich in Ereignisfolgen und Joins von Streamingdaten mit Referenztabellen. Sie schreiben Flink-Code in Java, Python oder Scala, während Flink Checkpoints und einen Exactly-once-Zustand verwaltet.

# 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);

Zwischen Streams und Firehose wählen

Für die SAA-C03-Prüfung müssen Sie die Entscheidungskriterien kennen. Verwenden Sie Kinesis Data Streams, wenn Sie eine Latenz im Subsekundenbereich, mehrere gleichzeitig lesende Konsumenten oder benutzerdefinierte Verarbeitungslogik mit vollständiger Kontrolle über Aufbewahrung und Replay benötigen. Verwenden Sie Firehose, wenn Sie Streamingdaten lediglich zuverlässig mit wenig Code, optionaler Transformation während der Übertragung und automatischer Skalierung an S3, Redshift oder OpenSearch übermitteln möchten und eine höhere Latenz (60+ Sekunden) akzeptabel ist.

Kinesis vs. SQS: Die klassische Prüfungsentscheidung

Eine häufige SAA-C03-Frage verlangt die Entscheidung zwischen Kinesis und SQS. Die wichtigsten Unterschiede: Kinesis bewahrt die Nachrichtenreihenfolge innerhalb eines Shards, unterstützt mehrere Konsumenten, die dieselben Daten gleichzeitig lesen, und bewahrt Aufzeichnungen für ein Replay auf. SQS entfernt Nachrichten nach dem Konsum (kein Replay), der FIFO-Modus garantiert eine strikt eingehaltene Reihenfolge, und SQS eignet sich besser zur Entkopplung von Microservices. Wenn das Szenario Echtzeitanalysen oder Replay erwähnt, wählen Sie Kinesis.

Kinesis-Produzenten: SDK und KPL

Die Kinesis Producer Library (KPL) ist ein Client mit hohem Durchsatz zum Schreiben in Kinesis Data Streams aus Anwendungen heraus. KPL aggregiert automatisch mehrere kleine Aufzeichnungen in einem einzigen API-Aufruf (bis zu 1 MB) und verarbeitet Wiederholungsversuche mit Backoff. Dadurch werden die Kosten für PUTs pro Aufzeichnung deutlich reduziert und der Durchsatz pro Shard erhöht. Verwenden Sie KPL für Produzenten mit hohem Datenvolumen, etwa bei Web-Clickstreams, IoT-Sensoren oder Log-Pipelines.

# 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() { ... });

Firehose-zu-Redshift-Muster

Eine häufige Architektur verwendet Firehose als verwaltete Pipeline von Kinesis-Streams oder direkten Produzenten zu Amazon Redshift für Data Warehousing. Firehose schreibt die Daten zunächst in einen temporären S3-Staging-Bucket und führt anschließend einen COPY-Befehl aus, um die Daten in Redshift zu laden. Dies ist die effizienteste Methode, Streamingdaten per Bulk-Load in Redshift zu laden – direkte zeilenweise Einfügungen in Redshift wären aufgrund des Overheads pro Zeile äußerst langsam.

# 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"
  }
}'

Kinesis mit CloudWatch überwachen

Wichtige CloudWatch-Metriken für Kinesis sind: IncomingBytes und IncomingRecords zur Messung des Produzentendurchsatzes, GetRecords.IteratorAgeMilliseconds zur Messung der Verzögerung bei Konsumenten (ein hoher Wert bedeutet, dass die Konsumenten nicht mithalten können) und WriteProvisionedThroughputExceeded, um zu erkennen, wann Produzenten an die Shard-Limits stoßen. Richten Sie CloudWatch-Alarme für das Iteratoralter und überschrittenen bereitgestellten Durchsatz ein, um über Application Auto Scaling die automatische Skalierung der Shards auszulösen.

# 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-alerts

Kurzer Wissenstest

Testen Sie Ihr Verständnis der Konzepte von AWS Solutions Architect (SAA-C03) aus dieser Lektion.

Zusammenfassung der Lektion

In dieser Lektion haben Sie gelernt: Kinesis Data Streams bietet dauerhaftes, geordnetes Streaming mit Shards, mehreren Konsumenten und Replay-Funktion, Kinesis Data Firehose ermöglicht die vollständig verwaltete, codefreie Übermittlung an S3, Redshift und OpenSearch und Managed Service for Apache Flink ermöglicht zustandsbehaftete Echtzeitanalysen auf Streams. Als Nächstes sehen wir uns die 7 Rs der Cloud-Migrationsstrategie an.

Kostenlos starten

Lerne Cloud & IT Cert Prep mit einem KI-Tutor — kostenlos

Schreibe und führe echten Code in deinem Browser aus, bekomme sofortige Hilfe von einem 24/7 KI-Tutor und setze dein Lernen im Web oder in der App fort.

Kurse
150
Lektionen
600

Häufig gestellte Fragen

Ist die Lektion „Kinesis Streams, Firehose und Echtzeitanalysen“ kostenlos?

Ja — der vollständige Text von „Kinesis Streams, Firehose und Echtzeitanalysen“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des Cloud & IT Cert Prep-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der Cloud & IT Cert Prep-Kurs umfasst insgesamt 4 Lektionen.

Was lerne ich in „Kinesis Streams, Firehose und Echtzeitanalysen“?

Nehmen Sie Streaming-Daten mit Kinesis Data Streams auf, liefern Sie sie mit Firehose an S3 oder Redshift und analysieren Sie sie mit Managed Service for Apache Flink in Echtzeit. Du übst Cloud & IT Cert Prep mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.

Brauche ich Erfahrung, um Cloud & IT Cert Prep zu starten?

Keine Vorkenntnisse erforderlich. Cloud & IT Cert Prep auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 4 von 4.

Wie lange dauert die Lektion „Kinesis Streams, Firehose und Echtzeitanalysen“?

Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.

Kann ich in dieser Cloud & IT Cert Prep-Lektion Code schreiben und ausführen?

Ja. Jede Cloud & IT Cert Prep-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.

Alle Lektionen in diesem Kurs

  1. Einen Data Lake auf S3 erstellen
  2. AWS Glue: ETL und Datenkatalog
  3. Amazon Athena: Server-SQL auf S3
  4. Kinesis Streams, Firehose und Echtzeitanalysen
← Zurück zu Cloud & IT Cert Prep