Pular para conteúdo

Proposta 3 — Arquivo de eventos + replay via query

Status: Brainstorming de consultoria — ver aviso metodológico Resolve: arquivamento de eventos além da retenção curta do Kafka, com replay pontual via consulta SQL Não resolve: dual-write nem DLQ em tempo real (ver Proposta 1, complementar a esta)

Contexto

Esta é a intenção real por trás do pedido inicial de "Event Sourcing": não reconstruir o estado de um agregado, mas conseguir arquivar eventos além da retenção curta do Kafka e reprocessar (replay) um evento ou uma janela de tempo específica via consulta SQL, quando necessário — sem a complexidade de um Event Store append-only. O modelo de referência trazido na discussão foi Kafka → Kinesis Firehose → Parquet (tabela Iceberg) → Athena → replay via query.

A Energisa já opera uma variante deste padrão hoje, dependente de Azure: Kafka → Kafka Connect Sink + Azure Runtime Integrator → Azure Data Lake Storage (retenção uniforme de 7 dias — risco R37 já registrado, com incidente real de perda de dado) → Databricks (Bronze/Silver/Gold) → SQL Warehouse/Delta Sharing → Power BI (ver engenharia-dados-adms.md). A variante abaixo é a alternativa Azure-independente, pedida explicitamente na discussão.

Diagrama C4 — Contêineres (variante Azure-independente)

Proposta 3 — Arquivo de eventos + replay via query

Design patterns aplicados

Padrão Papel nesta proposta
Event Archive via CDC Sink Kafka Connect S3 Sink move eventos do Kafka (retenção curta) para armazenamento de longo prazo, sem exigir mudança no produtor nem no consumidor original
Open Table Format (Iceberg) Formato de tabela sobre Parquet com catálogo próprio — permite consulta SQL sem depender de um motor proprietário (equivalente aberto ao papel do Delta Lake no pipeline Azure existente)
Partições fechadas habilitando backup convencional Mesma lógica já proposta por Castellani na Ata 10 para o Event Store do barramento: partições de 15min fechadas podem ser copiadas com backup de volume tradicional, ao contrário de um tópico Kafka em escrita contínua
Replay-by-requery Diferença central frente a Event Sourcing (Proposta 2): não reconstrói o estado de um agregado — apenas seleciona e reenvia eventos brutos já publicados, como se tivessem chegado agora
Inbox / Idempotent Consumer Necessário do lado do consumidor para absorver o reenvio do replay sem duplicar efeito — mesmo padrão da Proposta 1, reforçando que as duas propostas compartilham essa peça

Cenários de aplicação

Cenário Proposta se aplica?
Precisa recuperar de uma falha de processamento passada, reprocessando uma janela de tempo específica Sim — é o caso central
Auditoria ou consulta histórica via SQL, sem sobrecarregar a retenção do Kafka Sim
Não há necessidade de reconstruir o estado atual de um agregado a partir do zero Sim — é justamente o que diferencia esta proposta da Proposta 2
Requisito é impedir que o dual-write aconteça em primeiro lugar, em tempo real Não — isso é a Proposta 1; esta proposta assume que o evento já foi publicado com sucesso ao menos uma vez
Requisito de negócio exige reconstruir o estado de um agregado ponto-a-ponto no tempo Não — volta a ser Event Sourcing genuíno (Proposta 2)

Mecânica do sink: Iceberg vs Delta Lake (Databricks)

O componente "Kafka Connect S3 Sink" do diagrama acima esconde uma diferença de mecanismo entre as duas famílias de tabela aberta discutidas nesta e em outras sessões de consultoria — vale destrinchar, porque ela reforça (não substitui) a recomendação já registrada em "Alternativas de stack e trade-offs" abaixo com um "porquê" mecânico, não só de custo.

Iceberg Kafka Connect Sink (apache/iceberg, ou o fork databricks/iceberg-kafka-connect — mesma linhagem original da Tabular, adquirida pela Databricks em 2024; vale confirmar qual build empacotar antes de produção) escreve direto no object storage configurado, sem staging intermediário e sem depender de nenhum serviço gerenciado externo. A arquitetura interna separa dois papéis dentro do próprio cluster Connect: Workers (uma task por partição, convertem os registros Kafka em arquivos Parquet e os gravam no object storage, mas não commitam nada) e um Coordinator único, eleito entre as tasks, que só escuta um tópico de controle interno, agrupa os arquivos reportados pelos Workers a cada intervalo configurável e faz um commit atômico (uma nova snapshot) no catálogo Iceberg. O exactly-once vem de rastrear offsets em dois consumer groups paralelos (o padrão do Kafka Connect e um grupo próprio do sink). Resultado prático: o único requisito de colocation é entre o Connect e o object storage — já satisfeito pela topologia on-prem (OpenShift Data Foundation) desta proposta.

