Kinesis Data Streams para processamento de eventos em tempo real
Produza e consuma fluxos de eventos de alto throughput com o Kinesis Data Streams, gerencie fragmentos para controlar o throughput e use o Lambda como consumidor.
Kinesis Data Streams para processamento de eventos em tempo real é uma aula grátis de AWS Solutions Architect no CoddyKit. Esta é a aula 3 de 4. Você pode ler a aula completa abaixo gratuitamente — depois pratica ao vivo no navegador com um editor de código integrado e um tutor de IA 24/7. Faz parte do caminho de aprendizado de AWS Solutions Architect, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de AWS Solutions Architect inclui 4 aulas no total.
Conceitos fundamentais de Kinesis Data Fluxos
Kinesis Data Fluxos (KDS) é um serviço durável e ordenado de transmissão de dados em tempo real. Os dados são organizados em um fluxo composto por um ou mais fragmentos. Cada fragmento é uma sequência ordenada de registros de dados. Os produtores colocam registros nos fragmentos usando uma chave de partição que determina qual fragmento receberá o registro. Os consumidores leem os registros dos fragmentos e os processam na ordem em que chegaram dentro de cada fragmento.
Capacidade dos fragmentos e limites de taxa de transferência
Cada fragmento oferece suporte a uma taxa de gravação de 1 MB/s ou 1.000 registros/s e a uma taxa de leitura de 2 MB/s (compartilhada entre todos os consumidores padrão desse fragmento). A capacidade total do fluxo aumenta linearmente com a quantidade de fragmentos. Use a fórmula: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Se os produtores atingirem o limite de gravação, serão exibidos erros ProvisionedThroughputExceededException — resolva o problema dividindo os fragmentos ou distribuindo as chaves de partição de maneira mais 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 5Chaves de partição e distribuição de dados
A chave de partição é uma cadeia de caracteres que o Kinesis transforma usando hash (MD5) para determinar qual fragmento receberá um registro. Uma chave de partição bem escolhida distribui os registros uniformemente entre os fragmentos (prevenção de fragmentos sobrecarregados). Para IoT, use o ID do dispositivo. Para fluxos de cliques, use o ID da sessão ou o ID do usuário. Evite chaves com baixa cardinalidade (por exemplo, o nome do país com apenas 5 valores), pois elas causam fragmentos sobrecarregados: um fragmento recebe uma parcela desproporcional do tráfego de gravação enquanto os outros ficam ociosos.
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 padrão versus Enhanced Fan-Out
Os consumidores padrão compartilham a taxa de leitura de 2 MB/s por fragmento usando GetRecords com sondagem. Se você tiver 3 consumidores em um fragmento e cada um precisar de 2 MB/s, eles sofrerão limitação ao compartilhar o total de 2 MB/s. O Enhanced Fan-Out (EFO) fornece a cada consumidor registrado seu próprio canal de leitura dedicado de 2 MB/s por meio de uma conexão persistente de envio HTTP/2 (SubscribeToShard). O EFO adiciona um custo por consumidor-fragmento-hora, mas elimina completamente a contenção de leitura.
# 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 do Kinesis
O Lambda integra-se nativamente ao Kinesis Data Fluxos por meio de um mapeamento de origem de eventos. O Lambda consulta o fluxo, lê lotes de registros e invoca sua função. Configure BatchSize (de 1 a 10.000 registros), StartingPosition (TRIM_HORIZON para os mais antigos, LATEST para os mais recentes) e BisectBatchOnFunctionError para dividir lotes com falha. O fator de paralelização (de 1 a 10) permite que o Lambda inicie várias invocações simultâneas por fragmento para acompanhar fluxos rápidos.
# 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)
A Kinesis Client Library (KCL) é uma estrutura de aplicação para criar consumidores robustos do Kinesis em Java (ou em várias linguagens por meio do MultiLangDaemon). A KCL gerencia a enumeração de fragmentos, o gerenciamento de concessões (distribuindo fragmentos entre instâncias de Worker), o registro de pontos de verificação do progresso no DynamoDB e as divisões ou mesclagens graduais de fragmentos. Cada Worker da KCL processa um ou mais fragmentos, e a KCL redistribui automaticamente os fragmentos à medida que os Workers são ampliados ou falham. A KCL é preferível à sondagem direta do SDK em aplicações de consumidores de produção.
# 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();Retenção e reprodução de dados
O Kinesis Data Fluxos armazena registros por 24 horas por padrão (com possibilidade de extensão para 7 dias ou até 365 dias com retenção de longo prazo, mediante custo adicional). Diferentemente do SQS, os registros consumidos não são excluídos após o consumo — eles permanecem disponíveis até o fim do período de retenção. Isso permite que vários consumidores leiam os mesmos registros de forma independente e possibilita a reprodução ao redefinir o ponto de verificação de um consumidor para um número de sequência anterior — algo extremamente útil para correções de erros ou para preencher dados de novos serviços.
# 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 capacidade sob demanda
O Kinesis Data Fluxos oferece suporte a dois modos de capacidade. Modo provisionado: você gerencia manualmente o número de fragmentos e paga por fragmento-hora. Modo sob demanda: o Kinesis ajusta automaticamente a capacidade dos fragmentos com base na taxa de entrada (até 200 MB/s de gravação e 400 MB/s de leitura por padrão), e você paga por GB de dados gravados e recuperados. O modo sob demanda é ideal para padrões de tráfego variáveis ou imprevisíveis nos quais você não deseja gerenciar o ajuste da capacidade dos fragmentos.
# 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_DEMANDGarantias de ordenação dentro dos fragmentos
O Kinesis garante a ordenação dentro de um fragmento — os registros com a mesma chave de partição sempre vão para o mesmo fragmento e são lidos na ordem em que foram gravados. No entanto, não há garantia de ordenação entre fragmentos. Se sua aplicação exigir ordenação global entre todos os registros, use um único fragmento (limitando a taxa de transferência a 1 MB/s) ou reprojete-a para que a ordenação seja necessária apenas dentro de um grupo de chaves de partição (por exemplo, ordenação por dispositivo). Essa é uma distinção comum em provas em comparação com o SQS FIFO, que oferece desduplicação e ordenação rigorosas.
Segurança: criptografia e VPC
O Kinesis Data Fluxos criptografa os registros no lado do servidor em repouso usando o AWS KMS (CMK ou chave gerenciada pela AWS) quando a criptografia no lado do servidor está habilitada. Todos os dados em trânsito são criptografados com TLS. Para aplicações executadas em uma VPC que não devem enviar dados do fluxo pela internet pública, use um endpoint de interface da VPC (PrivateLink) para o Kinesis, para que o tráfego permaneça totalmente dentro da rede principal da AWS — algo importante para cargas de trabalho regulamentadas.
# 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 Fluxos versus SQS versus Kafka
Para a prova SAA-C03, compare o Kinesis Data Fluxos com as alternativas. Kinesis versus SQS: o Kinesis preserva a ordem dentro de um fragmento e permite que vários consumidores leiam os mesmos dados; o SQS exclui as mensagens após o consumo. Kinesis versus MSK (Kafka): o MSK é o Apache Kafka gerenciado — use-o quando precisar de compatibilidade com o protocolo Kafka, configurações avançadas de tópicos ou estiver migrando do Kafka local. Use o Kinesis para transmissão nativa da AWS, com integração mais estreita ao Lambda, ao Firehose e ao Flink. Escolha o Kinesis, a menos que a pergunta mencione especificamente Kafka ou requisitos de compatibilidade com Kafka.
Verificação rápida
Teste sua compreensão dos conceitos de AWS Solutions Architect (SAA-C03) desta lição.
Recapitulação da lição
Nesta lição, você aprendeu que: Kinesis Data Fluxos fornece fluxos particionados, ordenados e duráveis, com retenção padrão de 24 horas e capacidade de reprodução; Enhanced Fan-Out fornece a cada consumidor 2 MB/s dedicados por fragmento, eliminando a contenção de leitura; e o modo sob demanda ajusta automaticamente a capacidade dos fragmentos para tráfego imprevisível. Em seguida, exploraremos os padrões de coreografia e orquestração na arquitetura orientada a eventos.
Perguntas Frequentes
A aula “Kinesis Data Streams para processamento de eventos em tempo real” é grátis?
Sim — o texto completo de “Kinesis Data Streams para processamento de eventos em tempo real” é grátis para ler aqui na web. Para praticá-la interativamente (um editor de código integrado e um tutor de IA 24/7) e desbloquear o restante do curso de AWS Solutions Architect, atualize para CoddyKit PRO. O curso de AWS Solutions Architect inclui 4 aulas no total.
O que vou aprender em “Kinesis Data Streams para processamento de eventos em tempo real”?
Produza e consuma fluxos de eventos de alto throughput com o Kinesis Data Streams, gerencie fragmentos para controlar o throughput e use o Lambda como consumidor. Você pratica AWS Solutions Architect com código prático que executa diretamente no navegador, e um tutor de IA 24/7 responde suas dúvidas enquanto trabalha na aula.
Preciso ter experiência prévia para começar AWS Solutions Architect?
Nenhuma experiência prévia é necessária. AWS Solutions Architect no CoddyKit é estruturado para alunos iniciantes até avançados, então você pode começar aqui ou desde o início e aprender no seu ritmo. Esta é a aula 3 de 4.
Quanto tempo leva a aula “Kinesis Data Streams para processamento de eventos em tempo real”?
A maioria das aulas CoddyKit leva cerca de 5–10 minutos. Cada uma é compacta e interativa, então você faz progresso constante e retoma exatamente de onde parou entre web e app.
Posso escrever e executar código nesta aula de AWS Solutions Architect?
Sim. Cada aula de AWS Solutions Architect inclui um editor de código integrado, então você escreve e executa código real direto no navegador e recebe feedback de IA instantaneamente — nenhuma configuração local necessária.
Todas as aulas deste curso
- EventBridge: barramento de eventos e regras
- Step Functions: orquestração de fluxos de trabalho sem servidor
- Kinesis Data Streams para processamento de eventos em tempo real
- Padrões de coreografia versus orquestração