Pular para o conteúdo

Além da ementa · tema 23 de 30

Orquestração de pipelines

35 min · 2 vídeos · 14 cards · 2 drill

Por que cai

Sabatina de engenharia de dados quase sempre entra em falha e reprocessamento, porque é lá que o candidato mostra se já operou pipeline em produção. Quem só sabe desenhar o fluxo feliz trava na primeira pergunta sobre retry duplicando lançamento.

Pré-teste · 1 de 2

responda antes de ver

O que torna uma task de pipeline segura para ser reexecutada automaticamente?

Confiança:

Vídeo · português · 14 min

Vai direto ao ponto que a banca cobra primeiro: qual problema o orquestrador resolve que o cron não resolve. Anote os motivos que você conseguiria defender com um exemplo bancário próprio.

Em uma frase

Orquestrador é o componente que decide o que roda, em que ordem, quando e o que fazer quando falha — e que não deve, ele mesmo, processar o dado.

O que ele resolve

Cinco problemas, sempre os mesmos:

  1. 1Dependência: a agregação de fraude só pode começar depois que a carga de cartão terminou
  2. 2Ordem: o grafo declara a sequência, em vez de deixá-la implícita em horários
  3. 3Retry: falha transitória de rede não pode virar chamado às 3h da manhã
  4. 4Agendamento: por relógio, por evento ou por chegada de dado
  5. 5Visibilidade: qual janela rodou, quando, com qual duração e qual resultado

Cron entrega só o quarto item. Todo o resto vira gambiarra em shell script.

Airflow em cinco peças

Uma DAG é o grafo dirigido acíclico das tasks e das dependências entre elas. Ela declara o fluxo; não é o lugar do processamento.

Um operator é o molde de uma task que executa uma ação: submeter um job no Databricks, rodar um SQL, chamar uma API. Um sensor é uma task que espera uma condição — arquivo do parceiro no S3, partição criada, tabela liberada.

O scheduler lê as DAGs, calcula quais task instances estão prontas e as enfileira. O executor define onde a task enfileirada roda: processos locais, workers Celery ou pods Kubernetes.

Idempotência é o requisito número um

idempotênciaRodar a mesma task sobre a mesma janela uma ou várias vezes deixa o destino no mesmo estado final.

Por que é o mais importante: o orquestrador vai reexecutar. É o que ele faz de madrugada, sozinho, sem perguntar. Se a task não for idempotente, a funcionalidade que existe para dar resiliência passa a ser a maior fonte de dado duplicado do lake.

Duas receitas práticas:

  • Sobrescrever a partição da data de referência em vez de dar append. Barato, mas exige janela fechada.
  • MERGE pela chave de negócio (identificador da transação, por exemplo). Aguenta dado atrasado e correção, mas lê o destino, então custa mais.

E o detalhe que derruba backfill: parametrize sempre pela data lógica da execução, nunca por current_date. No reprocessamento de janeiro, o relógio marca hoje.

Retry, backoff e backfill

Retry cobre falha transitória: rede, throttling, cluster subindo. Sem backoff, tentar de novo a cada 30 segundos pressiona uma origem já degradada. Com número máximo de tentativas, falha permanente para de se disfarçar de amarelo eterno.

Backfill é reexecutar janelas passadas — porque falhou ou porque a regra mudou. Reprocessar seis meses de extrato só é operação de rotina se cada janela for independente e idempotente. Caso contrário é projeto.

SLA e alerta

SLA é o prazo declarado da entrega. Sem ele, o pipeline travado às 2h é descoberto às 9h pelo usuário do relatório regulatório. Alerta útil dispara em dois eventos: falha da task e SLA estourado — este último pega o pipeline que está rodando, mas devagar demais.

Anti-padrões

Anti-padrãoPor que quebraO que fazer
DAG que processa o dadoWorker do orquestrador não escala e vira ponto único de falhaTask submete job no Databricks ou EMR e acompanha o estado
Dependência por horárioNo dia de pico o upstream atrasa e o downstream lê dado incompletoSensor de partição, trigger do upstream ou evento
Task não idempotenteRetry automático duplica lançamentoSobrescrita de partição ou MERGE por chave
Retry sem backoff nem limiteDerruba origem degradada e esconde falha permanenteBackoff exponencial com tentativas máximas

