Onde parar o streaming, como construir idempotência em cada fronteira e cinco lições de uma plataforma de ticketing aéreo em produção.

Engenharia de dados · Azure Databricks · Leitura técnica

Neste artigo
  1. 01 A decisão vem antes da arquitetura
  2. 02 O salto extra que comprou reprocessabilidade
  3. 03 Idempotência não é uma caixa marcada uma vez
  4. 04 Onde paramos o streaming
  5. 05 Por que o Bronze virou tabela
  6. 06 Ingestão incremental sem outra infraestrutura de eventos
  7. 07 Consistência agregada acima da precisão individual
  8. 08 O INNER JOIN que devolveu zero linhas
  9. 09 A deduplicação que mudava entre execuções
  10. 10 O standard oficial que duplicou uma dimensão
  11. 11 A dívida silenciosa dos arquivos pequenos
  12. 12 Desmontar o legado também é engenharia
  13. 13 A framework fica na borda da lógica
  14. 14 O que transforma um pipeline em plataforma
  15. 15 Como falar de latência sem inventar um número
  16. 16 O resultado que importa
O requisito era tempo real. A primeira coisa que fizemos não foi arquitetar. Foi perguntar por quê.
01

A decisão vem antes da arquitetura

Processamento contínuo tem custo permanente: computação ativa, complexidade operacional e uma classe de falhas que só existe em streaming. A arquitetura só se paga quando uma decisão real perde valor ao esperar pelo dia seguinte.

Reporting comercial, isoladamente, não justificava espalhar streaming por toda a solução. Já a capacidade de detectar padrões anômalos de emissão enquanto ainda são acionáveis muda a conversa. Usamos esse critério para definir onde a baixa latência era necessária e, com a mesma disciplina, onde ela não era.

Essa pergunta virou uma restrição de projeto. Medimos o valor pela latência do sistema completo, não pela quantidade de caixas marcadas como streaming no diagrama.

Uma arquitetura em tempo real precisa começar pela decisão que não pode esperar.

02

O salto extra que comprou reprocessabilidade

A alternativa óbvia era consumir a fila diretamente no Spark Structured Streaming. Recusamos essa opção e materializamos cada mensagem como arquivo na landing zone antes de iniciar o processamento analítico.

Uma fila tem retenção limitada. O arquivo preserva exatamente o que a origem enviou e permite corrigir um parser meses depois sem pedir uma nova carga. Também desacopla a disponibilidade do pipeline do comportamento da fila e cria uma fronteira operacional simples: o arquivo chegou ou não chegou?

Em um modelo que suporta análise financeira, linhagem e auditoria não são conveniências. Cada linha precisa ser rastreável ao payload que a originou. Poupar um salto de rede valia menos do que garantir replay e diagnóstico.

Em pipelines financeiros, poder reprocessar vale mais do que economizar segundos de latência.

03

Idempotência não é uma caixa marcada uma vez

A entrega ponta a ponta não é magicamente exactly-once. O desenho combina entrega at-least-once com idempotência no consumidor, em duas fronteiras independentes.

Na integração, a mensagem permanece bloqueada até a escrita terminar. Se um lote falha no meio, ele volta para a fila. Um ledger registra apenas os itens já escritos naquele caminho de falha e impede a criação de arquivos duplicados. Depois de um lote bem-sucedido, esse estado transitório é removido.

Dentro do lakehouse existe outro problema. A origem pode enviar revisões legítimas do mesmo fato de negócio. Por isso, a deduplicação usa chave de negócio e uma ordenação explícita para escolher a versão correta. Reentrega técnica e revisão de negócio parecem duplicidade, mas exigem defesas diferentes.

A idempotência precisa existir em cada fronteira onde a entrega pode se repetir.

04

Onde paramos o streaming

A ingestão e a normalização funcionam continuamente. O modelo de negócio é materializado de forma incremental. Essa divisão não reduz a ambição do sistema. Ela protege sua correção.

As entidades finais exigem agregações e joins entre várias fontes. Forçar semântica de streaming nessa camada introduziria watermarks em múltiplos fluxos, modos de saída incompatíveis e um comportamento difícil de explicar durante uma falha. O ganho marginal de latência não compensava.

Near real-time é uma propriedade ponta a ponta. O melhor limite é aquele que entrega a janela necessária ao negócio sem importar complexidade para camadas que não se beneficiam dela.

Escolher onde parar o streaming é uma decisão de engenharia, não uma falha de ambição.

05

Por que o Bronze virou tabela

A camada Bronze alimentava dezenas de tabelas. Mantê-la como view parecia mais leve, mas faria cada consumidor abrir seu próprio leitor sobre a mesma origem, com checkpoints independentes e uma ordem de leitura impossível de garantir.

A materialização criou um único checkpoint de ingestão e um fan-out consistente. O custo de persistir essa fronteira comprou previsibilidade para todas as transformações seguintes.

Uma decisão que parece de performance pode, na verdade, ser uma decisão de correção.

06

Ingestão incremental sem outra infraestrutura de eventos

