0Pricing
AWS Solutions Architect · Leçon

Kinesis Data Streams pour le traitement d’événements en temps réel

Produisez et consommez des flux d’événements à haut débit avec Kinesis Data Streams, gérez les fragments pour contrôler le débit et utilisez Lambda comme consommateur.

Kinesis Data Streams pour le traitement d’événements en temps réel est une leçon AWS Solutions Architect gratuite sur CoddyKit. Ceci est la leçon 3 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage AWS Solutions Architect, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours AWS Solutions Architect comprend 4 leçons au total.

Concepts fondamentaux de Kinesis Data Streams

Kinesis Data Streams (KDS) est un service durable et ordonné de diffusion de données en temps réel. Les données sont organisées dans un flux composé d'une ou plusieurs partitions. Chaque partition est une séquence ordonnée d'enregistrements de données. Les producteurs placent les enregistrements dans les partitions à l'aide d'une clé de partition qui détermine quelle partition reçoit l'enregistrement. Les consommateurs lisent les enregistrements des partitions et les traitent dans l'ordre de leur arrivée au sein de chaque partition.

Capacité des partitions et limites de débit

Chaque partition prend en charge un débit d'écriture de 1 MB/s ou 1 000 enregistrements/s et un débit de lecture de 2 MB/s (partagé entre tous les consommateurs standard de cette partition). La capacité totale du flux évolue linéairement avec le nombre de partitions. Utilisez la formule suivante : shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Si les producteurs atteignent la limite d'écriture, des erreurs ProvisionedThroughputExceededException apparaissent : résolvez le problème en divisant les partitions ou en répartissant les clés de partition plus uniformément.

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

Clés de partition et distribution des données

La clé de partition est une chaîne que Kinesis hache (MD5) pour déterminer quelle partition reçoit un enregistrement. Une clé de partition bien choisie répartit uniformément les enregistrements entre les partitions (prévention des partitions saturées). Pour l'IoT, utilisez l'ID de l'appareil. Pour les flux de clics, utilisez l'ID de session ou l'ID utilisateur. Évitez les clés à faible cardinalité (par exemple, le nom du pays avec seulement 5 valeurs), car elles provoquent des partitions saturées : une partition reçoit une part disproportionnée du trafic d'écriture tandis que les autres restent inactives.

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
)

Consommateurs standard et Enhanced Fan-Out

Les consommateurs standard se partagent le débit de lecture de 2 MB/s par partition à l'aide de GetRecords et d'interrogations périodiques. Si vous avez 3 consommateurs sur une partition et que chacun nécessite 2 MB/s, ils seront limités et devront se partager un total de 2 MB/s. Enhanced Fan-Out (EFO) fournit à chaque consommateur enregistré son propre canal de lecture dédié de 2 MB/s au moyen d'une connexion persistante HTTP/2 avec transmission en continu (SubscribeToShard). EFO entraîne un coût par consommateur et par heure de partition, mais élimine entièrement la concurrence en lecture.

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

Lambda comme consommateur Kinesis

Lambda s'intègre nativement à Kinesis Data Streams grâce à un mappage de source d'événements. Lambda interroge le flux, lit des lots d'enregistrements et invoque votre fonction. Configurez BatchSize (1 à 10 000 enregistrements), StartingPosition (TRIM_HORIZON pour les plus anciens, LATEST pour les plus récents) et BisectBatchOnFunctionError pour diviser les lots ayant échoué. Le facteur de parallélisation (1 à 10) permet à Lambda de lancer plusieurs invocations simultanées par partition afin de suivre les flux rapides.

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

La Kinesis Client Library (KCL) est un framework applicatif permettant de créer des consommateurs Kinesis Java robustes (ou multilingues via MultiLangDaemon). KCL gère l'énumération des partitions, la gestion des baux (répartition des partitions entre les instances Worker), l'enregistrement des points de contrôle de la progression dans DynamoDB, ainsi que les divisions et fusions progressives des partitions. Chaque Worker KCL traite une ou plusieurs partitions, et KCL rééquilibre automatiquement les partitions lorsque les Workers sont ajoutés ou tombent en panne. KCL est préférable aux interrogations brutes du SDK pour les applications de consommation en production.

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

Conservation et rejeu des données