Airflow, Databricks Workflows e Step Functions

FerramentaOnde brilhaCusto de escolher
AirflowFluxo que cruza muitos sistemas; DAG em Python versionada; ecossistema grande de operatorsVocê opera o cluster de controle: scheduler, banco de metadados, workers
Databricks WorkflowsQuase tudo é notebook ou job Spark na própria plataforma; integra com Unity Catalog e cluster jobFica preso ao perímetro Databricks quando o fluxo puxa serviços de fora
AWS Step FunctionsFluxo feito de serviços AWS, serverless, com máquina de estados e retry declarativosOrquestração fica em JSON e a lógica se espalha entre estados e Lambdas

Se quiser outro ângulo

Mostra a DAG sendo escrita, então você associa nome a coisa: operator, dependência entre tasks e agendamento. Repare em onde o código coloca a lógica de negócio.

Como cai na sabatina

Seu pipeline diário falhou às 3h e rodou de novo às 6h. O que precisa ser verdade para isso não duplicar dados?

Erros comuns

  • Colocar o processamento pesado dentro da task do Airflow. O orquestrador vira gargalo e ponto único de falha; o certo é ele disparar o job no Databricks ou no EMR e acompanhar o estado.
  • Encadear pipelines por horário em vez de por sinal. Se o upstream atrasa 20 minutos, o downstream lê dado incompleto e ninguém percebe até o relatório sair errado.
  • Escrever task não idempotente. Qualquer retry, backfill ou clique em rerun duplica dado, e retry é justamente o que o orquestrador faz sozinho de madrugada.
  • Configurar retry infinito ou sem backoff. Se a origem caiu, martelar a cada 30 segundos derruba de vez e ainda esconde o incidente atrás de uma task que fica amarela para sempre.
  • Tratar SLA como enfeite. Sem prazo declarado, o pipeline que travou às 2h só é descoberto às 9h pelo usuário do relatório, e aí o custo já é reputacional.
  • Confundir data lógica da execução com a hora do relógio. Se a task usa current_date em vez da data da janela, o backfill de ontem grava tudo com a data de hoje.

Flashcards

Card 1 de 14

0 certos · 0 a rever

Drill

Ache o defeito na DAG

Cada item descreve uma DAG real. Diga qual é o defeito principal antes de abrir a resposta. Use exatamente uma das categorias.

  1. 1.A DAG de extrato lê o CSV do dia e faz INSERT INTO na tabela de destino, sem chave nem deduplicação.

  2. 2.A DAG de fraude roda às 6h porque a DAG de cartão costuma terminar às 5h40.

  3. 3.A task de agregação de PIX carrega 400 GB em pandas dentro do worker do Airflow e escreve o resultado no S3.

  4. 4.A DAG de cadastro tem retries igual a zero e nenhuma notificação configurada; quando falha, ninguém sabe até alguém abrir a interface.

  5. 5.A task de reprocessamento usa current_date para montar o caminho de escrita no lake.

Drill

É idempotente?

Para cada operação de escrita, decida se rodar duas vezes sobre a mesma janela deixa o destino igual.

  1. 1.INSERT das transações de cartão do dia numa tabela sem chave.

  2. 2.Sobrescrever a partição data_ref igual a 2026-03-11 com o resultado recalculado.

  3. 3.MERGE por id_transacao atualizando quando existe e inserindo quando não existe.

  4. 4.UPDATE saldo igual a saldo mais valor para cada linha do arquivo do dia.

  5. 5.Copiar o arquivo bruto do parceiro para a zona raw usando o nome original como destino.

Quiz · 1 de 2

Qual afirmação descreve corretamente a divisão de papéis entre scheduler e executor no Airflow?

Explique para um gerente

Explique para o gerente de operações por que o time precisa de um orquestrador, se o cron do servidor já dispara os scripts todo dia às 2h.

Responda em voz alta antes de seguir. Se travar numa palavra técnica, é sinal de que ainda não entendeu essa parte.

Perguntas de sabatina deste tema

  • · Seu pipeline diário de conciliação falhou às 3h e o orquestrador rodou de novo às 6h. O que precisa ser verdade para isso não duplicar dados?
  • · Um time diz que não precisa de orquestrador porque o cron do servidor já dispara os scripts na ordem certa. Como você responde?
Responder no simulado

Para ir além