Usamos notificações de arquivo gerenciadas pela camada de storage governado. Isso retirou da solução um serviço adicional de eventos, sua fila, credenciais e o risco de diferenças entre quatro ambientes.

O schema foi declarado explicitamente porque inferência sobre XML aninhado em streaming é instável. Registros fora do schema são preservados em uma coluna de resgate, e metadados de linhagem entram no momento da ingestão. O Bronze não perde o payload só porque ainda não sabe interpretá-lo.

stream = (
    spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "xml")
        .option("cloudFiles.useManagedFileEvents", "true")
        .option("rescuedDataColumn", "_rescued")
        .schema(explicit_schema)
        .load(landing_zone_path)
)

Menos infraestrutura também significa menos credenciais, menos drift e menos pontos de falha.

07

Consistência agregada acima da precisão individual

Uma entidade distribuía o valor de um documento entre segmentos de forma proporcional. A referência usada nesse cálculo não cobria todos os pares possíveis. Aplicar um fallback apenas ao segmento sem referência fazia a soma final deixar de bater com o total do documento.

Mudamos a unidade da regra. Quando uma única referência faltava, o documento inteiro passava para o método alternativo. A precisão de um segmento poderia diminuir, mas a reconciliação do documento permanecia exata em todos os casos.

Essa mudança eliminou a reconciliação manual porque a regra passou a operar no mesmo nível em que o negócio verifica a consistência.

O fallback deve ser desenhado na unidade de reconciliação, não no registro com dado ausente.

08

O INNER JOIN que devolveu zero linhas

Uma das entidades juntava documentos à fonte considerada principal para obter valor e moeda. Em produção, o resultado tinha zero linhas. Para aquele tipo de evento, a informação nunca aparecia nessa fonte. Ela existia exclusivamente em uma alternativa.

A correção foi trocar o INNER JOIN por LEFT JOIN e declarar uma hierarquia de coalesce entre fontes. O problema não era sintaxe. O join afirmava que todo documento tinha uma correspondência obrigatória, uma hipótese de negócio escondida dentro de SQL comum.

Depois da correção, passamos a revisar joins como contratos de cardinalidade e presença. Antes de escolher o tipo, perguntamos o que a ausência realmente significa no domínio e se eliminar uma linha é uma ação válida.

INNER ou LEFT não é preferência de estilo. É uma afirmação verificável sobre o domínio.

09

A deduplicação que mudava entre execuções

A primeira versão usava dropDuplicates. A operação remove duplicados, mas não garante qual linha sobrevive quando existem duas versões da mesma chave. Reprocessar a mesma janela podia produzir um resultado diferente sem que nenhum dado de entrada tivesse mudado.

Substituímos a operação por uma janela com partição na chave de negócio e ordenação explícita pelo instante do evento. A versão escolhida deixou de depender da distribuição física do DataFrame.

Esse detalhe é central para replay. Uma plataforma reprocessável não pode apenas terminar sem erro. Ela precisa produzir a mesma resposta sobre a mesma entrada, inclusive depois de uma mudança de cluster ou de particionamento.

window = Window.partitionBy("business_key").orderBy(
    F.col("event_timestamp").desc()
)

deduped = (
    events.withColumn("_row_number", F.row_number().over(window))
          .filter(F.col("_row_number") == 1)
          .drop("_row_number")
)

Idempotência é a diferença entre reprocessar com confiança e ter medo de executar o pipeline outra vez.

10

O standard oficial que duplicou uma dimensão

Uma dimensão era enriquecida com hierarquia geográfica por meio de uma referência baseada em um standard internacional. O join parecia naturalmente um para um, mas aumentava a contagem de linhas.

A referência preservava a história de códigos que haviam sido atribuídos a entidades diferentes ao longo do tempo. Oficial não significava único. Filtramos as entradas históricas e validamos a cardinalidade da chave antes do join.

Também adotamos uma convenção: um DataFrame só recebe o sufixo de lookup depois de estar reduzido e validado para a cardinalidade esperada. O nome passou a comunicar uma garantia, não apenas uma intenção.

Standards carregam história. Valide a cardinalidade antes de assumir que um LEFT JOIN preserva linhas.

11

A dívida silenciosa dos arquivos pequenos

Em batch, um número inadequado de partições deixa um conjunto ruim de arquivos. Em execução contínua, cada micro-lote repete o problema. A tabela começa saudável e se degrada lentamente, até que o custo de listar e abrir arquivos domina as leituras.

Ajustamos execução adaptativa, coalescência de partições, tamanho-alvo e otimização automática. O objetivo não era buscar um número perfeito por lote, mas impedir que a dívida crescesse sozinha durante semanas de operação.

Monitorar apenas a duração do job não detecta esse padrão cedo. A distribuição de tamanho e quantidade de arquivos faz parte da saúde operacional de uma tabela contínua.

Em streaming, arquivos pequenos não são um incômodo pontual. São uma dívida que capitaliza sozinha.

12

Desmontar o legado também é engenharia