Kinesis Data Streams stocke les enregistrements pendant 24 heures par défaut (durée extensible à 7 jours ou jusqu'à 365 jours avec la conservation à long terme, moyennant un coût supplémentaire). Contrairement à SQS, les enregistrements consommés ne sont pas supprimés après leur consommation : ils restent disponibles jusqu'à l'expiration de la durée de conservation. Cela permet à plusieurs consommateurs de lire indépendamment les mêmes enregistrements et autorise le rejeu en réinitialisant le point de contrôle d'un consommateur sur un numéro de séquence antérieur, ce qui est particulièrement utile pour corriger des bogues ou alimenter de nouveaux services avec des données historiques.

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

Mode de capacité ON_DEMAND

Kinesis Data Streams prend en charge deux modes de capacité. Mode Provisioned : vous gérez manuellement le nombre de partitions et payez par partition-heure. Mode ON_DEMAND : Kinesis ajuste automatiquement la capacité des partitions en fonction du débit entrant (jusqu'à 200 MB/s en écriture et 400 MB/s en lecture par défaut), et vous payez par GB de données écrites et récupérées. Le mode à la demande est idéal pour les modèles de trafic variables ou imprévisibles lorsque vous ne souhaitez pas gérer la mise à l'échelle des partitions.

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

Garanties d'ordre au sein des partitions

Kinesis garantit l'ordre au sein d'une partition : les enregistrements ayant la même clé de partition vont toujours dans la même partition et sont lus dans l'ordre où ils ont été écrits. En revanche, aucun ordre n'est garanti entre les partitions. Si votre application nécessite un ordre global pour tous les enregistrements, utilisez une seule partition (ce qui limite le débit à 1 MB/s) ou repensez votre conception afin que l'ordre ne soit nécessaire qu'au sein d'un groupe de clés de partition (par exemple, un ordre par appareil). Il s'agit d'une distinction fréquente à l'examen par rapport à SQS FIFO, qui fournit une déduplication et un ordre stricts.

Sécurité : chiffrement et VPC

Kinesis Data Streams chiffre les enregistrements côté serveur au repos à l'aide de AWS KMS (CMK ou clé gérée par AWS) lorsque le chiffrement côté serveur est activé. Toutes les données en transit sont chiffrées avec TLS. Pour les applications exécutées dans un VPC qui ne doivent pas envoyer les données du flux sur Internet public, utilisez un point de terminaison d'interface VPC (PrivateLink) pour Kinesis afin que le trafic reste entièrement au sein du réseau fédérateur AWS, ce qui est important pour les charges de travail réglementées.

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

Kinesis Data Streams, SQS et Kafka

Pour l'examen SAA-C03, comparez Kinesis Data Streams aux solutions alternatives. Kinesis par rapport à SQS : Kinesis préserve l'ordre au sein d'une partition et prend en charge plusieurs consommateurs lisant les mêmes données ; SQS supprime les messages après leur consommation. Kinesis par rapport à MSK (Kafka) : MSK est Apache Kafka géré : utilisez-le lorsque vous avez besoin de la compatibilité avec le protocole Kafka, de configurations avancées de rubriques ou que vous migrez depuis Kafka sur site. Utilisez Kinesis pour une diffusion propre à AWS, avec une intégration plus étroite à Lambda, Firehose et Flink. Choisissez Kinesis, sauf si la question mentionne spécifiquement Kafka ou des exigences de compatibilité avec Kafka.

Vérification rapide

Testez votre compréhension des concepts AWS Solutions Architect (SAA-C03) abordés dans cette leçon.

Récapitulatif de la leçon

Dans cette leçon, vous avez appris que Kinesis Data Streams fournit des flux partitionnés, ordonnés et durables, avec une conservation par défaut de 24 heures et une capacité de rejeu, que Enhanced Fan-Out fournit à chaque consommateur un débit dédié de 2 MB/s par partition, éliminant ainsi la concurrence en lecture, et que le mode ON_DEMAND met automatiquement les partitions à l'échelle pour les trafics imprévisibles. Nous allons maintenant découvrir les modèles de chorégraphie et d'orchestration dans une architecture orientée événements.

Questions Fréquemment Posées

La leçon « Kinesis Data Streams pour le traitement d’événements en temps réel » est-elle gratuite ?

Oui — le texte complet de « Kinesis Data Streams pour le traitement d’événements en temps réel » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours AWS Solutions Architect, passe à CoddyKit PRO. Le cours AWS Solutions Architect comprend 4 leçons au total.

Qu'est-ce que j'apprendrai dans « Kinesis Data Streams pour le traitement d’événements en temps réel » ?

Produisez et consommez des flux d’événements à haut débit avec Kinesis Data Streams, gérez les fragments pour contrôler le débit et utilisez Lambda comme consommateur. Tu pratiques AWS Solutions Architect avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.

Dois-je avoir de l'expérience pour commencer AWS Solutions Architect ?

Aucune expérience préalable n'est requise. AWS Solutions Architect sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 3 sur 4.

Combien de temps prend la leçon « Kinesis Data Streams pour le traitement d’événements en temps réel » ?

La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.

Peux-tu écrire et exécuter du code dans cette leçon AWS Solutions Architect ?

Oui. Chaque leçon AWS Solutions Architect inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.

Toutes les leçons de ce cours

  1. EventBridge : bus d’événements et règles
  2. Step Functions : orchestrer des flux de travail sans serveur
  3. Kinesis Data Streams pour le traitement d’événements en temps réel
  4. Modèles de chorégraphie ou d’orchestration
← Retour à AWS Solutions Architect