Lab 63 — Ingestão em streaming: shard, ordem e reprocesso
O problema, e a empresa que o tem
A Cadência entra na Semana Cadência — a campanha anual de maior volume — com o painel de risco alimentado por um consumidor de eventos de pedido: pedido_criado, pedido_pago, pedido_separado, pedido_cancelado e pedido_entregue. Cada evento sai da Checkout API (L01/L03) direto para um stream Kinesis assim que acontece — não há tabela no meio, como havia no L29. É a aplicação que decide publicar.
No primeiro dia de pico, dois sintomas apareceram juntos e pareceram o mesmo bug. O time de risco relatou pedidos cancelados minutos depois do pagamento — um padrão clássico de estorno fraudulento — sem que o alerta automático disparasse. E quando o alerta disparava, às vezes chegava com a ordem trocada: o sistema via pedido_cancelado antes de qualquer registro de pedido_pago para o mesmo pedido, o que não deveria ser logicamente possível.
A investigação levou dias porque nenhum log dizia "evento perdido". Os logs do produtor mostravam sucesso em toda chamada. A causa real estava numa única linha de código, escrita com uma intenção razoável — agrupar eventos do mesmo tipo — e que gerava os dois sintomas ao mesmo tempo, pelo mesmo motivo.
O que este laboratório NÃO é
Não é o L29. Lá o Streams captura mudança de ESTADO no banco — a tabela muda, o CDC reage. Aqui não existe necessariamente uma escrita em banco por trás do evento: é a aplicação que decide emitir "pedido separado" como parte do fluxo, e a chave de partição é uma decisão do código do produtor, não algo herdado da chave primária de uma tabela. Confundir os dois faz alguém procurar uma tabela que não existe.
O que você vai conseguir fazer
Objetivos verificáveis: cada um se prova com um comando na seção de implantação.
- Explicar por que PutRecords devolve HTTP 200 mesmo com registros que falharam.
- Derivar, a partir do hash MD5 da chave, por que a mesma chave sempre cai no mesmo shard.
- Diagnosticar um hot shard a partir de métricas de throughput por ShardId.
- Escolher chave de partição pela cardinalidade E pela pergunta de ordem que ela precisa responder.
- Explicar por que resharding não corrige uma chave de partição genuinamente enviesada.
- Configurar retenção estendida e reprocessar um trecho do stream a partir de TRIM_HORIZON ou AT_TIMESTAMP.
- Tratar falha parcial de PutRecords reenviando só os registros marcados como falhos.
- Configurar um destino de falha no event source mapping para que um registro ruim não trave o shard.
O que a certificação cobra disto
| Conceito | Certificação | Como aparece aqui | O que dominar |
|---|---|---|---|
| Capacidade fixa por shard | DVA-C02, SAP-C02 | 1.000 registros/s ou 1 MiB/s por shard, o que vier primeiro | por que dobrar shard_count não ajuda se a chave concentra tudo em um só |
| Hash da chave de partição | DVA-C02, SAP-C02 | MD5 mapeado contra o intervalo de hash do shard | a mesma chave produz sempre o mesmo hash, e cai sempre no mesmo shard |
| PutRecords com falha parcial | DVA-C02 | FailedRecordCount e ErrorCode por registro | por que HTTP 200 não é prova de que o lote inteiro foi gravado |
| Ordem dentro do shard | SAP-C02, MLA-C01 | garantida só lá — nunca entre shards | por que "ordem global do stream" não é um conceito que exista |
| ParallelizationFactor | DVA-C02 | até 10 lotes concorrentes por shard, preservando ordem por chave | a diferença entre concorrência (entre chaves) e ordenação (dentro da chave) |
| Tipos de shard iterator | SAP-C02 | TRIM_HORIZON, LATEST, AT_TIMESTAMP | quando cada um é a escolha certa para começar a ler |
| Retenção e reprocesso | SAP-C02, MLA-C01 | 24h padrão, extensível até 8.760h (365 dias) | a retenção é o que HABILITA reprocessar; sem ela o dado já se foi |
| Resharding e chave enviesada | SAP-C02 | on-demand reparte igualmente e não isola a chave quente | por que a correção real é na chave, não no número de shards |
Onde isto costuma ser cobrado errado
A pergunta clássica dá N shards provisionados e pede o throughput total do stream, esperando a soma de N × 1 MiB/s. A resposta certa depende da distribuição real da chave: sob uma chave enviesada, o throughput útil é limitado pelo shard mais quente, não pela soma — os outros shards sobram capacidade que ninguém usa.
Requisitos, e como cada um muda o desenho
Requisito que não vira uma linha de configuração é intenção. A coluna da direita é onde cada um deixou marca no Terraform ou no código.
| Requisito | Valor declarado | O que ele decide no desenho |
|---|---|---|
| Nenhum evento perdido em silêncio | obrigatório | produtor inspeciona FailedRecordCount e reenvia só os registros que falharam |
| Ordem por pedido reconstituível | obrigatório | chave de partição = id do pedido, nunca o tipo do evento |
| Pico sustentado de até 3× o volume normal | declarado | 8 shards provisionados com margem, mais SplitShard manual como plano de contingência |
| Corrigir bug sem reingestão pela aplicação | desejado | retenção estendida a 168h (7 dias) + mapping de reprocesso a partir de TRIM_HORIZON/AT_TIMESTAMP |
| Alerta de risco em segundos, não por varredura | obrigatório | Lambda via event source mapping, substituindo a função agendada anterior |
| Registro ruim não pode travar o shard inteiro | obrigatório | bisect_batch_on_function_error habilitado, com destino de falha em SQS |
| Nenhum segredo em variável de ambiente | obrigatório | ARN de tópico/fila por referência; criptografia do stream e da fila via KMS |
Arquitetura mínima: a chave é o tipo do evento
Este é o desenho que a Cadência tinha no primeiro dia de campanha, e ele é legítimo como ponto de partida: publica de verdade, com uma decisão de particionamento que parece razoável. O laboratório começa medindo essa decisão — porque um número torna o hot shard discutível, e "parece que está perdendo evento" não.
- → inicia e evolui o pedido
- → PutRecords, chave = tipo_evento
- → PutRecords, chave = tipo_evento
- → PutRecords, chave = tipo_evento
- → GetRecords, ordem só garantida aqui dentro
- → GetRecords, ordem só garantida aqui dentro
- → GetRecords, ordem só garantida aqui dentro
- → IncomingBytes e IncomingRecords, por shard
- Fora da AWS
- Compute
- Conceito de arquitetura
- Gestão e governança
Este desenho publica de verdade, com uma decisão só: a chave de partição é o TIPO do evento. Parece organização — agrupar eventos parecidos no mesmo lugar — e é exatamente essa lógica que concentra o pico inteiro em um shard e espalha os quatro estágios de um MESMO pedido em shards diferentes. Percorra os passos: o produtor nunca vê o próprio erro.
- A chave é o tipo do evento, não o pedido. Parece uma decisão razoável: agrupar "todo pedido_pago" num lugar só soa organizado. O problema é a cardinalidade — só existem 5 valores possíveis de tipo de evento, contra milhares de pedidos por minuto no pico.
- O hash decide, e ele é determinístico. O Kinesis aplica MD5 à chave de partição e mapeia o resultado contra o intervalo de hash de cada shard. O hash da string `"pedido_criado"` é sempre o MESMO valor — então todo registro com essa chave cai, sempre, no mesmo shard.
- Um shard carrega a criação de pedido do dia inteiro. Cada shard tem capacidade de escrita FIXA: 1.000 registros por segundo ou 1 MiB/s, o que vier primeiro — independente de quantos shards o stream tem no total. Concentrar um tipo de evento num shard só é concentrar o throughput do stream inteiro nesse teto.
- PutRecords devolve 200 mesmo quando um registro falhou. A chamada é um lote de até 500 registros, e a API processa cada um de forma independente: um registro throttled não derruba a chamada. O corpo da resposta traz `FailedRecordCount` e, por registro, um `ErrorCode` — e nada disso aparece se o código só confere o status HTTP.
- O mesmo pedido, em shards diferentes. Como a chave é o TIPO do evento, o "criado" e o "cancelado" do mesmo pedido vão para shards distintos. Cada shard é lido no seu próprio ritmo — a fila de leitura da Lambda em um shard não espera a do outro — então nada impede o consumidor de processar o cancelamento antes de ter visto o pagamento.
- Os shards ociosos não emprestam capacidade a ninguém. O throughput do stream não é a SOMA dos três shards — é limitado pelo pior deles no tráfego real. Enquanto o shard 0 satura, os shards 1 e 2 sobram capacidade e não há como redirecionar tráfego para eles sem mudar a chave.
- Por que alguém particiona por tipo de evento. Porque parece a chave "certa" do ponto de vista de quem vai LER — separar por tipo facilita filtrar depois. É o raciocínio invertido: a chave de partição decide DISTRIBUIÇÃO de escrita, não organização de leitura, e as duas coisas raramente pedem a mesma chave.
#!/usr/bin/env bash
# medir-shard-quente.sh — prova o hot shard com numero, nao com suspeita
set -euo pipefail
STREAM="${STREAM:?defina STREAM}"
INICIO=$(date -u -d '5 minutes ago' +%Y-%m-%dT%H:%M:%S) # GNU date; em macOS use -v-5M
for SHARD_ID in $(aws kinesis list-shards --stream-name "$STREAM" \
--query 'Shards[].ShardId' --output text); do
BYTES=$(aws cloudwatch get-metric-statistics \
--namespace AWS/Kinesis --metric-name IncomingBytes \
--dimensions Name=StreamName,Value="$STREAM" Name=ShardId,Value="$SHARD_ID" \
--start-time "$INICIO" --end-time "$(date -u +%Y-%m-%dT%H:%M:%S)" \
--period 300 --statistics Sum --query 'Datapoints[0].Sum' --output text)
echo "$SHARD_ID: ${BYTES:-0} bytes em 5 min (teto do shard: 300.000.000 bytes)"
done
# Esperado ANTES da correcao: um shard perto do teto (1 MiB/s x 300s ~ 300 MB),
# os outros dois em uma fracao pequena disso. Depois de trocar a chave para o
# id do pedido, os oito shards do desenho de producao ficam dentro de 15% um
# do outro — e essa e a prova de que a distribuicao mudou, nao so a contagem.
A perda real não está no Kinesis — está no que o produtor faz com a resposta
Se o código do produtor não confere `FailedRecordCount` e também não tem retry algum, um registro throttled é perdido de verdade e para sempre: não há como recuperá-lo depois, porque ele nunca chegou a existir no stream. Isso é distinto de um dado que existe e está apenas fora da janela de retenção — aqui, ele nunca foi gravado. É a perda irreversível deste laboratório, e ela mora inteiramente no lado do produtor, não do Kinesis.
Arquitetura para produção
Cada peça nova abaixo rastreia a uma linha da tabela de requisitos. A diferença estrutural em relação ao desenho anterior não é um shard a mais — é a chave de partição decidindo tudo de novo, mais dois caminhos de consumo que não existiam.
- → inicia e evolui o pedido
- → PutRecords, chave = id do pedido
- → PutRecords, chave = id do pedido
- → GetRecords, agora em ordem por pedido
- → GetRecords, agora em ordem por pedido
- → publica alerta de risco
- → lote esgotou as tentativas
- → GetShardIterator(TRIM_HORIZON), só após a correção
- → grava resultado corrigido
- → IteratorAgeMilliseconds e throughput excedido
- Fora da AWS
- Compute
- Conceito de arquitetura
- Integração de apps
- Banco de dados
- Gestão e governança
A diferença estrutural não é mais um shard: é a IDENTIDADE que decide o shard. Com a chave = id do pedido, os eventos de UM pedido caem sempre no MESMO shard, e a ordem entre eles passa a existir de verdade. Um caminho novo e temporário — o reprocesso — existe só porque a retenção guardou o dado. Percorra os passos: cada peça nova rastreia a um requisito da seção anterior.
- A chave passa a ser o id do pedido. Cardinalidade alta: cada pedido tem um id distinto, então o hash MD5 espalha o tráfego quase uniformemente pelos 8 shards. Não existe mais um shard "dono" de um tipo de evento — cada shard carrega uma fatia proporcional de PEDIDOS.
- O produtor lê a resposta antes de dar por certo. Depois de cada PutRecords, o código confere `FailedRecordCount`. Se for maior que zero, reconstrói um novo lote só com os registros que falharam — não reenvia os 500 originais — e tenta de novo com espera crescente.
- Ordem por pedido existe porque a chave garante o mesmo shard. Como todo evento de um pedido cai no mesmo shard, e o Kinesis preserva ordem dentro de um shard, a sequência criado → pago → cancelado chega ao consumidor na ordem em que aconteceu — não por sorte, por garantia do serviço.
- Falha do consumidor não é descarte: é redirecionamento. O mapping bisecta o lote com erro e tenta de novo um número limitado de vezes; esgotadas as tentativas, o lote inteiro vai para o destino de falha em vez de travar o shard — os registros seguintes continuam sendo processados.
- A retenção estendida é o que torna o reprocesso possível. Com o stream configurado para reter os dados por mais tempo que o padrão de 24 horas, um bug descoberto dias depois ainda encontra o evento original — sem essa janela, o dado já teria sido descartado e não haveria o que reler.
- Depois da correção, um consumidor temporário relê desde o início. Um novo event source mapping, com `starting_position = TRIM_HORIZON`, lê o shard do ponto mais antigo ainda retido. Ele escreve numa tabela separada — não toca o produtor nem o caminho em tempo real, que continuam publicando e alertando.
- Os alarmes dizem se o consumidor está acompanhando o pico. `IteratorAgeMilliseconds` alto significa que o consumidor está processando registros antigos e ficando para trás; throughput excedido no shard significa que o gargalo voltou a ser capacidade, não velocidade de leitura.
O ajuste com maior efeito, e ele cabe numa linha de código
Trocar a chave de partição de `evento.Tipo` para `evento.IdPedido` resolve os dois sintomas relatados pela Cadência ao mesmo tempo — distribuição de escrita e ordem por pedido — porque os dois tinham a MESMA causa. Não é coincidência: cardinalidade da chave é o parâmetro que decide as duas coisas simultaneamente.
O caminho de um evento, ponta a ponta
O nome de cada etapa não é jargão: é a prova de onde a ordem existe e onde ela para de existir. Dentro de um shard, ela é garantia de serviço; entre shards, ela não é garantida por absolutamente nada.
// O payload que o Checkout API publica a cada transicao de estado do pedido.
// IdPedido e a CHAVE DE PARTICAO — nao Tipo, que foi o desenho minimo.
{
"idPedido": "6f2a9e10-8b3c-4a71-9e2d-1c9f7a44b201",
"tipo": "pedido_cancelado",
"instante": "2026-08-08T14:03:12.418Z",
"valor": 349.90,
// Motivo so existe em eventos de cancelamento; campos que nao se aplicam ao
// tipo simplesmente nao aparecem, em vez de irem como null.
"motivo": "estorno_solicitado_pelo_cliente"
}Por que o motivo do cancelamento importa para o consumidor
É o dado que diferencia um cancelamento operacional (endereço inválido, item em falta) de um pedido de estorno depois do pagamento — que é o padrão que o alerta de fraude deste laboratório existe para pegar. Sem esse campo, toda a lógica do consumidor teria que inferir a partir de tempo decorrido, um sinal muito mais fraco.
As decisões, e o que se perde em cada uma
📋 Ingerir eventos de ciclo de vida do pedido (criado, pago, separado, cancelado, entregue) durante uma campanha com pico de até 3× o volume normal, com a exigência de que nenhum evento se perca em silêncio e que a ordem por pedido seja reconstituível.
O id do pedido tem a cardinalidade que o problema pede: distribui a escrita quase uniformemente E garante que os eventos de um mesmo pedido caiam no mesmo shard, que é a única forma de ter ordem entre eles sem um serviço de coordenação à parte. Modo provisionado importa aqui porque, sob chave enviesada, o modo on-demand reparte o shard igualmente sem isolar a chave quente — quem quer corrigir um shard específico precisa do controle granular que só o provisionado oferece.
Alt: Amazon MSK (Kafka gerenciado) — Resolveria o mesmo problema com partições em vez de shards, mas cobra operação de cluster — nós, brokers, storage — que este time de duas pessoas não tem capacidade de manter. Justificável quando já existe investimento em Kafka.
Alt: SQS FIFO por grupo de mensagem — Resolveria a ordem por pedido com `MessageGroupId`, mas SQS não é reprocessável por janela de tempo: uma vez consumida e apagada, a mensagem não existe mais. Perde exatamente a capacidade de reprocesso que este laboratório exige.
Alt: Kinesis em modo on-demand — Remove a gestão manual de shards, mas — com uma chave ainda enviesada — sofre a mesma saturação: o modo divide o tráfego uniformemente ao escalar, sem detectar nem isolar a chave quente. É mais simples de operar, não corrige o desenho.
Alt: DynamoDB Streams (L29) — Resolve captura de mudança de ESTADO no banco, não emissão de evento de APLICAÇÃO. Aqui não há necessariamente uma escrita no banco por trás de cada evento — "pedido separado" pode ser um evento só de fluxo, sem persistência.
| Decisão | Escolha | Alternativas | Motivo | O que se perde |
|---|---|---|---|---|
| Chave de partição | id do pedido | tipo do evento; chave aleatória por chamada; id do cliente | resolve distribuição de escrita e ordem por pedido com a MESMA linha de código | eventos do mesmo cliente em pedidos diferentes não ficam agrupados — raramente é o que se quer aqui |
| Modo do stream | provisionado, 8 shards | on-demand | permite SplitShard mirado numa chave específica se um hot shard reaparecer | exige acompanhar utilização e reprovisionar manualmente antes de picos maiores |
| Tratamento de falha do produtor | reenviar só os registros com ErrorCode | descartar o lote inteiro; reenviar tudo de novo; ignorar e seguir | não reprocessa quem já teve sucesso, e não perde quem falhou | mais código no produtor do que um `PutRecords` simples |
| Retenção do stream | 168h (7 dias) | 24h (padrão); 8.760h (365 dias, o teto) | cobre o horizonte real de detecção de bug de lógica sem pagar pelo teto | um bug descoberto no oitavo dia já não tem o dado original para reprocessar |
| Falha do consumidor | bissecção de lote + destino de falha em SQS | sem destino de falha (padrão); DLQ direto sem bisseção | isola o registro ruim sem travar o shard nem descartar o lote inteiro sem investigar | mensagens na fila de falha exigem um processo humano de triagem |
| Reprocesso | mapping temporário, separado, com destino próprio | reaproveitar o mesmo consumidor em tempo real | evita duplicar ação de negócio (alerta, estorno) por escrever no mesmo lugar duas vezes | mais um recurso temporário para lembrar de remover depois |
A dívida que este desenho ainda tem, e que não é resolvida aqui
O estado "último evento visto por pedido" no consumidor vive em memória, dentro da execução da Lambda. Ele não sobrevive a um ambiente reciclado, e duas invocações concorrentes do mesmo shard — o que o ParallelizationFactor permite — podem não compartilhar esse cache. Corrigir isso pede um estado externo persistido, e é uma decisão explicitamente fora do escopo, para não misturar "ordem de entrega" com "onde guardar estado de negócio".
Construir: o stream, com a lente do shard
O Terraform abaixo não tem nada de exótico — é a métrica por shard e a retenção estendida que fazem a diferença entre este desenho e o mínimo.
# kinesis.tf — o stream, com a lente do shard
resource "aws_kinesis_stream" "eventos_pedido" {
name = "${var.projeto}-eventos-pedido"
# Provisionado, nao on-demand: sob chave enviesada, o on-demand reparte o
# trafego igualmente ao escalar e NAO isola a chave quente. So o modo
# provisionado da o controle granular (SplitShard mirado) para corrigir UM
# shard especifico sem recriar o stream inteiro.
stream_mode_details {
stream_mode = "PROVISIONED"
}
shard_count = 8 # dimensionado com margem sobre o pico de 3x o volume normal
# Retencao alem do padrao de 24h: e o que HABILITA o reprocesso depois de
# corrigir um bug no consumidor. 168h = 7 dias cobre o horizonte real de
# deteccao de um bug de logica; o teto documentado e 8760h (365 dias).
retention_period = 168
encryption_type = "KMS"
kms_key_id = aws_kms_key.eventos.key_id
# Sem isto, CloudWatch so mostra metricas agregadas do STREAM. E a metrica
# POR SHARD que torna um shard quente visivel em vez de suspeita.
shard_level_metrics = [
"IncomingBytes",
"IncomingRecords",
"WriteProvisionedThroughputExceeded",
"ReadProvisionedThroughputExceeded",
"IteratorAgeMilliseconds",
]
}
# O produtor so PRECISA escrever. Ler e responsabilidade do consumidor, com
# um papel proprio — e e por isso que as duas policies abaixo sao distintas.
data "aws_iam_policy_document" "publicar_eventos" {
statement {
effect = "Allow"
actions = ["kinesis:PutRecords"]
resources = [aws_kinesis_stream.eventos_pedido.arn]
}
statement {
effect = "Allow"
actions = ["kms:GenerateDataKey", "kms:Decrypt"]
resources = [aws_kms_key.eventos.arn]
}
}
# Anexada ao MESMO task role da Checkout API do L01/L03 — nao criamos um
# papel novo para o publicador, porque a identidade que escreve continua
# sendo a mesma aplicacao.
resource "aws_iam_role_policy" "checkout_publica_eventos" {
name = "${var.projeto}-publica-eventos"
role = aws_iam_role.checkout_task_role.id
policy = data.aws_iam_policy_document.publicar_eventos.json
}
output "stream_eventos_pedido_arn" {
value = aws_kinesis_stream.eventos_pedido.arn
}
Por que provisionado, e não on-demand, neste laboratório especificamente
O modo on-demand escala automaticamente, mas divide o tráfego IGUALMENTE entre os shards novos ao crescer — sem detectar qual chave está concentrando a carga. Sob uma chave verdadeiramente enviesada, ele continua devolvendo exceção de throughput. O modo provisionado permite `SplitShard` mirado no shard específico que está quente, que é o controle que este laboratório quer demonstrar.
Construir: o produtor que lê a resposta antes de dar por certo
A mudança central do produtor não é a chave — é o que ele faz DEPOIS de chamar `PutRecords`. O código anterior da Cadência parava na primeira linha; este vai até o fim do corpo da resposta.
// PublicadorEventoPedido.cs — le a resposta de PutRecords antes de dar por certo
using Amazon.Kinesis;
using Amazon.Kinesis.Model;
public sealed class PublicadorEventoPedido
{
private readonly IAmazonKinesis _kinesis;
private readonly string _stream;
private readonly ILogger<PublicadorEventoPedido> _log;
public PublicadorEventoPedido(IAmazonKinesis kinesis, string stream,
ILogger<PublicadorEventoPedido> log)
{
_kinesis = kinesis;
_stream = stream;
_log = log;
}
// A CHAVE e o id do pedido — nao o tipo do evento. E a linha que resolve
// os dois sintomas do laboratorio: distribui a escrita (cardinalidade
// alta) e garante que os eventos de UM pedido caiam no MESMO shard
// (ordem). Trocar esta linha de volta para "evento.Tipo" reintroduz os
// dois problemas ao mesmo tempo, porque tem a MESMA causa.
private static string ChaveDeParticao(EventoPedido evento) => evento.IdPedido.ToString();
public async Task PublicarLoteAsync(IReadOnlyList<EventoPedido> eventos,
CancellationToken ct = default)
{
var pendentes = eventos
.Select(e => new PutRecordsRequestEntry
{
PartitionKey = ChaveDeParticao(e),
Data = new MemoryStream(JsonSerializer.SerializeToUtf8Bytes(e)),
})
.ToList();
// Ate 4 tentativas: a primeira e as chamadas originais, as demais so
// reenviam quem falhou. Backoff exponencial com jitter evita que o
// proprio retry sincronizado piore o throttling que o causou.
for (var tentativa = 1; tentativa <= 4 && pendentes.Count > 0; tentativa++)
{
var resposta = await _kinesis.PutRecordsAsync(new PutRecordsRequest
{
StreamName = _stream,
Records = pendentes,
}, ct);
// O PONTO CENTRAL do laboratorio: a chamada acima NAO lanca
// excecao por falha parcial. HTTP 200 e devolvido mesmo com
// FailedRecordCount > 0 — quem so confere o status conclui
// sucesso com registros que nunca entraram no stream.
if (resposta.FailedRecordCount == 0)
return;
var proximaRodada = new List<PutRecordsRequestEntry>();
for (var i = 0; i < resposta.Records.Count; i++)
{
var item = resposta.Records[i];
if (item.ErrorCode is null) continue; // este teve sucesso
_log.LogWarning(
"registro {Indice} falhou na tentativa {Tentativa}: {Codigo} — {Mensagem}",
i, tentativa, item.ErrorCode, item.ErrorMessage);
proximaRodada.Add(pendentes[i]);
}
pendentes = proximaRodada;
if (pendentes.Count > 0)
await Task.Delay(BackoffComJitter(tentativa), ct);
}
if (pendentes.Count > 0)
{
// Esgotadas as tentativas do PRODUTOR — nao do stream. Isto e
// perda real e precisa ser tratado como incidente, nunca como
// log silencioso: metrica customizada + alarme, nao so LogError.
_log.LogError(
"{Quantidade} eventos descartados apos {Tentativas} tentativas — dado perdido",
pendentes.Count, 4);
throw new EventosPedidoNaoPublicadosException(pendentes.Count);
}
}
private static TimeSpan BackoffComJitter(int tentativa)
{
var baseMs = Math.Pow(2, tentativa) * 100;
var jitter = Random.Shared.Next(0, 100);
return TimeSpan.FromMilliseconds(baseMs + jitter);
}
}
public sealed record EventoPedido(Guid IdPedido, string Tipo, DateTimeOffset Instante, decimal Valor);
public sealed class EventosPedidoNaoPublicadosException(int quantidade)
: Exception($"{quantidade} eventos de pedido nao foram publicados apos esgotar as tentativas");
Backoff sem jitter sincroniza os retries e piora o throttling
Se vários produtores throttlam no mesmo instante — o que é comum num pico — e todos esperam exatamente o mesmo tempo antes de tentar de novo, a segunda tentativa colide de novo, em massa. O jitter aleatório espalha as tentativas no tempo e é o que torna o backoff exponencial eficaz na prática, não só na teoria.
Construir: o consumidor e o event source mapping
O mapping é onde a garantia de ordem por chave, a paralelização e o destino de falha se encontram — os três Terraform diferentes, um objeto só.
# consumidor.tf — a Lambda em tempo real, o destino de falha, os alarmes
resource "aws_sqs_queue" "eventos_pedido_dlq" {
name = "${var.projeto}-eventos-pedido-dlq"
# 14 dias: maior que a retencao do stream, para dar tempo de investigar um
# lote que falhou sem que a fila em si vire novo ponto de perda.
message_retention_seconds = 1209600
kms_master_key_id = aws_kms_key.eventos.key_id
}
resource "aws_lambda_function" "alerta_fraude" {
function_name = "${var.projeto}-alerta-fraude"
role = aws_iam_role.lambda_alerta_fraude.arn
handler = "FraudeFunction::FraudeFunction.Function::FunctionHandler"
runtime = "dotnet8"
timeout = 30
memory_size = 256
environment {
variables = {
SNS_TOPIC_ARN = aws_sns_topic.alerta_risco.arn
}
}
}
resource "aws_lambda_event_source_mapping" "eventos_pedido_tempo_real" {
event_source_arn = aws_kinesis_stream.eventos_pedido.arn
function_name = aws_lambda_function.alerta_fraude.arn
starting_position = "LATEST" # so o consumo normal comeca do fim; o replay usa outro mapping
batch_size = 100
maximum_batching_window_in_seconds = 2
# Ate 10 lotes concorrentes POR SHARD. O Lambda continua garantindo ordem
# por CHAVE de particao mesmo com paralelizacao > 1 — a concorrencia e
# entre chaves diferentes do mesmo shard, nao dentro da mesma chave.
parallelization_factor = 2
# Um erro no lote e re-tentado ate 3 vezes, bissectando o lote pela metade
# a cada tentativa — isola o registro ruim em vez de re-tentar o lote
# inteiro do mesmo jeito repetidamente.
bisect_batch_on_function_error = true
maximum_retry_attempts = 3
maximum_record_age_in_seconds = 3600
destination_config {
on_failure {
# Esgotadas as tentativas, o lote vai para ca — e o shard SEGUE lendo
# os proximos registros. Sem isto, um registro ruim travaria o shard
# inteiro, e essa e a diferenca entre "um pedido com problema" e
# "todo o stream parado".
destination_arn = aws_sqs_queue.eventos_pedido_dlq.arn
}
}
}
resource "aws_sns_topic" "alerta_risco" {
name = "${var.projeto}-alerta-risco"
kms_master_key_id = aws_kms_key.eventos.key_id
}
data "aws_iam_policy_document" "lambda_le_stream" {
statement {
effect = "Allow"
actions = [
"kinesis:GetRecords",
"kinesis:GetShardIterator",
"kinesis:DescribeStreamSummary",
"kinesis:ListShards",
]
resources = [aws_kinesis_stream.eventos_pedido.arn]
}
statement {
effect = "Allow"
actions = ["kms:Decrypt"]
resources = [aws_kms_key.eventos.arn]
}
statement {
effect = "Allow"
actions = ["sqs:SendMessage"]
resources = [aws_sqs_queue.eventos_pedido_dlq.arn]
}
statement {
effect = "Allow"
actions = ["sns:Publish"]
resources = [aws_sns_topic.alerta_risco.arn]
}
}
resource "aws_iam_role_policy" "lambda_le_stream" {
name = "${var.projeto}-lambda-le-stream"
role = aws_iam_role.lambda_alerta_fraude.id
policy = data.aws_iam_policy_document.lambda_le_stream.json
}
# Os dois alarmes que decidem se o consumidor acompanha o pico, e nao a
# capacidade de escrita — que ja tem os seus proprios (nao mostrados aqui,
# ja cobertos pelo shard_level_metrics do kinesis.tf).
resource "aws_cloudwatch_metric_alarm" "iterator_age" {
alarm_name = "${var.projeto}-iterator-age-alta"
namespace = "AWS/Lambda"
metric_name = "IteratorAge"
statistic = "Maximum"
period = 60
evaluation_periods = 3
# 30s de atraso e tolerado; acima disso o consumidor esta perdendo o pico,
# nao so processando devagar.
threshold = 30000
comparison_operator = "GreaterThanThreshold"
treat_missing_data = "notBreaching"
dimensions = {
FunctionName = aws_lambda_function.alerta_fraude.function_name
}
alarm_actions = [aws_sns_topic.alerta_operacao.arn]
}
resource "aws_cloudwatch_metric_alarm" "throughput_excedido" {
alarm_name = "${var.projeto}-shard-throughput-excedido"
namespace = "AWS/Kinesis"
metric_name = "WriteProvisionedThroughputExceeded"
statistic = "Sum"
period = 60
evaluation_periods = 2
threshold = 0
comparison_operator = "GreaterThanThreshold"
treat_missing_data = "notBreaching"
dimensions = {
StreamName = aws_kinesis_stream.eventos_pedido.name
}
alarm_actions = [aws_sns_topic.alerta_operacao.arn]
}
| Parâmetro | Valor aqui | Por quê |
|---|---|---|
| `starting_position` | `LATEST` no mapping principal | consumo em tempo real começa do agora; o reprocesso usa um mapping separado |
| `parallelization_factor` | 2 | dobra o throughput de leitura por shard sem abrir mão de ordem por chave |
| `bisect_batch_on_function_error` | true | isola o registro ruim em vez de re-tentar o lote inteiro do mesmo jeito |
| `maximum_retry_attempts` | 3 | limite explícito; sem ele, o padrão tenta indefinidamente até expirar o registro |
| `destination_config.on_failure` | fila SQS | o shard segue lendo os próximos registros mesmo com um lote isolado como falho |
Construir: o consumidor em C#, e o limite do que ele sabe
O handler processa um lote de um shard, na ordem em que o shard entregou. É simples de propósito: o laboratório prova ORDEM DE ENTREGA, não desenho de máquina de estados de fraude — misturar os dois esconderia qual dos dois está sendo testado.
// FraudeFunction.cs — le em ordem por pedido, e admite o limite do que sabe
using Amazon.Lambda.Core;
using Amazon.Lambda.KinesisEvents;
using Amazon.SimpleNotificationService;
namespace FraudeFunction;
public class Function
{
private readonly IAmazonSimpleNotificationService _sns;
private readonly string _topicoArn = Environment.GetEnvironmentVariable("SNS_TOPIC_ARN")!;
// Cache em memoria, valido so ENQUANTO a execucao ficar quente. Nao e um
// substituto para estado persistido: se o ambiente reciclar entre dois
// eventos do mesmo pedido, o ultimo estado visto se perde. Para producao
// de verdade isso pede uma tabela (DynamoDB), e essa decisao fica FORA
// do escopo deste laboratorio, que e sobre ordem de ENTREGA, nao sobre
// onde guardar estado de negocio — misturar os dois esconderia qual dos
// dois esta sendo provado.
private static readonly Dictionary<Guid, string> UltimoEventoVisto = new();
public Function(IAmazonSimpleNotificationService sns) => _sns = sns;
public async Task FunctionHandler(KinesisEvent evento, ILambdaContext contexto)
{
// Os registros deste lote pertencem TODOS ao mesmo shard, e chegam
// na ordem em que o shard os recebeu — e ela so é confiável porque
// a chave de particao é o id do pedido: todo evento de um pedido
// sempre cai neste MESMO shard.
foreach (var registro in evento.Records)
{
var dado = JsonSerializer.Deserialize<EventoPedidoLido>(
registro.Kinesis.Data);
if (dado is null) continue;
if (UltimoEventoVisto.TryGetValue(dado.IdPedido, out var anterior)
&& anterior == "pedido_cancelado" && dado.Tipo == "pedido_pago")
{
// Deveria ser logicamente impossivel com a ordem preservada:
// pago depois de cancelado. Se isto disparar, o problema nao
// e o consumidor — e a garantia de shard sendo violada em
// algum ponto anterior (ex.: reprocesso escrevendo no
// caminho de tempo real por engano).
contexto.Logger.LogWarning(
$"ordem suspeita no pedido {dado.IdPedido}: pago apos cancelado");
}
if (anterior == "pedido_pago" && dado.Tipo == "pedido_cancelado")
{
await _sns.PublishAsync(_topicoArn,
$"pedido {dado.IdPedido}: cancelado {SegundosDesde(anterior)}s apos o pagamento — revisar estorno");
}
UltimoEventoVisto[dado.IdPedido] = dado.Tipo;
}
}
private static long SegundosDesde(string _) => 0; // ilustrativo: o calculo real usa o instante do evento anterior
}
public sealed record EventoPedidoLido(Guid IdPedido, string Tipo, DateTimeOffset Instante, decimal Valor);
Nunca aponte o mapping de reprocesso para o mesmo destino do tempo real
Se o consumidor de replay publicar no MESMO tópico SNS que o consumidor em tempo real, cada evento reprocessado dispara um segundo alerta — e se a ação do outro lado do alerta for automática (bloquear conta, estornar valor), ela executa duas vezes. É dinheiro se movendo por engano, não só ruído de log. O destino do replay tem de ser sempre distinto, e a reconciliação entre os dois, um passo manual deliberado.
Implantar, e provar com número — não com "chegou tudo"
Cinco provas. Nenhuma aceita "parece que funcionou": cada uma tem um comando e um resultado que aprova ou reprova.
# provas.sh — cinco medicoes; nenhuma conclusao vem de "parece que chegou"
PROJETO=ffv-lab-kinesis
STREAM="${PROJETO}-eventos-pedido"
# ── Prova 1: o hot shard, com numero (ver medir-shard-quente.sh acima) ───────
# ── Prova 2: PutRecords tem falha parcial mesmo devolvendo 200 ───────────────
# Gera trafego concentrado numa unica chave, o suficiente para passar de
# 1.000 registros/s NUM shard, e confere o corpo da resposta.
python3 - <<'PY'
import boto3, json, time
k = boto3.client("kinesis")
lote = [{"Data": json.dumps({"i": i}).encode(), "PartitionKey": "pedido_criado"}
for i in range(500)]
falhas_totais = 0
for _ in range(6): # 6 x 500 = 3.000 registros contra um shard de 1.000/s
r = k.put_records(StreamName="ffv-lab-kinesis-eventos-pedido", Records=lote)
falhas_totais += r["FailedRecordCount"]
print("HTTP 200, FailedRecordCount =", r["FailedRecordCount"])
print("total de falhas parciais:", falhas_totais)
PY
# Esperado: pelo menos uma das seis chamadas mostra FailedRecordCount > 0,
# e TODAS retornam sem lancar excecao — a falha mora no corpo, nao no status.
# ── Prova 3: os eventos do MESMO pedido chegam no MESMO shard ────────────────
aws kinesis list-shards --stream-name "$STREAM" \
--query 'Shards[].{Id:ShardId,Faixa:HashKeyRange}' --output table
# Compare o ShardId reportado nos logs da Lambda para dois eventos do mesmo
# IdPedido: com chave = id do pedido, tem de ser o MESMO ShardId sempre.
# ── Prova 4: a retenção realmente guarda o dado além de 24h ──────────────────
aws kinesis describe-stream-summary --stream-name "$STREAM" \
--query 'StreamDescriptionSummary.RetentionPeriodHours' --output text
# Esperado: 168 (7 dias). Se voltar 24, a mudança de retenção não foi aplicada
# e o reprocesso da prova 5 vai falhar por falta de dado.
# ── Prova 5: reprocesso a partir de TRIM_HORIZON encontra o evento antigo ────
SHARD=$(aws kinesis list-shards --stream-name "$STREAM" \
--query 'Shards[0].ShardId' --output text)
ITERADOR=$(aws kinesis get-shard-iterator --stream-name "$STREAM" \
--shard-id "$SHARD" --shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' --output text)
aws kinesis get-records --shard-iterator "$ITERADOR" \
--query 'Records[0].{Sequencia:SequenceNumber,Chave:PartitionKey}' --output table
# Esperado: um registro com SequenceNumber baixo (o mais antigo ainda retido),
# nao o mais recente — confirma que TRIM_HORIZON aponta para o inicio, e nao
# para "agora", que seria o comportamento de LATEST.
| Prova | O que confirma | O que reprova, e o que significa |
|---|---|---|
| 1 · Hot shard mensurado | um shard perto do teto, os outros ociosos | todos uniformes já indicaria chave corrigida — rode antes da correção |
| 2 · Falha parcial em PutRecords | FailedRecordCount > 0 em pelo menos uma chamada | zero falhas pode indicar volume insuficiente para saturar o shard de teste |
| 3 · Mesmo pedido, mesmo shard | o ShardId é idêntico para dois eventos do mesmo IdPedido | ShardId diferente indica que a chave de partição não é realmente o id do pedido |
| 4 · Retenção aplicada | `RetentionPeriodHours` retorna 168 | retornar 24 significa que a mudança de retenção não foi aplicada ao stream |
| 5 · Reprocesso encontra o dado antigo | SequenceNumber baixo, não o mais recente | erro "trim horizon" indicaria dado já fora da janela de retenção |
Reprocessar: o caminho que a retenção paga
Esta é a prova final do requisito central: corrigir um bug no consumidor sem tocar o produtor, e sem perder o que já tinha acontecido.
#!/usr/bin/env bash
# reprocessar.sh — depois de corrigir o bug, relê o stream inteiro sem tocar o produtor
set -euo pipefail
PROJETO="${PROJETO:?defina PROJETO}"
STREAM="${PROJETO}-eventos-pedido"
FUNCAO_REPLAY="${PROJETO}-reprocessamento"
# 1. Cria um mapping TEMPORARIO, separado do consumo em tempo real, que ja
# continua rodando com starting_position=LATEST sem interrupção.
UUID=$(aws lambda create-event-source-mapping \
--function-name "$FUNCAO_REPLAY" \
--event-source-arn "arn:aws:kinesis:us-east-1:$(aws sts get-caller-identity --query Account --output text):stream/$STREAM" \
--starting-position TRIM_HORIZON \
--batch-size 200 \
--query 'UUID' --output text)
echo "mapping de reprocesso criado: $UUID"
# 2. Acompanha o IteratorAge caindo: quando chegar perto de zero, o replay
# alcançou o tempo real e pode ser desligado.
echo "acompanhe com:"
echo " aws lambda get-event-source-mapping --uuid $UUID --query 'LastProcessingResult'"
echo " aws cloudwatch get-metric-statistics --namespace AWS/Lambda --metric-name IteratorAge \"
echo " --dimensions Name=FunctionName,Value=$FUNCAO_REPLAY --statistics Maximum --period 60 \"
echo " --start-time \$(date -u -d '10 minutes ago' +%Y-%m-%dT%H:%M:%S) --end-time \$(date -u +%Y-%m-%dT%H:%M:%S)"
read -rp "IteratorAge proximo de zero? apagar o mapping de reprocesso [s/N]: " ok
if [[ "$ok" == "s" ]]; then
aws lambda delete-event-source-mapping --uuid "$UUID"
echo "mapping de reprocesso removido — o replay foi temporario, como deveria ser"
fi
Reprocessar fora da janela de retenção não tem correção possível
Se o bug for descoberto no dia oito com retenção de sete, o dado do primeiro dia já foi descartado pelo Kinesis — de forma definitiva e irreversível. Não existe comando nem tipo de iterador que recupera isso; `TRIM_HORIZON` aponta para o mais antigo AINDA retido, não para o início absoluto do stream. A única prevenção é configurar a retenção pela cadência real de detecção de bugs, antes de precisar dela.
Um serviço publica 500 registros com PutRecords, recebe HTTP 200 e considera o lote publicado. Duas horas depois, descobre que 40 registros nunca chegaram ao stream. Qual é a explicação mais provável?
Quebrar de propósito: três falhas e o diagnóstico
| Falha | Como provocar | Sintoma | Onde olhar | Correção |
|---|---|---|---|---|
| Chave de baixa cardinalidade | publique com chave = tipo_evento durante um teste de carga | um ShardId concentra quase todo IncomingBytes; os outros ficam perto de zero | CloudWatch, métrica por shard (`shard_level_metrics`) | trocar a chave para algo de alta cardinalidade e que agrupe o que precisa de ordem |
| Produtor ignora FailedRecordCount | remova a checagem do campo no código e sature um shard de propósito | logs do produtor mostram sucesso; o consumidor nunca recebe parte dos eventos | comparar contagem publicada (log do produtor) com contagem consumida (métrica da Lambda) | reintroduzir a checagem de `FailedRecordCount` com reenvio dos registros falhos |
| Mapping sem destino de falha | force uma exceção não tratada no handler e desabilite `destination_config` | o shard para de avançar; `IteratorAge` cresce sem limite para aquele shard | `get-event-source-mapping` → `LastProcessingResult` | habilitar destino de falha e `bisect_batch_on_function_error` |
A falha que demora mais para ser notada
Um shard travado por registro ruim não gera erro visível em nenhum painel padrão — ele simplesmente para de avançar, silenciosamente, enquanto os outros shards seguem normais. Sem um alarme específico em `IteratorAge` POR FUNÇÃO, o sintoma só aparece quando alguém nota que um subconjunto de pedidos parou de gerar alerta.
Segurança: o que passa a trafegar por um caminho novo
Eventos de pedido carregam valor monetário e, potencialmente, dado que identifica o cliente. Um caminho de streaming novo é um caminho novo de exposição, mesmo sem mudar o banco de dados.
| Risco | Probabilidade | Impacto | Prevenção | Detecção | Resposta |
|---|---|---|---|---|---|
| Dado de pedido exposto na fila de falha | média | alto | criptografia KMS na fila; retenção da fila só o tempo de investigação | CloudTrail em `ReceiveMessage` fora do papel de investigação esperado | purgar mensagens após triagem; girar a chave KMS se houve exposição |
| Produtor com permissão além de escrita | baixa | médio | política do produtor restrita a `kinesis:PutRecords` no ARN do stream | IAM Access Analyzer sobre uso real do papel | remover a ação não utilizada da política |
| Reprocesso duplicando ação de negócio | média | alto | destino do mapping de replay sempre distinto do caminho em tempo real | comparar contagem de registros processados nos dois destinos | pausar o mapping de replay; reconciliar manualmente antes de reativar |
| Chave de partição com dado sensível | média | médio | usar identificador interno (id do pedido), nunca e-mail ou documento, como chave | Macie sobre o payload armazenado no stream | trocar a chave; os registros antigos ainda carregam o dado até expirar pela retenção |
| Alarme de throughput ignorado | alta | médio | alarme com ação em SNS de operação, não só visível no console | ausência de confirmação no canal de plantão dentro do prazo | escalonamento automático após N minutos sem resposta |
O papel de execução da Lambda de reprocesso não é o mesmo do tempo real
Mesmo com política igualmente restrita, manter os dois papéis SEPARADOS torna possível revogar o acesso do reprocesso assim que o replay termina, sem afetar o consumo normal — que nunca deveria ficar sem permissão, nem por um segundo.
Observabilidade: as perguntas que o painel de ingestão responde
| Pergunta | Métrica ou consulta | O que significa mudar | Limiar inicial |
|---|---|---|---|
| O consumidor está acompanhando o pico? | `IteratorAge` (Lambda) | crescendo, é atraso acumulando — o consumidor está processando dado cada vez mais velho | > 30.000 ms por 3 períodos de 1 min |
| Algum shard está saturado na escrita? | `WriteProvisionedThroughputExceeded` por `ShardId` | qualquer valor acima de zero é throttling acontecendo agora | > 0, qualquer período |
| A leitura está no limite do shard? | `ReadProvisionedThroughputExceeded` | mais de 5 `GetRecords`/s ou 2 MB/s sustentado no mesmo shard | tendência de subida, não só valor absoluto |
| Quantos registros o produtor está perdendo? | métrica customizada a partir de `FailedRecordCount` | aumento sustentado indica throttling de escrita persistente, não picos isolados | > 1% do volume publicado em 5 min |
| O destino de falha está recebendo mensagem? | `ApproximateNumberOfMessagesVisible` (SQS) | qualquer valor sustentado indica lotes esgotando as tentativas do mapping | > 0 por mais de 10 min |
| A distribuição entre shards está uniforme? | `IncomingBytes` por `ShardId`, comparado entre shards | desvio grande entre o mais e o menos carregado é sinal de chave enviesada de novo | desvio > 30% entre o shard mais e o menos carregado |
| O reprocesso está avançando? | `IteratorAge` da função de replay | caindo é o replay se aproximando do tempo real | acompanhar até ficar próximo de zero, então desligar o mapping |
A métrica que engana quando só um consumidor é observado
Um stream pode ter mais de um consumidor lendo o mesmo shard — o tempo real e, temporariamente, o replay. `IteratorAge` é por FUNÇÃO, não por stream: olhar só a métrica do consumidor em tempo real esconde um replay que travou silenciosamente.
Escala: 10, 10 mil, 1 milhão de eventos por dia
| Volume | O que acontece | O que passa a doer | O que fazer |
|---|---|---|---|
| 10 eventos/dia | um shard sobra capacidade o tempo todo | nada — é desperdício de shard-hora, não risco | considerar on-demand aqui especificamente, já que não há chave enviesada em jogo com tão pouco tráfego |
| 10 mil eventos/dia (~0,1/s de média) | um shard ainda cobre a média | rajadas de campanha podem saturar por minutos, mesmo com média baixa | dimensionar pelo PICO, não pela média — é o erro mais comum nesta faixa |
| 1 milhão eventos/dia (~12/s de média, picos bem maiores) | múltiplos shards se tornam obrigatórios | chave de alta cardinalidade deixa de ser boa prática e vira pré-condição | chave = id do pedido, 8+ shards, alarmes de throughput ativos desde o início |
| Pico de campanha (3× o normal, em rajada) | mesmo com chave boa, o throughput agregado pode faltar | shard_count dimensionado para o normal não cobre o pico | `UpdateShardCount` ANTES do evento conhecido, nunca durante |
| Falha de AZ | o Kinesis replica entre AZs automaticamente — o stream em si não cai | quem quebra é a aplicação produtora, se ela não for multi-AZ | é o plano de endereçamento e disponibilidade do L01/L02 cobrando de novo, não algo novo deste laboratório |
O gargalo que só aparece quando alguém dimensiona pela média
Um volume de "1 milhão de eventos por dia" soa como 12 eventos por segundo, dentro da folga de um shard só. Mas eventos de pedido não chegam uniformemente: eles se concentram no horário comercial e explodem em campanha. Dimensionar shard_count pela média diária é o erro que a Semana Cadência expôs no primeiro dia de pico.
Custo: o que este laboratório acrescenta à fatura
O modo provisionado cobra pelo que você RESERVA, não só pelo que usa — é a diferença central em relação ao on-demand, e o motivo de dimensionar com cuidado.
| Cenário | Volume | O que acrescenta | Tendência | Otimização |
|---|---|---|---|---|
| Protótipo | 1 shard, tráfego de teste | shard-hora e unidade de PUT payload; retenção padrão de 24h | desprezível | nenhuma; otimizar aqui gasta atenção onde não há dinheiro em jogo |
| Produção pequena | 8 shards, campanha ocasional | shard-hora × 8, retenção estendida a 7 dias, invocações Lambda proporcionais ao volume | previsível e ligada ao número de shards, não ao tráfego real dentro deles | ajustar shard_count pelo pico real medido, não por margem arbitrária |
| Alta escala | 20+ shards, múltiplas campanhas por mês | shard-hora relevante na fatura; considerar enhanced fan-out se houver múltiplos consumidores | cresce com shard_count permanente, mesmo em dias sem campanha | reduzir shard_count fora de campanha, se o processo operacional suportar a mudança |
| Dimensão | Cobra por | Cuidado |
|---|---|---|
| Shard provisionado | shard-hora, reservado independente de uso | shard ocioso fora de campanha continua cobrando — não é elástico como on-demand |
| Unidade de PUT payload | a cada 25 KB de dado publicado, arredondado para cima | muitos registros pequenos custam mais que poucos registros grandes de mesmo total de bytes |
| Retenção estendida | shard-hora adicional acima de 24h retidas | reter 365 dias "por segurança" sem uso real é o desperdício mais comum aqui |
| Enhanced fan-out (se usado) | consumer-shard-hora e dado consumido | só se justifica com múltiplos consumidores competindo pelos mesmos 2 MB/s compartilhados |
| Lambda (tempo real + replay) | invocações e duração | o replay soma ao mesmo medidor enquanto está ativo — mais um motivo para ser temporário |
| SQS e SNS | requisições e notificações | valor baixo por unidade, mas cresce com volume de falha — que deveria ser exceção, não regra |
O ganho de custo que não aparece na fatura de infraestrutura
A diferença entre descobrir a fraude no minuto do cancelamento e descobrir na conciliação do mês seguinte é o valor do estorno em si — não uma linha de AWS. O custo que este laboratório ataca é o de decisão tardia, e ele nunca aparece no Cost Explorer.
Well-Architected nos seis pilares
| Pilar | Situação ao fim deste laboratório | Risco que fica | Melhoria | Prioridade |
|---|---|---|---|---|
| Excelência operacional | reprocesso roteirizado e reversível, com prova de que alcançou o tempo real | triagem da fila de falha ainda é manual | automatizar reclassificação de erro conhecido versus erro novo | média |
| Segurança | chave de partição sem PII, criptografia KMS ponta a ponta, papéis separados | nenhum controle de schema impede um produtor futuro de incluir dado sensível na chave | validação de schema na borda de publicação | média |
| Confiabilidade | falha parcial tratada no produtor; falha do consumidor isolada, não bloqueante | estado de negócio do consumidor ainda em memória, não persistido | mover o cache de "último evento visto" para um armazenamento externo | alta |
| Eficiência de performance | chave de alta cardinalidade distribui quase uniformemente; ParallelizationFactor ajustado | shard_count fixo não se adapta sozinho a picos maiores que o previsto | automação de `UpdateShardCount` antes de campanhas conhecidas | média |
| Otimização de custos | retenção dimensionada pela cadência real de correção de bug, não pelo teto | shard_count permanente cobre o pico mesmo em dias normais | reduzir shards fora de campanha, se o processo operacional suportar | baixa |
| Sustentabilidade | shard ocioso fora de pico é capacidade não usada, não trabalho desperdiçado | nenhum específico deste laboratório | acompanha a otimização de custos acima | baixa |
Evolução em níveis: o que muda, e o que passa a doer
A terceira arquitetura não é um desenho: é a resposta a QUANDO trocar de desenho. Cada nível resolve um risco e compra outro.
Chave de partição por tipo de evento, produtor sem checar `FailedRecordCount`, um shard só. É onde a Cadência estava no primeiro dia de campanha.Chave = id do pedido, retry no produtor a partir de `FailedRecordCount`, destino de falha no consumidor, retenção estendida com reprocesso comprovado.Enhanced fan-out para consumidores que competiriam pelos mesmos 2 MB/s compartilhados; schema versionado no payload.O "último evento visto por pedido" sai da memória da Lambda e vai para um armazenamento externo, com escrita idempotente por `SequenceNumber`.Registro de schema centralizado, cotas de shard geridas por equipe, streams federados entre times com política de recurso.O histórico de eventos retido vira DADOS de treino: em vez de uma regra fixa ("cancelado logo após pago"), um MODELO aprende sequências de transição que precedem fraude confirmada, com IA graduando o risco em vez de só sinalizar sim/não.A ordem não é negociável, e o motivo é concreto
Um modelo de detecção no nível 6 depende de sequência de eventos CONFIÁVEL — que só existe a partir do nível 2, quando a chave de partição passou a garantir ordem por pedido. Treinar um modelo sobre dados do nível 1, onde a ordem pode estar trocada, ensinaria o modelo a aprender o BUG, não o padrão de fraude.
Onde IA entra nesta arquitetura, e onde não entra
Neste módulo, IA não resolve o problema central, e forçá-la seria o antipadrão que a própria série critica. "Como garantir que nenhum evento se perca e que a ordem por pedido exista" tem resposta determinística: chave de alta cardinalidade, checagem de falha parcial, retenção dimensionada. Um modelo não melhora nenhuma das três.
Há um lugar onde IA acrescentaria valor real, e é o Nível 6 da evolução: dado o histórico de transições de estado — agora confiavelmente ORDENADO por pedido, graças ao que este laboratório construiu — um classificador poderia graduar o risco de um padrão de eventos em vez de depender só de uma regra fixa como "cancelado logo após pago".
| Pergunta | Resposta honesta para este módulo |
|---|---|
| Qual problema a IA resolveria? | reconhecer padrões de sequência que precedem fraude e que uma regra fixa não cobre |
| Por que uma regra não bastaria, eventualmente? | porque padrões de fraude mudam; uma regra fixa só pega o que já foi catalogado |
| De onde viriam os dados? | do próprio stream retido e do destino corrigido do reprocesso — nenhuma coleta nova é necessária |
| Qual o risco? | aprender de poucos casos confirmados e marcar pedido legítimo como suspeito; exige avaliação com dado retido e caminho de revisão humana |
| Por que não agora? | porque a ordem confiável por pedido é PRÉ-CONDIÇÃO, e ela só passou a existir neste laboratório; não há histórico limpo suficiente ainda |
O uso de IA que parece atraente e é armadilha aqui
Pedir a um modelo para "decidir se o lote de PutRecords teve sucesso" troca um sinal determinístico — `FailedRecordCount`, que é exato — por um probabilístico, sem nenhum ganho: o campo já responde a pergunta com certeza. IA sobre ingestão faz sentido quando o sinal é ambíguo; aqui ele não é.
Anti-padrões deste laboratório
| Anti-padrão | Por que alguém faz | Por que é problema | Sintoma em produção | Forma correta | Quando é aceitável |
|---|---|---|---|---|---|
| Chave de partição = tipo do evento | parece organização: agrupar eventos parecidos no mesmo lugar | baixa cardinalidade concentra o tráfego real num shard e espalha o mesmo pedido entre shards | um ShardId concentra quase todo o IncomingBytes; ordem entre estágios de um pedido some | chave de alta cardinalidade que também mantenha junto o que precisa de ordem | nunca em produção; aceitável só em prova de conceito de baixíssimo volume |
| Ignorar FailedRecordCount | o caminho feliz do SDK deixa isso fácil — `PutRecordsAsync` não lança exceção por falha parcial | HTTP 200 vira sinônimo de sucesso total, e parte do lote nunca chega ao stream | contagem publicada no produtor maior que a contagem consumida, sem erro visível | inspecionar `FailedRecordCount` e reenviar só os registros com `ErrorCode` | nunca; é o ponto cego mais caro da API |
| Aumentar shard_count achando que resolve hot shard | é o menor número de linhas de Terraform que parece mexer no problema | resharding divide o INTERVALO de hash; uma chave constante cai sempre no mesmo lado da divisão | o número de shards sobe, o hot shard continua exatamente igual | diagnosticar a causa por métrica de shard ANTES de reprovisionar; corrigir a chave | só depois de confirmar que a chave já é de alta cardinalidade e o volume real cresceu |
| Consumidor sem destino de falha | parece que "sempre vai dar certo", e configurar destino de falha soa como pessimismo | um registro ruim trava o shard inteiro em retry até expirar, sem isolar o problema | `IteratorAge` de um shard específico cresce sem limite enquanto os outros seguem normais | `bisect_batch_on_function_error` + destino de falha configurado desde o início | nunca em produção; aceitável em teste local onde parar tudo é o comportamento desejado |
| Retenção no padrão de 24h em ambiente que precisa reprocessar | ninguém pensa em retenção até o dia em que precisa dela | um bug descoberto no segundo dia já não tem mais o dado do primeiro para reler | reprocesso falha com erro relacionado à janela de retenção | dimensionar a retenção pela cadência real de detecção de bug, com margem | só quando o consumidor é comprovadamente idempotente à distância de segundos, nunca de dias |
| Reprocessamento escrevendo no mesmo destino do tempo real | parece reaproveitar código e economizar um recurso | duplica ação de negócio disparada pelo alerta — cada evento reprocessado dispara de novo | contagem de alertas maior que a contagem de eventos únicos processados | destino do replay sempre separado, com reconciliação manual antes de mesclar | nunca; é o antipadrão com maior potencial de dano financeiro deste módulo |
Quando algo não funciona
| Sintoma | Causa provável | Como investigar | Onde olhar | Correção |
|---|---|---|---|---|
| Métrica de um shard em 100% e throughput total baixo | hot shard por chave de baixa cardinalidade | compare `IncomingBytes` entre todos os `ShardId` do stream | CloudWatch, `shard_level_metrics` | trocar a chave de partição por algo de alta cardinalidade |
| Produtor loga sucesso, evento nunca chega ao consumidor | `FailedRecordCount` ignorado no código do produtor | inspecione o corpo completo da resposta de `PutRecords`, não só o status | código do produtor | checar `FailedRecordCount` e reenviar os registros com `ErrorCode` |
| Fraude vê "cancelado" antes de "pago" para o mesmo pedido | chave de baixa cardinalidade espalhando estágios do mesmo pedido em shards diferentes | compare o `ShardId` de dois eventos do mesmo `IdPedido` nos logs | logs do consumidor Lambda | chave de partição = id do pedido |
| `IteratorAge` crescendo sem parar | consumidor mais lento que a taxa de chegada | CloudWatch `IteratorAgeMilliseconds` por função | métricas da Lambda | aumentar `ParallelizationFactor` ou otimizar a função |
| `GetRecords` lança `ProvisionedThroughputExceededException` | mais de 5 chamadas/s ou 2 MB/s de leitura no mesmo shard | métrica de leitura por shard, correlacionada com número de consumidores | `ReadProvisionedThroughputExceeded` | considerar enhanced fan-out se houver múltiplos consumidores competindo |
| Reprocesso não encontra o evento esperado | evento fora da janela de retenção | `describe-stream-summary` → `RetentionPeriodHours`, comparado com a idade do evento | configuração do stream | aumentar a retenção ANTES de precisar; não há recuperação retroativa |
| Mapping de replay parece nunca terminar | `IteratorAge` da função de replay não estava sendo acompanhado | consulte a métrica antes de assumir que travou | CloudWatch, função de reprocessamento | acompanhar até `IteratorAge` ficar próximo de zero, então desligar o mapping |
| `GetShardIterator` falha com iterador expirado | mais de 5 minutos entre a emissão do iterador e o uso em `GetRecords` | confira o intervalo entre as duas chamadas nos logs | código do consumidor manual (fora do event source mapping gerenciado) | solicitar um novo iterador; o event source mapping da Lambda já lida com isso sozinho |
A pergunta que resolve metade destes casos
Antes de mexer em qualquer parâmetro, pergunte: o problema é de ESCRITA (o produtor não consegue publicar) ou de LEITURA (o consumidor não consegue acompanhar)? As duas famílias de métrica — `Write*` e `Read*`/`IteratorAge` — apontam para lados opostos do pipeline, e confundir uma pela outra faz alguém ajustar o lado que já funcionava.
Limpeza: o que o destroy não leva
Este laboratório cria menos recurso que parece, mas dois deles sobrevivem ao terraform destroy se você esqueceu de um passo antes.
#!/usr/bin/env bash
# limpar.sh — o que o destroy não leva
set -euo pipefail
PROJETO=ffv-lab-kinesis
# 1. Derrube o que o Terraform administra.
terraform destroy -auto-approve
# 2. Se houve um mapping de reprocesso e ele nao foi apagado no fim do
# replay, ele continua lendo e pagando invocacao de Lambda para sempre.
aws lambda list-event-source-mappings --function-name "${PROJETO}-reprocessamento" \
--query 'EventSourceMappings[].UUID' --output text | \
xargs -n1 -I{} aws lambda delete-event-source-mapping --uuid {} 2>/dev/null || true
# 3. FILA DE FALHA: guarda ate 14 dias de mensagens mesmo depois do stream
# sumir. Se houver algo pendente de investigar, salve antes de apagar.
aws sqs receive-message --queue-url "$(aws sqs get-queue-url --queue-name ${PROJETO}-eventos-pedido-dlq --query QueueUrl --output text)" \
--max-number-of-messages 10 --query 'Messages[].Body' --output json 2>/dev/null || true
# 4. GRUPOS DE LOGS da Lambda: tem ciclo proprio, nao cai com a funcao.
aws logs delete-log-group --log-group-name "/aws/lambda/${PROJETO}-alerta-fraude" 2>/dev/null || true
aws logs delete-log-group --log-group-name "/aws/lambda/${PROJETO}-reprocessamento" 2>/dev/null || true
# 5. Prova final: nada com a tag do projeto de pe.
aws resourcegroupstaggingapi get-resources \
--tag-filters Key=Projeto,Values=${PROJETO} \
--query "ResourceTagMappingList[].ResourceARN" --output table
| Recurso | Sai no destroy? | Cobra parado? | Por que fica |
|---|---|---|---|
| Stream Kinesis | sim | sim, por shard-hora até ser destruído | nada de especial — mas confirme que não há mapping ativo apontando para ele |
| Mapping de reprocesso | não, se não foi apagado ao fim do replay | sim, por invocação | não pertence ao ciclo de vida do Terraform se foi criado via script/console à parte |
| Mensagens na fila de falha | a fila sim; o conteúdo não é recuperável depois | sim, até expirar | guarde o que precisa investigar antes de destruir — 14 dias de retenção não são infinitos |
| Grupos de logs da Lambda | depende de `skip_destroy` | sim, por retenção | ciclo próprio, sobrevive à função que os alimentava |
| Tabela de resultado do replay | sim, se em Terraform | sim, por armazenamento | confirme que o resultado já foi consumido pelo processo de negócio antes de apagar |
Resumo: problema, peça e motivo
| Problema | Peça | Por que ela, e não outra |
|---|---|---|
| Evento perdido em silêncio no pico | produtor lê `FailedRecordCount` e reenvia só os que falharam | HTTP 200 não é prova de gravação; a falha mora no corpo da resposta |
| Um shard satura enquanto outros ficam ociosos | chave de partição de alta cardinalidade | a mesma chave sempre cai no mesmo shard — cardinalidade baixa concentra por definição |
| Estágios do mesmo pedido em shards diferentes | chave = id do pedido | ordem só existe DENTRO de um shard; a chave decide quem compartilha shard |
| Registro ruim travando o shard inteiro | bissecção de lote + destino de falha em SQS | isola o problema sem parar os registros seguintes nem descartar sem investigar |
| Bug de consumidor descoberto dias depois | retenção estendida + reprocesso a partir de TRIM_HORIZON/AT_TIMESTAMP | retenção é o que HABILITA reler; sem ela o dado já se foi de forma irreversível |
| Reprocesso duplicando ação de negócio | destino do replay sempre separado do tempo real | evita disparar o mesmo alerta ou estorno duas vezes |
- O checkout serializa o evento e escolhe a chave de partição — id do pedido, não tipo.
- PutRecords grava até 500 registros por chamada, cada um de forma independente.
- O Kinesis aplica MD5 à chave e mapeia o resultado para o intervalo de hash de um shard.
- A mesma chave cai sempre no mesmo shard — é isso que garante ordem por pedido.
- O produtor confere FailedRecordCount e reenvia só os registros que falharam.
- A Lambda lê por shard, em lotes, na ordem em que o shard recebeu.
- Um lote com sucesso publica o alerta; um lote que falha repetidamente vai à fila de falha.
- Se um bug é corrigido depois, um mapping temporário relê desde TRIM_HORIZON ou AT_TIMESTAMP.
- O replay escreve num destino separado, e é desligado assim que alcança o tempo real.
Perguntas frequentes
❓ Por que o Kinesis PutRecords não lança exceção quando um registro falha?
❓ Adicionar mais shards resolve um hot shard?
❓ Qual a diferença entre TRIM_HORIZON, LATEST e AT_TIMESTAMP?
❓ Por que separar o consumidor de reprocesso do consumidor em tempo real?
❓ O que acontece com um registro que falha repetidamente no consumidor Lambda?
❓ Quanto tempo um evento fica disponível para reprocesso?
❓ O ParallelizationFactor da Lambda quebra a ordem de processamento?
❓ Vale a pena usar Kinesis para um volume baixo, tipo 10 eventos por dia?
Fixando
Um stream Kinesis tem chave de partição = tipo_evento, com "pedido_criado" e "pedido_cancelado" mapeados para shards diferentes. Um mesmo pedido gera um evento de cada tipo. O que garante a ordem entre esses dois eventos, do ponto de vista do consumidor?
Um time corrige um bug no consumidor seis dias depois de ele começar a processar errado. O stream tem retenção de 7 dias. Qual configuração de start position permite reprocessar exatamente o que foi afetado, sem tocar no produtor?
Conhecimentos, próximo módulo e documentação
| Item | Conteúdo |
|---|---|
| Conhecimentos anteriores necessários | Checkout API no ar (L01/L03), IAM básico, CloudWatch. Contraste útil (não pré-requisito): L29 |
| Conhecimentos adquiridos | capacidade fixa por shard e o hash MD5 que mapeia chave a shard; falha parcial de PutRecords e o campo FailedRecordCount; ordem garantida só dentro do shard; retenção como pré-condição de reprocesso; os três tipos de shard iterator |
| Limitação que fica | o estado "último evento visto por pedido" no consumidor vive em memória, não sobrevive a reciclagem de ambiente — decisão explícita, fora do escopo |
| Próximo módulo da série | L64 — entrega e formato: Parquet, partição, arquivo pequeno. Onde os eventos que aqui aprenderam a chegar em ordem, sem perda, viram arquivo colunar consultável sem varrer tudo |
| Data da última validação técnica | 8 de agosto de 2026 |
Documentação oficial consultada: Amazon Kinesis Data Streams — Terminology and concepts — hash MD5 da chave contra o intervalo de hash do shard; Quotas and limits — capacidade de escrita e leitura por shard, retenção padrão e máxima; PutRecords (API Reference) — FailedRecordCount e ErrorCode por registro; GetShardIterator (API Reference) — os três tipos de iterador e a expiração em 5 minutos; Choose the right mode to stream in — por que o modo on-demand não isola chave enviesada ao escalar; e Lambda parameters for Kinesis event source mappings — ParallelizationFactor preservando ordem por chave. Os valores de preço não aparecem neste módulo por decisão: use o AWS Pricing Calculator, porque preço varia por região e envelhece mais rápido que o conteúdo.
O que não foi verificado, e você deve conferir na sua conta
O volume de 80 eventos/s em dia normal e o pico de 3× são os medidos na conta de exemplo da Cadência, e servem como ordem de grandeza, não como referência. O número de shards (8) e a retenção (168h) devem ser derivados do SEU volume real e da SUA cadência de detecção de bug — não copiados. O comportamento de resharding descrito para o modo on-demand está confirmado na documentação de dimensionamento; o comportamento específico do seu stream sob a sua distribuição de chave, meça antes de assumir.
Terminou de ler?
Marcar como concluído registra o XP, mantém sua sequência e coloca 3 cartas deste módulo na fila de revisão espaçada.
Próximos passos sugeridos
Temas deste módulo
Discussão
Carregando comentários…