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 verO que torna uma task de pipeline segura para ser reexecutada automaticamente?
Vídeo · português · 14 min
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:
- 1Dependência: a agregação de fraude só pode começar depois que a carga de cartão terminou
- 2Ordem: o grafo declara a sequência, em vez de deixá-la implícita em horários
- 3Retry: falha transitória de rede não pode virar chamado às 3h da manhã
- 4Agendamento: por relógio, por evento ou por chegada de dado
- 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ência— Rodar 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ão | Por que quebra | O que fazer |
|---|---|---|
| DAG que processa o dado | Worker do orquestrador não escala e vira ponto único de falha | Task submete job no Databricks ou EMR e acompanha o estado |
| Dependência por horário | No dia de pico o upstream atrasa e o downstream lê dado incompleto | Sensor de partição, trigger do upstream ou evento |
| Task não idempotente | Retry automático duplica lançamento | Sobrescrita de partição ou MERGE por chave |
| Retry sem backoff nem limite | Derruba origem degradada e esconde falha permanente | Backoff exponencial com tentativas máximas |
Airflow, Databricks Workflows e Step Functions
| Ferramenta | Onde brilha | Custo de escolher |
|---|---|---|
| Airflow | Fluxo que cruza muitos sistemas; DAG em Python versionada; ecossistema grande de operators | Você opera o cluster de controle: scheduler, banco de metadados, workers |
| Databricks Workflows | Quase tudo é notebook ou job Spark na própria plataforma; integra com Unity Catalog e cluster job | Fica preso ao perímetro Databricks quando o fluxo puxa serviços de fora |
| AWS Step Functions | Fluxo feito de serviços AWS, serverless, com máquina de estados e retry declarativos | Orquestração fica em JSON e a lógica se espalha entre estados e Lambdas |
Se quiser outro ângulo
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 reverDrill
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.A DAG de extrato lê o CSV do dia e faz INSERT INTO na tabela de destino, sem chave nem deduplicação.
2.A DAG de fraude roda às 6h porque a DAG de cartão costuma terminar às 5h40.
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.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.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.INSERT das transações de cartão do dia numa tabela sem chave.
2.Sobrescrever a partição data_ref igual a 2026-03-11 com o resultado recalculado.
3.MERGE por id_transacao atualizando quando existe e inserindo quando não existe.
4.UPDATE saldo igual a saldo mais valor para cada linha do arquivo do dia.
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?
Para ir além