Reconstruir o que aconteceu antes de repetir
Uma equipa fictícia APS prepara um ficheiro diário de posições de fundos. A extração terminou, a transformação parece concluída e o carregamento perdeu a resposta de confirmação. O relógio aproxima-se da hora de entrega. O operador recebe o pedido de repetir tudo. Antes de executar, identifica a unidade de trabalho: data de negócio, partição de entrada, versão da transformação, destino e identificador remoto de cada job. Sem estes elementos, repetir pode significar processar outros dados ou acrescentar linhas já carregadas. Separa três perguntas no registo do incidente: o pedido chegou ao serviço, o job terminou e o resultado está correto? Em BigQuery, DONE significa que o job já não está a executar; verifica errorResult para saber se falhou. A ausência desse erro ainda pode coexistir com erros não fatais, que devem ser comparados com os critérios de aceitação. Um carregamento que tolerou linhas inválidas pode estar tecnicamente concluído e continuar incompleto para o negócio. O objetivo do suporte é recuperar uma entrega reconciliada, preservando a evidência que permite explicar cada decisão.
Controlar recursos antes da execução
Um analista propõe reduzir o custo de uma consulta pesada acrescentando LIMIT 10. Numa tabela não clustered, limitar linhas devolvidas não reduz por si só os bytes lidos. Analisa projeção de colunas, filtros e partições efetivamente necessárias. Para o modelo on-demand, maximum bytes billed pode impedir uma consulta cuja estimativa ultrapasse o limite. Não confundas esse controlo com uma promessa de duração, nem o transportes sem análise para faturação por capacidade. No exemplo, o limite aprovado é 200 unidades de bytes didáticas e a estimativa é 260. A ação correta é rever o plano ou obter a alteração de orçamento aplicável, e não aumentar automaticamente o limite em cada retry. Para tabelas clustered, a estimativa pode ser um limite superior e rejeitar uma consulta cujo custo real seria menor. Regista essa incerteza ao justificar o ajuste. Se vários passos recalculam a mesma transformação, considera materializar um resultado intermédio com versão, retenção e verificação de integridade. Compara o custo de processamento evitado com armazenamento e manutenção; uma tabela intermédia sem proprietário pode criar outro problema operacional.
Fixar o significado de uma nova tentativa
Uma tarefa Airflow que lê sempre a partição mais recente pode produzir resultados diferentes quando é repetida. O relógio now() usado para escolher a data de negócio tem o mesmo problema. Usa o intervalo de dados ou outro identificador estável da execução para selecionar entradas e saídas. Regista alterações legítimas da origem como uma nova versão de processamento, em vez de as esconder numa repetição supostamente equivalente. Esta escolha torna possível comparar tentativas e explicar diferenças ao negócio. Para um carregamento BigQuery com WRITE_APPEND, conserva o job ID escolhido pelo cliente. Se a submissão perder a resposta e uma repetição com esse ID devolver duplicate, consulta o job existente. Não cries logo um ID aleatório para contornar o erro. Um novo ID previsível pode fazer sentido depois de confirmar que as tentativas anteriores falharam. O identificador remoto, o identificador do lote e a chave dos registos têm funções diferentes. O primeiro ajuda a gerir submissões; não demonstra sozinho que dois ficheiros diferentes não contêm os mesmos movimentos. A reconciliação de negócio continua necessária antes de disponibilizar o resultado.
Limitar concorrência pela dependência real
Adicionar workers ao orquestrador pode aumentar pressão sobre uma base ou API já saturada. Um pool Airflow limita paralelismo de um conjunto de tarefas. Associa as tarefas que partilham a mesma dependência ao mesmo controlo e estima o seu peso. Se uma tarefa pesada ocupa dois slots num pool de dois, as tarefas leves desse pool aguardam. A existência de workers livres não cria slots adicionais nesse pool. No serviço fictício, a base de posições suporta uma extração pesada ou duas pequenas em simultâneo durante a janela acordada. Esta é uma hipótese do exercício, que precisaria de medição no sistema real. Documenta a relação entre essa capacidade e pool_slots; não assumes que um slot corresponde automaticamente a uma ligação de base de dados. Verifica também se tarefas deferred contam como ocupação, conforme a configuração do pool. Libertar recursos locais enquanto um job externo continua ativo pode admitir trabalho adicional sobre a mesma dependência. O plano de capacidade deve considerar onde o trabalho realmente continua, o prazo de negócio e a prioridade das filas. Um retry elegível ainda pode ter de esperar por capacidade.
Orçamento de tentativas e espera
Uma política de retry tem pelo menos três decisões: quais os erros elegíveis, quantas repetições admitir e quanto esperar entre elas. Em Workflows, max_retries não inclui a execução inicial. O laboratório usa maxAttempts, que a inclui. Assim, max_retries=3 permite até quatro execuções, enquanto maxAttempts=3 permite três. Ao traduzir configurações entre ferramentas, escreve uma sequência de tentativas para evitar este erro de contagem. Escolhe a política HTTP segundo a semântica do passo. A política predefinida para passos não idempotentes cobre menos situações do que a política para passos idempotentes; um timeout pode deixar o resultado remoto desconhecido. O laboratório é deliberadamente conservador: estado unknown exige reconciliação mesmo se safeRepeat for true. Para um estado retryable, aplica espera exponencial limitada e calcula o fim estimado da próxima tentativa, incluindo validação. Não modela jitter, filas ou falhas futuras. Num serviço real, a dispersão das tentativas evita que muitos clientes regressem ao mesmo instante. Aumentar o número máximo sem considerar o prazo pode apenas prolongar um incidente que já precisa de outra mitigação.
Executar o exercício de decisão
Executa python3 run.py < case.json no laboratório pde-retry-window. Os tempos são segundos num relógio didático comum, não timestamps reais. A fixture usa asOf=100, deadline=140 e idade máxima de observação de dez segundos. A política permite quatro tentativas totais, começa com espera de quatro segundos, duplica a espera e limita-a a doze. Cada tarefa declara estado, número de tentativas iniciadas, instante de observação, fim confirmado da última tentativa e durações estimadas de execução e validação. report falhou duas vezes e terminou a última tentativa em 95. A espera é oito; o próximo início possível é 103. Com vinte segundos de execução e cinco de validação, termina em 128 e recebe wait-backoff. load pode começar em 100 e precisa de 35+5 segundos: chega exatamente ao prazo e recebe retry-eligible. append tem resultado desconhecido e recebe reconcile-remote. archive esgotou tentativas. A observação de monitor está antiga e deve ser atualizada. transform terminou com sucesso, mas a saída ainda precisa de validação. O programa só calcula decisões; não espera, não observa serviços e não submete jobs.
Testar fronteiras e contestar a previsão
Reduz deadline de 140 para 139: load deixa de caber no prazo. Mantém 140 mas avança asOf para 101: também deixa de caber, porque a espera humana consumiu um segundo. A validação não pode ser removida do cálculo só para obter um resultado favorável. Em report, avançar asOf para 103 termina a espera de backoff e torna a tentativa elegível, desde que a observação continue dentro da idade aceite. Testa ainda o limite de observação: uma amostra de 90 é aceite em 100 com limite dez; uma amostra de 89 exige atualização. Um instante de observação futuro é uma entrada inválida, não evidência mais fresca. Altera safeRepeat para false num erro retryable: o modelo exige reconciliação dos efeitos. Esse booleano é uma afirmação do chamador; o programa não prova idempotência. Os testes permutam as seis tarefas em 720 ordens e exigem o mesmo resultado, porque cada decisão é independente. Isto não demonstra que todas possam executar em simultâneo. Dependências, capacidade e filas teriam de entrar noutro modelo e nos ensaios do serviço.
Mitigar sem perder a autonomia de RUN
Um pedido de cancelamento BigQuery não garante que o job tenha sido cancelado; este pode ter terminado entretanto. Custos também podem permanecer. Antes de iniciar uma substituição, consulta o estado e verifica os efeitos. Distingue cancelar trabalho em execução de desfazer resultados já publicados. Se a entrega original já não cabe no prazo, comunica impacto, população afetada e alternativa operacional, em vez de esconder a previsão atrás de novas tentativas. Na passagem a RUN, entrega um procedimento com identificadores, localização dos jobs, fonte do estado, classificação de erros, limite de tentativas, critérios de validação e responsáveis pelo escalamento. Inclui exemplos de duplicate após resposta perdida, DONE com erro e deadline-infeasible no exercício. Pede a outro operador que explique a decisão sem ajuda do autor. Para a reunião internacional, prepara uma atualização curta em inglês: estado confirmado, resultado ainda desconhecido, ação seguinte e prazo da próxima evidência. O fecho deve apoiar-se nos dados reconciliados e na aceitação do consumidor. O laboratório ensina a tornar a decisão explícita; a validação do comportamento real e a revisão especializada continuam a ser trabalho separado.
"""Original offline retry-decision model. Does not submit, cancel or observe jobs."""
import json
import re
import sys
def require(condition, message):
if not condition:
raise ValueError(message)
def integer(value, low, high):
return type(value) is int and low <= value <= high
def keys(value, fields):
require(type(value) is dict and set(value) == set(fields.split()), 'Unexpected fields')
def evaluate(payload):
keys(payload, 'asOf deadline maxObservationAge policy tasks')
for name in ['asOf', 'deadline', 'maxObservationAge']:
require(integer(payload[name], 0, 1000000000), 'Invalid global time')
policy = payload['policy']
keys(policy, 'maxAttempts initialDelay maxDelay multiplier')
require(integer(policy['maxAttempts'], 1, 10), 'Invalid maxAttempts')
require(integer(policy['initialDelay'], 1, 1000000), 'Invalid initialDelay')
require(integer(policy['maxDelay'], policy['initialDelay'], 1000000), 'Invalid maxDelay')
require(integer(policy['multiplier'], 1, 5), 'Invalid multiplier')
tasks = payload['tasks']
require(type(tasks) is list and 1 <= len(tasks) <= 20, 'Need 1..20 tasks')
ids, results = set(), []
for task in tasks:
keys(task, 'id state safeRepeat attempts observedAt lastFinishedAt duration validation')
identifier = task['id']
require(type(identifier) is str and re.fullmatch(r'[A-Za-z0-9_-]{1,48}', identifier) is not None, 'Invalid task id')
require(identifier not in ids, 'Duplicate task id')
ids.add(identifier)
require(task['state'] in ['running', 'unknown', 'succeeded', 'retryable', 'terminal'], 'Invalid state')
require(type(task['safeRepeat']) is bool, 'safeRepeat must be boolean')
require(integer(task['attempts'], 1, 10), 'Invalid attempts')
require(integer(task['observedAt'], 0, payload['asOf']), 'Invalid observation time')
require(integer(task['duration'], 1, 1000000), 'Invalid duration')
require(integer(task['validation'], 0, 1000000), 'Invalid validation duration')
ended = task['lastFinishedAt']
if task['state'] in ['running', 'unknown']:
require(ended is None, 'Unconfirmed completion must use null')
else:
require(integer(ended, 0, task['observedAt']), 'Invalid completion time')
result = {'id': identifier, 'decision': None, 'backoff': None,
'earliestStart': None, 'estimatedFinish': None}
if payload['asOf'] - task['observedAt'] > payload['maxObservationAge']:
result['decision'] = 'refresh-observation'
elif task['state'] == 'running':
result['decision'] = 'observe-running'
elif task['state'] == 'unknown':
result['decision'] = 'reconcile-remote'
elif task['state'] == 'succeeded':
result['decision'] = 'validate-output'
elif task['state'] == 'terminal':
result['decision'] = 'repair-cause'
elif task['attempts'] >= policy['maxAttempts']:
result['decision'] = 'attempts-exhausted'
elif not task['safeRepeat']:
result['decision'] = 'reconcile-effects'
else:
delay = min(policy['maxDelay'], policy['initialDelay'] * policy['multiplier'] ** (task['attempts'] - 1))
start = max(payload['asOf'], ended + delay)
finish = start + task['duration'] + task['validation']
result.update(backoff=delay, earliestStart=start, estimatedFinish=finish)
if finish > payload['deadline']:
result['decision'] = 'deadline-infeasible'
elif start > payload['asOf']:
result['decision'] = 'wait-backoff'
else:
result['decision'] = 'retry-eligible'
results.append(result)
return {'tasks': sorted(results, key=lambda item: item['id']),
'deadlinePassed': payload['asOf'] > payload['deadline'],
'jobsExecuted': False, 'idempotencyProven': False,
'completionGuaranteed': False, 'productionApproval': False}
def main():
raw = sys.stdin.read(1000001)
require(len(raw) <= 1000000, 'Input too large')
print(json.dumps(evaluate(json.loads(raw)), sort_keys=True))
if __name__ == '__main__':
try:
main()
except (ValueError, TypeError, RecursionError):
print('Invalid retry fixture', file=sys.stderr)
sys.exit(2)
No relógio didático, report aguarda até 103 e termina estimadamente em 128; load chega ao prazo 140 incluindo validação. append exige reconciliação remota.
Armadilhas comuns
Tratar timeout como falha remota; DONE como sucesso; max_retries como tentativas totais; workers livres como slots livres; estimativa favorável como autorização de produção.
Tópicos relacionados: Operação e capacidade · Ingestão e quarentena · Recuperação de conjuntos de dados
Repetir exige estado suficientemente recente, efeitos compreendidos, capacidade e tempo para validar a entrega.
Referência: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)