A migração não terminava quando o pipeline novo produzia o resultado certo. Dezenas de objetos externos geridos manualmente precisavam sair do catálogo antes da nova operação, sem apagar os dados subjacentes e sem atingir tabelas já promovidas.

Criamos um processo de limpeza versionado com alvos explícitos e uma lista de proteção. A retirada passou a ser revisável e repetível como qualquer outro artefato, em vez de depender de uma sessão manual de comandos.

Esse processo permitiu encerrar a coexistência sem perda de dados nem janela de indisponibilidade. A operação nova só se tornou simples depois que a antiga deixou de disputar nomes, responsabilidades e atenção.

A migração acaba quando o legado sai de cena de forma auditável.

13

A framework fica na borda da lógica

Pipelines declarativos não precisam prender a lógica de negócio ao runtime da framework. Mantivemos os decoradores em wrappers finos, responsáveis apenas por resolver as origens. As transformações vivem em funções puras que recebem e devolvem DataFrames.

Esse padrão permite testar regras com dados sintéticos e pytest sem iniciar o pipeline completo. Não afirmamos cobertura total. O valor demonstrável é estrutural: cada regra pode ser exercitada isoladamente, e o mesmo código é promovido entre desenvolvimento, qualidade, pré-produção e produção apenas por configuração.

@pipeline.table(name="business_entity")
def business_entity():
    return build_business_entity(
        spark.table(source_a),
        spark.table(source_b),
    )

def build_business_entity(df_a, df_b):
    # Toda a regra de negócio vive nesta função pura.
    ...

Frameworks devem orquestrar a lógica, não escondê-la.

14

O que transforma um pipeline em plataforma

Os pipelines são definidos em configuração versionada, com catálogos, caminhos e políticas de computação resolvidos por ambiente. O mesmo arquivo de transformação atravessa quatro ambientes sem receber condicionais ou nomes locais. Isso reduz a distância entre o que foi testado e o que chega à produção.

A observabilidade também usa os artefatos da própria plataforma. O log de eventos dos pipelines é materializado no catálogo e consultado por SQL. Tempos de execução, contagens, atualizações e linhagem ficam disponíveis sem espalhar chamadas de logging por cada transformação.

Essa escolha não elimina alertas nem responsabilidade operacional. Ela cria uma fonte comum para investigação. Quando um fluxo atrasa, quem está de plantão consegue separar falha de entrega, atraso de ingestão e problema de modelagem sem abrir vários sistemas antes de formular uma hipótese.

Convenções de nomes, deploy reproduzível e uma fronteira clara de responsabilidade raramente aparecem na demonstração. Ainda assim, são essas decisões que definem se a plataforma pode ser mantida por outro time depois que o projeto termina.

Operabilidade começa quando a pessoa de plantão consegue descobrir o que aconteceu sem depender de quem escreveu o pipeline.

15

Como falar de latência sem inventar um número

A latência completa tem três partes: da mensagem disponível até o arquivo, do arquivo até o Bronze e do Bronze até a entidade materializada. Medimos apenas a primeira de forma transferível. Em dois workflows independentes de desenvolvimento, a consulta ao ledger, a escrita e a confirmação completaram em menos de um segundo na mediana.

Esse número demonstra que o custo do componente stateless é estável. Ele não demonstra a latência ponta a ponta, o throughput de produção nem a frequência de mensagens. Somar uma estimativa dos outros trechos produziria um número comercialmente atraente e tecnicamente indefensável.

Para fechar a medição, o caminho correto é cruzar o instante de modificação do arquivo com o timestamp de ingestão no Bronze e, depois, usar o log de eventos para medir a materialização das entidades. O resultado deve publicar p50, p95, tamanho da amostra, período e ambiente.

Até que essa amostra exista, descrevemos o sistema como near real-time e mantemos a métrica sub-segundo restrita ao componente que realmente foi observado. Precisão na linguagem é parte da engenharia, especialmente quando a página será lida por compradores técnicos.

Uma métrica só entra no case quando sua origem, seu perímetro e suas limitações podem ser explicados.

16

O resultado que importa

O sistema opera em produção e disponibiliza continuamente informação que antes chegava à análise com latência de dias. A mesma base de código atravessa quatro ambientes sem alterações na lógica.

A landing imutável eliminou reprocessamentos artesanais. A idempotência nas duas fronteiras eliminou correções manuais de duplicados. O fallback por documento eliminou reconciliações à mão. São resultados sem um número de ROI inventado e com uma relação direta entre decisão técnica e trabalho operacional removido.

A medição segura disponível cobre apenas o componente de integração: consulta ao ledger, escrita e confirmação da fila completam um ciclo em menos de um segundo na mediana em históricos de desenvolvimento. Não transformamos esse número em latência ponta a ponta. O sistema completo é descrito honestamente como near real-time.

Uma plataforma sobrevive quando reprocessamento, auditoria e operação são propriedades do desenho, não procedimentos de emergência.

A arquitetura completa, sem atalhos de marketing

O case reúne o contexto, o desenho sanitizado, os resultados confirmados e a relação entre cada decisão de engenharia e o trabalho manual que ela eliminou.

Ver o case completo