Kinesis Streams, Firehose y análisis en tiempo real
Ingesta datos en streaming con Kinesis Data Streams, entréguelos en S3 o Redshift con Firehose y analícelos en tiempo real con Managed Service for Apache Flink
Kinesis Streams, Firehose y análisis en tiempo real es una lección gratuita de Cloud & IT Cert Prep en CoddyKit. Esta es la lección 4 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de Cloud & IT Cert Prep, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de Cloud & IT Cert Prep incluye 4 lecciones en total.
Descripción general de la familia Kinesis
Amazon Kinesis es una familia de servicios para recopilar, procesar y analizar datos de streaming en tiempo real. Los tres servicios principales son: Kinesis Data Streams (procesamiento personalizado de baja latencia), Kinesis Data Firehose (entrega totalmente administrada a S3/Redshift/OpenSearch) y Managed Service for Apache Flink (anteriormente Kinesis Data Analytics), destinado a SQL y procesamiento con Flink en tiempo real. Cada servicio está orientado a un punto diferente de la canalización de streaming.
Arquitectura de Kinesis Data Streams
Un Kinesis Data Stream es un registro duradero y ordenado, dividido en shards. Cada shard proporciona un rendimiento de escritura de 1 MB/s y de lectura de 2 MB/s. Los registros de datos se conservan durante 24 horas de forma predeterminada (se puede ampliar a 7 o 365 días). Los productores escriben registros en un stream; los consumidores —Lambda, aplicaciones de KCL, Firehose o Flink— leen de uno o varios shards en paralelo. Los registros son inmutables una vez escritos.
# 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, rendimiento y escalado
El número de shards determina el rendimiento total de un stream. Puede dividir un shard para duplicar el rendimiento o combinar dos shards para reducir los costes. Utilice Enhanced Fan-Out para proporcionar a cada consumidor registrado su propio rendimiento de lectura de 2 MB/s, independientemente de los demás consumidores, y eliminar la limitación de lectura cuando varias aplicaciones consuman el mismo stream. Supervise GetRecords.IteratorAgeMilliseconds para detectar el retraso de los consumidores.
# 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: entrega administrada
Kinesis Data Firehose es un servicio totalmente administrado que captura, transforma y entrega datos de streaming a destinos como S3, Amazon Redshift, Amazon OpenSearch Service, Splunk y puntos de conexión HTTP. No hay shards que administrar: Firehose escala automáticamente. Configure un tamaño de búfer (1–128 MB) y un intervalo de búfer (60–900 segundos); Firehose realiza la entrega cuando se alcanza primero cualquiera de los dos límites.
# 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"
}'Transformación de datos de Firehose con Lambda
Firehose puede invocar una función Lambda para cada lote de registros antes de la entrega, con el fin de transformar, enriquecer o filtrar los datos durante el procesamiento. Entre los casos de uso habituales se incluyen convertir JSON a Parquet (mediante un esquema de Glue), enmascarar campos de información personal identificable (PII) o descartar eventos de poco valor. Los registros cuya transformación falla se pueden enviar opcionalmente a un prefijo de errores independiente en S3 para volver a procesarlos, de modo que no se pierdan datos.
# 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 (anteriormente Kinesis Data Analytics) ejecuta aplicaciones de Apache Flink en una infraestructura totalmente administrada. Utilícelo para análisis de estado en tiempo real: agregaciones en ventanas deslizantes, detección de anomalías, coincidencia de patrones en secuencias de eventos y combinación de datos de streaming con tablas de referencia. Escriba el código de Flink en Java, Python o Scala; Flink administra los puntos de control y el estado con procesamiento exactamente una vez.
# 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);Elección entre Streams y Firehose
Para el examen SAA-C03, debe conocer los criterios de decisión. Utilice Kinesis Data Streams cuando necesite una latencia inferior a un segundo, varios consumidores que lean simultáneamente o lógica de procesamiento personalizada con control total sobre la conservación y la reproducción. Utilice Firehose cuando solo necesite entregar datos de streaming de forma fiable a S3, Redshift u OpenSearch, con un mínimo de código, transformación opcional durante el procesamiento y escalado automático, aunque con una latencia mayor (60 segundos o más).
Kinesis frente a SQS: la elección clásica del examen
Una pregunta habitual del SAA-C03 le pide elegir entre Kinesis y SQS. Diferencias clave: Kinesis conserva el orden dentro de un shard, admite varios consumidores que leen los mismos datos simultáneamente y conserva los registros para su reproducción. SQS elimina los mensajes después de consumirlos (sin reproducción), el modo FIFO garantiza un orden estricto y es más adecuado para desacoplar microservicios. Si el escenario menciona análisis en tiempo real o reproducción, elija Kinesis.
Productores de Kinesis: SDK y KPL
La Kinesis Producer Library (KPL) es un cliente de alto rendimiento para escribir en Kinesis Data Streams desde aplicaciones. KPL agrega automáticamente varios registros pequeños en una única llamada a la API (hasta 1 MB) y gestiona los reintentos con retroceso. Esto reduce considerablemente el coste de PUT por registro y aumenta el rendimiento por shard. Utilice KPL para productores de gran volumen, como flujos de clics web, sensores de IoT o canalizaciones de registros.
# 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() { ... });Patrón de Firehose a Redshift
Una arquitectura habitual consiste en utilizar Firehose como canalización administrada desde streams de Kinesis o productores directos hacia Amazon Redshift para el almacenamiento de datos. Firehose escribe primero los datos en un bucket de preparación intermedio de S3 y, a continuación, emite un comando COPY para cargar los datos en Redshift. Es la forma más eficiente de cargar datos de streaming en bloque en Redshift; las inserciones directas fila por fila en Redshift serían extremadamente lentas debido a la sobrecarga de cada fila.
# 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"
}
}'Supervisión de Kinesis con CloudWatch
Métricas clave de CloudWatch para Kinesis: IncomingBytes e IncomingRecords para medir el rendimiento de los productores, GetRecords.IteratorAgeMilliseconds para medir el retraso de los consumidores (un valor alto significa que los consumidores no pueden mantener el ritmo) y WriteProvisionedThroughputExceeded para detectar cuándo los productores alcanzan los límites de los shards. Configure alarmas de CloudWatch sobre la antigüedad del iterador y el rendimiento aprovisionado excedido para activar el escalado automático de shards mediante 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-alertsComprobación rápida
Compruebe su comprensión de los conceptos de AWS Solutions Architect (SAA-C03) de esta lección.
Resumen de la lección
En esta lección ha aprendido que: Kinesis Data Streams proporciona streaming duradero, ordenado y dividido en shards, con varios consumidores y capacidad de reproducción, Kinesis Data Firehose ofrece una entrega totalmente administrada y sin código a S3, Redshift y OpenSearch, y Managed Service for Apache Flink permite realizar análisis de estado en tiempo real sobre streams. A continuación, exploraremos las 7 R de la estrategia de migración a la nube.
Preguntas frecuentes
¿La lección «Kinesis Streams, Firehose y análisis en tiempo real» es gratis?
Sí — el texto completo de «Kinesis Streams, Firehose y análisis en tiempo real» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de Cloud & IT Cert Prep, actualiza a CoddyKit PRO. El curso de Cloud & IT Cert Prep incluye 4 lecciones en total.
¿Qué aprenderé en «Kinesis Streams, Firehose y análisis en tiempo real»?
Ingesta datos en streaming con Kinesis Data Streams, entréguelos en S3 o Redshift con Firehose y analícelos en tiempo real con Managed Service for Apache Flink Practicas Cloud & IT Cert Prep con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.
¿Necesito experiencia previa para empezar Cloud & IT Cert Prep?
No se requiere experiencia previa. Cloud & IT Cert Prep en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 4 de 4.
¿Cuánto tiempo toma la lección «Kinesis Streams, Firehose y análisis en tiempo real»?
La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.
¿Puedo escribir y ejecutar código en esta lección de Cloud & IT Cert Prep?
Sí. Cada lección de Cloud & IT Cert Prep incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.
Todas las lecciones de este curso
- Creación de un lago de datos en S3
- AWS Glue: ETL y catálogo de datos
- Amazon Athena: SQL sin servidor en S3
- Kinesis Streams, Firehose y análisis en tiempo real