Databricks Delta Lake Sink Connector (mantido pela Confluent, disponível para Confluent Cloud/Platform) funciona de forma diferente: não escreve o par Parquet + transaction log do Delta diretamente. Ele faz staging dos registros num bucket S3 intermediário e só então comita esses dados na tabela via a API gerenciada de uma instância Databricks real. Duas implicações diretas para o caso da Energisa: (1) é append-only, e a integridade do bucket de staging é pré-requisito da garantia de exactly-once; (2) bucket de staging, tabela Delta e cluster Kafka precisam estar na mesma região de um serviço cloud gerenciado — ou seja, o mecanismo em si já pressupõe a dependência de Azure/Databricks que a variante "Azure-independente" desta proposta existe para evitar. Não é um detalhe de implementação secundário: é a mesma classe de dependência que já produziu o risco R37 (retenção uniforme de 7 dias na Landing, com perda de dado real) no pipeline hoje em produção.

Iceberg (Kafka Connect Sink) Delta Lake (Databricks Sink Connector)
Onde grava Direto no object storage configurado, sem intermediário Staging em bucket S3 → commit via API gerenciada do Databricks
Dependência de serviço gerenciado Nenhuma — só catálogo (Hive Metastore/Nessie) e object storage, ambos on-prem Sim — exige instância Databricks real, mesma região do bucket e do Kafka
Encaixe com a topologia on-prem (OpenShift) desta proposta Direto — mesmo desenho já registrado nesta proposta Reintroduz a dependência Azure que a variante "Azure-independente" existe para evitar
Garantia de exactly-once Commit atômico coordenado; offsets em dois consumer groups Depende da integridade do bucket de staging; sem UPDATE/MERGE nativo
Maturidade/proveniência Origem Tabular; hoje mantido tanto em apache/iceberg quanto em fork databricks/iceberg-kafka-connect Conector comercial da Confluent, documentado e suportado, mas atrelado ao ecossistema Databricks

Essa mecânica é o motivo técnico — não só de custo ou de risco já materializado (R37) — para a recomendação de Iceberg + OpenShift Data Foundation registrada abaixo: o sink de Delta não é "Delta rodando on-prem", é, por desenho, um caminho que sempre termina numa instância Databricks gerenciada. Se a decisão for aceitar essa dependência, o pipeline ADLS→Databricks já em produção cumpre esse papel — não haveria motivo para adicionar um segundo conector fazendo a mesma coisa em paralelo.

Alternativas de stack e trade-offs

  • Object storage: OpenShift Data Foundation (MCG/NooBaa) preferível a MinIO — MinIO community edition entrou em modo de manutenção em dezembro/2025 e foi arquivado em 25/04/2026 (distribuição apenas via source, sem binários oficiais nem novos patches), apesar de a Energisa já ter uma instância experimental de MinIO no datacenter da Paraíba (achado da Ata 10).
  • Motor de consulta: Trino como equivalente mais próximo do Athena (cluster distribuído); DuckDB como alternativa mais leve, single-process, se o volume não justificar um cluster Trino dedicado.
  • Catálogo Iceberg: exige um serviço à parte (Hive Metastore ou Nessie) — não é um componente opcional, é pré-requisito para qualquer consulta funcionar.
  • Manutenção de jobs: Tekton/OpenShift Pipelines como opção nativa do OpenShift para compactação, expiração de snapshots Iceberg e outras rotinas de manutenção da tabela.
  • Licenciamento: assim como a Proposta 1, depende de capacidade cuja licença não está confirmada — ver checklist de licenciamento (OpenShift Data Foundation e OpenShift Pipelines aparecem como "capacidade existe, licenciamento não confirmado").
  • Variante Azure: se a dependência de Azure for aceitável, o pipeline já existente (ADLS → Databricks → SQL Warehouse) cobre o mesmo papel — mas carrega os riscos já documentados: retenção uniforme de 7 dias na Landing (R37, com incidente real de perda de dado) e custo do Databricks já questionado como achado de Fase 2.