Kinesis Data Streams para el procesamiento de eventos en tiempo real
Produzca y consuma flujos de eventos de alto rendimiento con Kinesis Data Streams, gestione los shards para controlar el rendimiento y utilice Lambda como consumidor
Kinesis Data Streams para el procesamiento de eventos en tiempo real es una lección gratuita de Cloud & IT Cert Prep en CoddyKit. Esta es la lección 3 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.
Conceptos fundamentales de Kinesis Data Streams
Kinesis Data Streams (KDS) es un servicio de transmisión de datos duradero, ordenado y en tiempo real. Los datos se organizan en un stream compuesto por uno o varios shards. Cada shard es una secuencia ordenada de registros de datos. Los productores escriben registros en los shards mediante una partition key que determina qué shard recibe cada registro. Los consumidores leen los registros de los shards y los procesan en el orden en que llegaron a cada shard.
Capacidad y límites de rendimiento de los shards
Cada shard admite un rendimiento de escritura de 1 MB/s o 1.000 registros/s y un rendimiento de lectura de 2 MB/s (compartido entre todos los consumidores estándar de ese shard). La capacidad total del stream aumenta linealmente con el número de shards. Utilice la fórmula: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Si los productores alcanzan el límite de escritura, aparecen errores ProvisionedThroughputExceededException; resuélvalo dividiendo los shards o distribuyendo las partition keys de forma más uniforme.
# 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 5Partition keys y distribución de datos
La partition key es una cadena que Kinesis procesa mediante un hash (MD5) para determinar qué shard recibe un registro. Una partition key bien elegida distribuye los registros uniformemente entre los shards (prevención de hot shards). Para IoT, utilice el ID del dispositivo. Para flujos de clics, utilice el ID de sesión o el ID de usuario. Evite las keys con baja cardinalidad (por ejemplo, el nombre del país con solo 5 valores), ya que provocan hot shards: un shard recibe una cantidad desproporcionada de tráfico de escritura mientras los demás permanecen inactivos.
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
)Consumidores estándar frente a Enhanced Fan-Out
Los consumidores estándar comparten el rendimiento de lectura de 2 MB/s por shard mediante GetRecords con sondeo. Si tiene 3 consumidores en un shard y cada uno necesita 2 MB/s, se limitará su rendimiento al compartir un total de 2 MB/s. Enhanced Fan-Out (EFO) proporciona a cada consumidor registrado su propio canal de lectura dedicado de 2 MB/s mediante una conexión persistente de inserción HTTP/2 (SubscribeToShard). EFO añade un coste por consumidor, shard y hora, pero elimina por completo la contención de lectura.
# 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-telemetryLambda como consumidor de Kinesis
Lambda se integra de forma nativa con Kinesis Data Streams mediante un Event Source Mapping. Lambda sondea el stream, lee lotes de registros e invoca su función. Configure BatchSize (de 1 a 10.000 registros), StartingPosition (TRIM_HORIZON para los más antiguos, LATEST para los más recientes) y BisectBatchOnFunctionError para dividir los lotes que fallen. El factor de paralelización (de 1 a 10) permite que Lambda ejecute varias invocaciones simultáneas por shard para seguir el ritmo de los streams de alta velocidad.
# 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)
Kinesis Client Library (KCL) es un framework de aplicaciones para crear consumidores de Kinesis robustos en Java (o en varios lenguajes mediante MultiLangDaemon). KCL gestiona la enumeración de shards, la administración de leases (distribución de shards entre instancias de trabajo), los checkpoints del progreso en DynamoDB y las divisiones y fusiones ordenadas de shards. Cada worker de KCL procesa uno o varios shards, y KCL redistribuye automáticamente los shards cuando los workers se escalan horizontalmente o fallan. Para las aplicaciones consumidoras en producción, se recomienda KCL en lugar del sondeo directo mediante el SDK.
# 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();Retención y reproducción de datos
Kinesis Data Streams almacena los registros durante 24 horas de forma predeterminada (ampliable a 7 días o hasta 365 días con Long-Term Retention, con un coste adicional). A diferencia de SQS, los registros consumidos no se eliminan después del consumo: permanecen disponibles hasta que caduca la retención. Esto permite que varios consumidores lean los mismos registros de forma independiente y permite la reproducción al restablecer el checkpoint de un consumidor a un número de secuencia anterior, algo muy valioso para corregir errores o cargar datos históricos en servicios nuevos.
# 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 100Modo de capacidad bajo demanda
Kinesis Data Streams admite dos modos de capacidad. Modo Provisioned: usted administra manualmente el número de shards y paga por shard y hora. Modo On-Demand: Kinesis escala automáticamente la capacidad de los shards en función del rendimiento entrante (hasta 200 MB/s de escritura y 400 MB/s de lectura de forma predeterminada) y usted paga por cada GB de datos escrito y recuperado. El modo bajo demanda es ideal para patrones de tráfico variables o impredecibles en los que no desea administrar el escalado de los shards.
# 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_DEMANDGarantías de orden dentro de los shards
Kinesis garantiza el orden dentro de un shard: los registros con la misma partition key siempre se envían al mismo shard y se leen en el orden en que se escribieron. Sin embargo, no hay garantía de orden entre shards. Si su aplicación requiere un orden global para todos los registros, utilice un único shard (lo que limita el rendimiento a 1 MB/s) o rediseñe la aplicación para que el orden solo sea necesario dentro de un grupo de partition keys (por ejemplo, orden por dispositivo). Esta es una distinción habitual en los exámenes frente a SQS FIFO, que ofrece deduplicación y orden estrictos.
Seguridad: cifrado y VPC
Kinesis Data Streams cifra los registros en el servidor y en reposo mediante AWS KMS (CMK o una clave administrada por AWS) cuando se habilita el cifrado del lado del servidor. Todos los datos en tránsito se cifran mediante TLS. Para las aplicaciones que se ejecutan en una VPC y no deben enviar datos del stream a través de Internet público, utilice un VPC Interface Endpoint (PrivateLink) para Kinesis, de modo que el tráfico permanezca completamente dentro de la red troncal de AWS, algo importante para las cargas de trabajo reguladas.
# 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-abc123Kinesis Data Streams frente a SQS y Kafka
Para el examen SAA-C03, compare Kinesis Data Streams con las alternativas. Kinesis frente a SQS: Kinesis conserva el orden dentro de un shard y admite que varios consumidores lean los mismos datos; SQS elimina los mensajes después del consumo. Kinesis frente a MSK (Kafka): MSK es Apache Kafka administrado; utilícelo cuando necesite compatibilidad con el protocolo Kafka, configuraciones avanzadas de topics o esté migrando desde Kafka local. Utilice Kinesis para la transmisión nativa de AWS con una integración más estrecha con Lambda, Firehose y Flink. Elija Kinesis salvo que la pregunta mencione específicamente Kafka o requisitos de compatibilidad con Kafka.
Comprobació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 streams fragmentados, ordenados y duraderos, con una retención predeterminada de 24 horas y capacidad de reproducción; Enhanced Fan-Out proporciona a cada consumidor 2 MB/s dedicados por shard y elimina la contención de lectura; y el modo On-Demand escala automáticamente los shards para gestionar tráfico impredecible. A continuación, exploraremos los patrones de coreografía y orquestación en la arquitectura orientada a eventos.
Preguntas frecuentes
¿La lección «Kinesis Data Streams para el procesamiento de eventos en tiempo real» es gratis?
Sí — el texto completo de «Kinesis Data Streams para el procesamiento de eventos 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 Data Streams para el procesamiento de eventos en tiempo real»?
Produzca y consuma flujos de eventos de alto rendimiento con Kinesis Data Streams, gestione los shards para controlar el rendimiento y utilice Lambda como consumidor 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 3 de 4.
¿Cuánto tiempo toma la lección «Kinesis Data Streams para el procesamiento de eventos 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
- EventBridge: bus de eventos y reglas
- Step Functions: orquestación de flujos de trabajo sin servidor
- Kinesis Data Streams para el procesamiento de eventos en tiempo real
- Patrones de coreografía frente a orquestación