Streams do Kinesis, Firehose e análise em tempo real
Ingira dados em fluxo com o Kinesis Data Streams, entregue-os ao S3 ou ao Redshift com o Firehose e analise-os em tempo real com o Managed Service for Apache Flink.
Streams do Kinesis, Firehose e análise em tempo real é uma aula grátis de Cloud & IT Cert Prep no CoddyKit. Esta é a aula 4 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 Cloud & IT Cert Prep, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de Cloud & IT Cert Prep inclui 4 aulas no total.
Visão geral da família Kinesis
O Amazon Kinesis é uma família de serviços para coletar, processar e analisar dados transmitidos em tempo real. Os três serviços principais são: Kinesis Data Streams (processamento personalizado com baixa latência), Kinesis Data Firehose (entrega totalmente gerenciada ao S3/Redshift/OpenSearch) e Managed Service for Apache Flink (anteriormente Kinesis Data Analytics), para SQL e processamento com Flink em tempo real. Cada serviço atende a uma finalidade diferente no fluxo de processamento de dados.
Arquitetura do Kinesis Data Streams
Um Kinesis Data Stream é um registro durável e ordenado, particionado em shards. Cada shard fornece uma capacidade de gravação de 1 MB/s e de leitura de 2 MB/s. Os registros de dados são mantidos por 24 horas por padrão (período que pode ser ampliado para 7 ou 365 dias). Os produtores gravam registros em um fluxo; os consumidores — Lambda, aplicações KCL, Firehose ou Flink — leem de um ou mais shards em paralelo. Depois de gravados, os registros são imutáveis.
# 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, capacidade e escalabilidade
A quantidade de shards determina a capacidade total do fluxo. Você pode dividir um shard para dobrar a capacidade ou mesclar dois shards para reduzir custos. Use Enhanced Fan-Out para fornecer a cada consumidor registrado sua própria capacidade de leitura de 2 MB/s, independentemente dos demais consumidores, eliminando a limitação de leitura quando várias aplicações consomem o mesmo fluxo. Monitore GetRecords.IteratorAgeMilliseconds para detectar o atraso dos 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 gerenciada
O Kinesis Data Firehose é um serviço totalmente gerenciado que captura, transforma e entrega dados transmitidos a destinos como S3, Amazon Redshift, Amazon OpenSearch Service, Splunk e endpoints HTTP. Não há shards para gerenciar — o Firehose escala automaticamente. Você configura um tamanho do buffer (1–128 MB) e um intervalo do buffer (60–900 segundos); o Firehose faz a entrega assim que qualquer um dos dois limites é atingido.
# 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"
}'Transformação de dados do Firehose com Lambda
O Firehose pode invocar uma função Lambda em cada lote de registros antes da entrega, para transformar, enriquecer ou filtrar os dados durante o processamento. Casos de uso comuns incluem converter JSON em Parquet (por meio do esquema do Glue), mascarar campos de PII ou descartar eventos de baixo valor. Os registros que falham na transformação são opcionalmente enviados a um prefixo de erro separado no S3 para reprocessamento, portanto nenhum dado é perdido.
# 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
O Amazon Managed Service for Apache Flink (anteriormente Kinesis Data Analytics) executa aplicações Apache Flink em uma infraestrutura totalmente gerenciada. Use-o para análises com estado em tempo real: agregações em janelas deslizantes, detecção de anomalias, correspondência de padrões em sequências de eventos e junção de dados transmitidos com tabelas de referência. Você escreve o código Flink em Java, Python ou Scala, e o Flink gerencia pontos de verificação e o estado com garantia de processamento exatamente uma 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);Escolhendo entre Streams e Firehose
Para a prova SAA-C03, conheça os critérios de decisão. Use Kinesis Data Streams quando precisar de latência inferior a um segundo, vários consumidores lendo simultaneamente ou lógica de processamento personalizada com controle total sobre retenção e reprodução. Use o Firehose quando precisar apenas entregar dados transmitidos de forma confiável ao S3, Redshift ou OpenSearch, com pouco código, transformação opcional durante o processamento e escalabilidade automática, aceitando uma latência maior (60 segundos ou mais).
Kinesis versus SQS: a escolha clássica da prova
Uma pergunta comum da SAA-C03 pede que você escolha entre Kinesis e SQS. Principais diferenças: o Kinesis preserva a ordem dentro de um shard, permite que vários consumidores leiam os mesmos dados simultaneamente e mantém os registros para reprodução. O SQS remove as mensagens depois que são consumidas (sem reprodução), o modo FIFO garante ordenação estrita e o serviço é mais adequado para desacoplar microsserviços. Se o cenário mencionar análise em tempo real ou reprodução, escolha o Kinesis.
Produtores do Kinesis: SDK e KPL
A Kinesis Producer Library (KPL) é um cliente de alta capacidade para gravar em Kinesis Data Streams a partir de aplicações. A KPL agrega automaticamente vários registros pequenos em uma única chamada de API (até 1 MB) e gerencia novas tentativas com recuo progressivo. Isso reduz significativamente o custo de PUT por registro e aumenta a capacidade por shard. Use a KPL para produtores de alto volume, como fluxos de cliques da web, sensores de IoT ou fluxos 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() { ... });Padrão Firehose para Redshift
Uma arquitetura comum usa o Firehose como um fluxo de processamento gerenciado dos fluxos do Kinesis ou de produtores diretos para o Amazon Redshift, formando um armazém de dados. Primeiro, o Firehose grava os dados em um bucket intermediário de preparação no S3; em seguida, emite um comando COPY para carregar os dados no Redshift. Essa é a forma mais eficiente de carregar dados transmitidos em massa no Redshift — inserções diretas linha a linha seriam extremamente lentas devido à sobrecarga por linha.
# 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"
}
}'Monitoramento do Kinesis com CloudWatch
Métricas importantes do CloudWatch para o Kinesis: IncomingBytes e IncomingRecords para medir a capacidade dos produtores, GetRecords.IteratorAgeMilliseconds para medir o atraso dos consumidores (um valor alto significa que eles não conseguem acompanhar) e WriteProvisionedThroughputExceeded para detectar quando os produtores estão atingindo os limites dos shards. Configure alarmes do CloudWatch para a idade do iterador e para a capacidade provisionada excedida, acionando o Auto Scaling dos shards por meio do 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-alertsVerificação rápida
Teste sua compreensão dos conceitos de AWS Solutions Architect (SAA-C03) desta lição.
Resumo da lição
Nesta lição, você aprendeu que: o Kinesis Data Streams fornece transmissão durável, ordenada e particionada, com vários consumidores e capacidade de reprodução, o Kinesis Data Firehose oferece entrega totalmente gerenciada e sem código ao S3, Redshift e OpenSearch e o Managed Service for Apache Flink permite análises de transmissão em tempo real com estado. Em seguida, exploraremos os 7 Rs da estratégia de migração para a nuvem.
Perguntas Frequentes
A aula “Streams do Kinesis, Firehose e análise em tempo real” é grátis?
Sim — o texto completo de “Streams do Kinesis, Firehose e análise 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 Cloud & IT Cert Prep, atualize para CoddyKit PRO. O curso de Cloud & IT Cert Prep inclui 4 aulas no total.
O que vou aprender em “Streams do Kinesis, Firehose e análise em tempo real”?
Ingira dados em fluxo com o Kinesis Data Streams, entregue-os ao S3 ou ao Redshift com o Firehose e analise-os em tempo real com o Managed Service for Apache Flink. Você pratica Cloud & IT Cert Prep 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 Cloud & IT Cert Prep?
Nenhuma experiência prévia é necessária. Cloud & IT Cert Prep 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 4 de 4.
Quanto tempo leva a aula “Streams do Kinesis, Firehose e análise 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 Cloud & IT Cert Prep?
Sim. Cada aula de Cloud & IT Cert Prep 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
- Criação de um data lake no S3
- AWS Glue: ETL e catálogo de dados
- Amazon Athena: SQL sem servidor no S3
- Streams do Kinesis, Firehose e análise em tempo real