Flux Kinesis, Firehose et analyse en temps réel
Ingérez des données en continu avec Kinesis Data Streams, transmettez-les à S3 ou Redshift avec Firehose et analysez-les en temps réel avec Managed Service for Apache Flink.
Flux Kinesis, Firehose et analyse en temps réel est une leçon Cloud & IT Cert Prep gratuite sur CoddyKit. Ceci est la leçon 4 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 Cloud & IT Cert Prep, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours Cloud & IT Cert Prep comprend 4 leçons au total.
Présentation de la famille Kinesis
Amazon Kinesis est une famille de services destinés à collecter, traiter et analyser des données diffusées en temps réel. Les trois services principaux sont : Kinesis Data Streams (traitement personnalisé à faible latence), Kinesis Data Firehose (distribution entièrement gérée vers S3/Redshift/OpenSearch) et Managed Service for Apache Flink (anciennement Kinesis Data Analytics) pour le SQL et le traitement Flink en temps réel. Chaque service répond à un besoin différent du pipeline de diffusion.
Architecture de Kinesis Data Streams
Un Kinesis Data Stream est un journal durable et ordonné, partitionné en fragments. Chaque fragment fournit un débit d’écriture de 1 Mo/s et un débit de lecture de 2 Mo/s. Les enregistrements sont conservés 24 heures par défaut, avec une extension possible à 7 ou 365 jours. Les producteurs écrivent les enregistrements dans un flux ; les consommateurs — Lambda, applications KCL, Firehose ou Flink — les lisent en parallèle depuis un ou plusieurs fragments. Les enregistrements sont immuables une fois écrits.
# 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'Fragments, débit et mise à l’échelle
Le nombre de fragments détermine le débit total d’un flux. Vous pouvez scinder un fragment pour doubler le débit ou fusionner deux fragments afin de réduire les coûts. Utilisez Enhanced Fan-Out pour attribuer à chaque consommateur enregistré son propre débit de lecture de 2 Mo/s, indépendamment des autres consommateurs, et éviter ainsi la limitation des lectures lorsque plusieurs applications consomment le même flux. Surveillez GetRecords.IteratorAgeMilliseconds pour détecter le retard des consommateurs.
# 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 : distribution gérée
Kinesis Data Firehose est un service entièrement géré qui capture, transforme et distribue des données diffusées vers des destinations telles que S3, Amazon Redshift, Amazon OpenSearch Service, Splunk et des points de terminaison HTTP. Aucun fragment n’est à gérer : Firehose adapte automatiquement sa capacité. Vous configurez une taille de tampon (1 à 128 Mo) et un intervalle de tampon (60 à 900 secondes) ; Firehose distribue les données dès que l’une des deux limites est atteinte.
# 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"
}'Transformation des données Firehose avec Lambda
Firehose peut invoquer une fonction Lambda pour chaque lot d’enregistrements avant leur distribution, afin de transformer, d’enrichir ou de filtrer les données en transit. Les cas d’utilisation courants comprennent la conversion de JSON en Parquet (via un schéma Glue), le masquage des champs PII ou la suppression des événements de faible valeur. Les enregistrements dont la transformation échoue peuvent être envoyés à un préfixe S3 d’erreur distinct pour retraitement, afin qu’aucune donnée ne soit perdue.
# 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 (anciennement Kinesis Data Analytics) exécute des applications Apache Flink sur une infrastructure entièrement gérée. Utilisez-le pour les analyses avec état en temps réel : agrégations sur des fenêtres glissantes, détection d’anomalies, recherche de motifs dans des séquences d’événements et jointure de données diffusées avec des tables de référence. Vous écrivez le code Flink en Java, Python ou Scala, et Flink gère les points de contrôle ainsi que l’état avec une sémantique exactement une fois.
# 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);Choisir entre Streams et Firehose
Pour l’examen SAA-C03, vous devez connaître les critères de décision. Utilisez Kinesis Data Streams lorsque vous avez besoin d’une latence inférieure à la seconde, de plusieurs consommateurs lisant simultanément ou d’une logique de traitement personnalisée avec un contrôle total de la conservation et de la relecture. Utilisez Firehose lorsque vous devez simplement distribuer de manière fiable des données diffusées vers S3, Redshift ou OpenSearch, avec peu de code, une transformation facultative en transit et une mise à l’échelle automatique, au prix d’une latence plus élevée (au moins 60 secondes).
Kinesis ou SQS : le choix classique à l’examen
Une question fréquente du SAA-C03 vous demande de choisir entre Kinesis et SQS. Principales différences : Kinesis préserve l’ordre au sein d’un fragment, permet à plusieurs consommateurs de lire simultanément les mêmes données et conserve les enregistrements pour les relire. SQS supprime les messages après leur consommation (sans relecture), son mode FIFO garantit un ordre strict et il convient mieux au découplage des microservices. Si le scénario mentionne l’analyse en temps réel ou la relecture, choisissez Kinesis.
Producteurs Kinesis : SDK et KPL
La Kinesis Producer Library (KPL) est un client à haut débit permettant d’écrire dans Kinesis Data Streams depuis des applications. KPL regroupe automatiquement plusieurs petits enregistrements en un seul appel d’API (jusqu’à 1 Mo) et gère les nouvelles tentatives avec temporisation progressive. Cela réduit considérablement le coût PUT par enregistrement et augmente le débit par fragment. Utilisez KPL pour les producteurs à fort volume, comme les flux de clics web, les capteurs IoT ou les pipelines de journaux.
# 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() { ... });Schéma Firehose vers Redshift
Une architecture courante consiste à utiliser Firehose comme pipeline géré depuis des flux Kinesis ou des producteurs directs vers Amazon Redshift pour l’entreposage de données. Firehose écrit d’abord les données dans un bucket S3 intermédiaire de préparation, puis émet une commande COPY pour charger les données dans Redshift. Il s’agit de la méthode la plus efficace pour charger en bloc des données diffusées dans Redshift : les insertions directes ligne par ligne seraient extrêmement lentes en raison du surcoût associé à chaque ligne.
# 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"
}
}'Surveiller Kinesis avec CloudWatch
Principales métriques CloudWatch pour Kinesis : IncomingBytes et IncomingRecords pour mesurer le débit des producteurs, GetRecords.IteratorAgeMilliseconds pour mesurer le retard des consommateurs (une valeur élevée signifie qu’ils ne suivent pas le rythme) et WriteProvisionedThroughputExceeded pour détecter le dépassement des limites des fragments par les producteurs. Définissez des alarmes CloudWatch sur l’âge de l’itérateur et le dépassement du débit alloué afin de déclencher la mise à l’échelle automatique des fragments via 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-alertsVé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 une diffusion durable, ordonnée et partitionnée, avec plusieurs consommateurs et la possibilité de relire les données, que Kinesis Data Firehose permet une distribution entièrement gérée et sans code vers S3, Redshift et OpenSearch, et que Managed Service for Apache Flink permet d’effectuer des analyses avec état en temps réel sur des flux. Nous allons maintenant découvrir les 7 R de la stratégie de migration vers le cloud.
Apprends Cloud & IT Cert Prep avec un tuteur IA — gratuit
Écris et exécute du vrai code dans ton navigateur, obtiens de l'aide instantanée d'un tuteur IA disponible 24h/24, et reprends là où tu t'es arrêté sur le web ou dans l'app.
- Cours
- 150
- Leçons
- 600
Questions Fréquemment Posées
La leçon « Flux Kinesis, Firehose et analyse en temps réel » est-elle gratuite ?
Oui — le texte complet de « Flux Kinesis, Firehose et analyse 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 Cloud & IT Cert Prep, passe à CoddyKit PRO. Le cours Cloud & IT Cert Prep comprend 4 leçons au total.
Qu'est-ce que j'apprendrai dans « Flux Kinesis, Firehose et analyse en temps réel » ?
Ingérez des données en continu avec Kinesis Data Streams, transmettez-les à S3 ou Redshift avec Firehose et analysez-les en temps réel avec Managed Service for Apache Flink. Tu pratiques Cloud & IT Cert Prep 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 Cloud & IT Cert Prep ?
Aucune expérience préalable n'est requise. Cloud & IT Cert Prep 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 4 sur 4.
Combien de temps prend la leçon « Flux Kinesis, Firehose et analyse 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 Cloud & IT Cert Prep ?
Oui. Chaque leçon Cloud & IT Cert Prep 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
- Créer un lac de données sur S3
- AWS Glue : ETL et catalogue de données
- Amazon Athena : SQL sans serveur sur S3
- Flux Kinesis, Firehose et analyse en temps réel