1. Operar o resultado que o negócio utiliza
Um pipeline não está recuperado apenas porque o processo está ativo. Define qual é o resultado que o consumidor precisa de receber, em que instante e com que critérios de aceitação. No nosso caso fictício APS, o serviço entrega lotes reconciliados para reporting. Um lote produzido fora do prazo ou com discrepâncias não conta como o mesmo resultado aceite. Regista a unidade de negócio, a janela, os consumidores, o responsável pela aceitação e o tratamento de uma entrega provisória. Estes exemplos não representam procedimentos BNP Paribas. Compara custo e resultado na mesma população. Passar de 120 unidades de custo para 90 parece uma melhoria, mas entregar apenas seis lotes aceites em vez de doze muda o custo por resultado de 10 para 15. Investiga a perda de entregas antes de apresentar a alteração como eficiência. Para a passagem a RUN, associa cada sinal a uma decisão: quem recebe o alerta, qual é o primeiro diagnóstico e quando deve escalar. Não se exige uma métrica perfeita para agir; exige-se que as limitações da evidência acompanhem a decisão. Um indicador de capacidade ajuda a prever; uma reconciliação confirma o conteúdo entregue.
2. Reduzir desperdício sem perder controlo
Dimensiona compute a partir da carga e da janela útil. Um cluster permanente pode servir trabalho frequente e previsível; uma execução isolada pode justificar compute criado para esse trabalho. Num workflow com managed cluster Spark, prepara a saída para sobreviver ao fim do cluster. Se o único ficheiro está no disco local do worker, reduzir o tempo de vida do cluster pode destruir o resultado de que o consumidor precisa. Inclui publicação durável e validação no workflow, e mede a entrega completa. A documentação atual apresenta Managed Service for Apache Spark e Managed Service for Apache Airflow; o guia do exame ainda usa Dataproc e Cloud Composer. Conserva essa correspondência sem inferir uma nova edição do exame. No controlo financeiro, distingue orçamento Alerts-only de Spend cap. O primeiro, sem automação adicional, comunica o desvio; não suspende utilização por si. Confirma sempre o tipo e o âmbito configurados. Também valida a população de um relatório de consumo: somar um job SCRIPT com os filhos pode duplicar valores agregados. Uma diferença de custo deve levar a investigação de workload, configuração e unidade de medição.
3. Automatizar repetições com significado estável
Desenha a tarefa para repetir a mesma unidade lógica. Uma publicação pode ter sido concluída no destino mesmo que o orquestrador tenha perdido a confirmação. Nesse caso, repetir um append com uma identidade nova transforma uma recuperação numa duplicação. Guarda uma identidade estável do lote, o resultado persistido e a versão da transformação. O contrato do destino deve permitir reconhecer o resultado anterior ou aplicar uma operação idempotente. Um timeout sozinho não prova sucesso nem ausência de commit. O controlo começa antes da execução. Chamadas de rede feitas ao carregar um ficheiro DAG podem atrasar repetidamente o seu processamento. Separa definição e execução, usando configuração preparada ou consultas realizadas pela tarefa quando apropriado. Testa dependências: transform não deve começar antes de extract ter a saída necessária. A ordem textual no template não é essa garantia. Finalmente, corrigir um template não demonstra que uma execução já ativa recebeu a correção. Regista a versão da instância afetada e a ação que lhe foi aplicada. A evidência de recuperação deve identificar o que correu, com que entradas e qual foi o resultado publicado.
4. Ligar capacidade ao prazo de negócio
Organiza workloads segundo criticidade, concorrência e custo aceitável. Reservas e atribuições BigQuery permitem gerir capacidade para diferentes populações de trabalho. Nomes de datasets, labels e contas de serviço têm outras funções; não demonstram, isoladamente, capacidade dedicada. Valida o dimensionamento e o comportamento de partilha usados no ambiente. Uma reserva mal dimensionada continua a poder falhar o objetivo do consumidor. O prazo inclui espera, execução, publicação e validação. Não uses apenas a duração de uma execução anterior que começou imediatamente. Quando um job batch está em fila, estima a margem disponível e identifica que mudança pode realmente afetá-lo. Reatribuir um projeto da reserva A para B não desloca retroativamente os pedidos já em fila ou execução em A; novos pedidos seguem para B. Esta distinção evita anunciar uma mitigação que ainda não alcança o trabalho afetado. No relatório de incidente, separa os jobs antigos dos novos e observa ambos. Uma alteração de prioridade ou capacidade é uma hipótese de melhoria até existir evidência de entrega dentro do prazo e com conteúdo aceite.
5. Investigar sinais antes de aumentar recursos
Começa pelos sintomas e pela distribuição do trabalho. Backlog crescente com todos os workers ocupados sugere uma investigação diferente de backlog crescente com apenas um worker saturado. Uma chave muito frequente pode limitar paralelismo. Nesse caso, avalia uma redistribuição compatível com a agregação e a ordenação exigidas; acrescentar workers não cria automaticamente capacidade útil para essa chave. Confirma a hipótese com dados por etapa e worker, evitando que uma média esconda o problema. Lê o motivo dos erros. quotaExceeded não significa universalmente falta de slots: identifica a quota, o âmbito e a operação afetada antes de escolher a resposta. Também distingue frescura dos dados, latência de processamento e volume pendente. Nenhuma destas medidas, isolada, prova que o destino recebeu tudo. Se uma série deixa de chegar, regista uma lacuna de observação. Preencher o gráfico com zero pode esconder a lacuna. Verifica o tratamento de dados ausentes na política de alerta e procura evidência independente no destino. O fecho automático de um alerta não deve ser usado como certificado de recuperação funcional.
6. Laboratório: backlog, capacidade líquida e idade da amostra
Executa python3 content/labs/pde-backlog-budget/run.py < content/labs/pde-backlog-budget/case.json a partir da raiz. O programa original usa apenas Python local. now e deadline são segundos inteiros numa escala fictícia. maxSampleAge é a idade máxima aceite, inclusiva. stages contém entre uma e cem etapas únicas; samples admite até mil amostras com id, stage, at, backlog, inRate e outRate. Não usa credenciais nem consulta serviços cloud. Para cada etapa, seleciona o instante mais recente que não excede now. Sem amostra, devolve unknown/missing; com idade excessiva, unknown/stale. Amostras nesse instante com valores contraditórios produzem unknown/conflicting. Valores iguais conservam todos os IDs como proveniência. Define net=outRate−inRate e projeta backlogNow=max(0,backlog−net×idade), assumindo taxas constantes desde a amostra até ao prazo. O estado within-model exige net>=0 e backlogNow<=net×(deadline−now). A comparação inteira inclui o limite exato. Uma fila vazia com entrada superior à saída continua em risco porque voltará a crescer. A estimativa arredondada de segundos é apenas uma apresentação da mesma hipótese.
7. Interpretar o exercício e tentar refutar a previsão
No ficheiro de exemplo, now=1000 e deadline=1060. ingest tem backlog 1200 e taxas 80/100: a capacidade líquida 20 permite drenar em 60 segundos, exatamente no limite. fold tem o mesmo backlog e taxas 100/110: precisa de 120 segundos e fica em risk. export tem uma amostra de idade 20 quando o máximo é 10, pelo que permanece unknown. Uma amostra futura de fold é ignorada. O resultado global conserva simultaneamente a lista de riscos e a lista de etapas desconhecidas. Reduz o prazo em um segundo e observa ingest falhar. Move a amostra de ingest para 990: a projeção desconta dez segundos de drenagem, sem alegar que foi medida novamente. Cria duas amostras contraditórias no mesmo instante e confirma que mudar a ordem não escolhe uma delas. Usa backlog 3, entrada 0, saída 2 e um segundo disponível para verificar arredondamento. Mesmo que todas as etapas fiquem within-model, não está provado o prazo de ponta a ponta: há dependências, taxas variáveis, publicação e reconciliação fora do modelo. O programa mantém explícitas essas limitações.
8. Recuperar dados e devolver autonomia a RUN
Escolhe a recuperação segundo o tipo de falha. HA de Cloud SQL entre zonas da mesma região cobre um âmbito diferente de uma indisponibilidade regional. Uma réplica assíncrona noutra região pode apoiar recuperação, mas a promoção exige considerar transações ainda não replicadas. Não transformes a existência da réplica numa promessa de RPO zero. Identifica a evidência disponível sobre o ponto recuperado e comunica a incerteza que permanece. Uma alteração lógica incorreta pode já ter sido replicada. Para PITR PostgreSQL em Cloud SQL, planeia uma nova instância, a validação do instante escolhido e o encaminhamento de clientes. A recuperação técnica termina antes da aceitação funcional: valida conectividade, permissões, dados esperados e capacidade para a carga de RUN. No handover, entrega o procedimento repetível, responsáveis, limites de autoridade, evidência do último ensaio e critérios para escalar. Atualiza a comunicação de negócio com resultado entregue, atraso e discrepâncias conhecidas. O resumo operacional deve permitir a um colega continuar o trabalho sem adivinhar o que foi observado, o que foi inferido e o que ainda falta confirmar.
"""Original offline fluid-queue exercise; no cloud calls or recovery approval."""
import json
import re
import sys
def integer(value, maximum):
if type(value) is not int or not 0 <= value <= maximum:
raise ValueError('invalid non-negative integer')
return value
def identifier(value):
if not isinstance(value, str) or not re.fullmatch(r'[A-Za-z0-9_-]{1,64}', value):
raise ValueError('invalid identifier')
return value
def evaluate(data):
if not isinstance(data, dict) or set(data) != {'now', 'deadline', 'maxSampleAge', 'stages', 'samples'}:
raise ValueError('invalid input shape')
now = integer(data['now'], 10**9)
deadline = integer(data['deadline'], 10**9)
max_age = integer(data['maxSampleAge'], 10**6)
if deadline < now:
raise ValueError('deadline precedes evaluation')
stages, samples = data['stages'], data['samples']
if not isinstance(stages, list) or not 1 <= len(stages) <= 100:
raise ValueError('stages must contain 1..100 identifiers')
if not isinstance(samples, list) or len(samples) > 1000:
raise ValueError('too many samples')
for stage in stages:
identifier(stage)
if len(set(stages)) != len(stages):
raise ValueError('duplicate stage')
ids = set()
for row in samples:
if not isinstance(row, dict) or set(row) != {'id', 'stage', 'at', 'backlog', 'inRate', 'outRate'}:
raise ValueError('invalid sample shape')
identifier(row['id'])
if row['id'] in ids:
raise ValueError('duplicate sample identifier')
ids.add(row['id'])
identifier(row['stage'])
if row['stage'] not in stages:
raise ValueError('unknown stage')
integer(row['at'], 10**9)
for field in ['backlog', 'inRate', 'outRate']:
integer(row[field], 10**6)
results = []
for stage in sorted(stages):
rows = [r for r in samples if r['stage'] == stage and r['at'] <= now]
result = {'stage': stage, 'status': 'unknown', 'reason': 'missing',
'sampleIds': [], 'sampleAt': None, 'age': None,
'projectedBacklogNow': None, 'netDrainRate': None, 'drainSecondsCeil': None}
if rows:
at = max(r['at'] for r in rows)
latest = [r for r in rows if r['at'] == at]
result.update(sampleIds=sorted(r['id'] for r in latest), sampleAt=at, age=now-at)
if now - at > max_age:
result['reason'] = 'stale'
elif len({(r['backlog'], r['inRate'], r['outRate']) for r in latest}) > 1:
result['reason'] = 'conflicting'
else:
row = latest[0]
net = row['outRate'] - row['inRate']
backlog_now = max(0, row['backlog'] - net * (now-at))
seconds = ((backlog_now + net - 1) // net if net > 0
else 0 if net == 0 and backlog_now == 0 else None)
fits = net >= 0 and backlog_now <= net * (deadline-now)
reason = ('within-window' if fits else 'growing' if net < 0
else 'no-drain' if net == 0 else 'insufficient-window')
result.update(status='within-model' if fits else 'risk', reason=reason,
projectedBacklogNow=backlog_now, netDrainRate=net,
drainSecondsCeil=seconds)
results.append(result)
risk = [r['stage'] for r in results if r['status'] == 'risk']
unknown = [r['stage'] for r in results if r['status'] == 'unknown']
return {'results': results, 'riskStages': risk, 'unknownStages': unknown,
'status': 'risk' if risk else 'unknown' if unknown else 'within-model',
'endToEndDeadlineProven': False, 'dataCompletenessProven': False,
'cloudCapacityMeasured': False, 'productionApproval': False}
def unique_object(pairs):
result = {}
for key, value in pairs:
if key in result:
raise ValueError('duplicate JSON field')
result[key] = value
return result
if __name__ == '__main__':
try:
raw = sys.stdin.read(1000001)
if len(raw) > 1000000:
raise ValueError('input too large')
payload = json.loads(raw, object_pairs_hook=unique_object)
print(json.dumps(evaluate(payload), sort_keys=True))
except (ValueError, TypeError) as error:
print('invalid input: ' + str(error), file=sys.stderr)
sys.exit(2)
Backlog 1200, entrada 100/s e saída 110/s exigem 120 segundos; uma janela de 60 segundos é insuficiente.
Armadilhas comuns
Saída bruta como drenagem; ausência como zero; template novo como execução corrigida; HA zonal como DR regional; alerta como limite de despesa.
Tópicos relacionados: Ingestão e replay · Fidelidade e reconciliação · Continuidade operacional
Uma previsão depende das taxas e da evidência; confirma a entrega aceite antes de declarar recuperação.
Referência: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)