1. Saber o que entrou e o que ficou por tratar
Num exemplo fictício APS, o pipeline do fecho termina com sucesso, mas alguns movimentos foram descartados durante parsing. O dashboard verde descreve o processo, não a completude da entrega. Define três populações: entradas recebidas, eventos lógicos aceites e entradas rejeitadas. Um evento pode chegar várias vezes, pelo que comparar apenas contagens finais com contagens de transporte produz conclusões erradas. Estes casos não representam procedimentos BNP Paribas. Antes de escolher serviços, fixa o contrato de origem e destino: identidade da entrada, identidade do evento, representação dos campos, transformações permitidas e critérios de aceitação. Uma alteração silenciosa de vírgula para ponto decimal pode parecer conveniente, mas precisa de uma regra acordada. Para cada rejeição, conserva uma referência que permita localizar a evidência, uma razão compreensível e um responsável pelo tratamento. A equipa de produção deve conseguir responder a perguntas concretas: quantas entradas foram recebidas, quantas contribuíram para resultados aceites e quantas continuam pendentes? Em que prazo podem ser recuperadas? Quem aceita uma entrega parcial? O sucesso técnico só é uma parte desta resposta. Regista limites de observação e evita que um valor desconhecido seja apresentado como zero.
2. Separar filtragem, entrega e validação
Uma subscrição Pub/Sub com filtro entrega apenas as mensagens que correspondem aos atributos selecionados. As restantes são confirmadas automaticamente pelo serviço; o filtro não examina o JSON do campo data. Assim, um consumidor que validou 700 de 1000 mensagens publicadas não pode apresentar todas como validadas. Define o âmbito esperado antes de comparar números. Um filtro sintaticamente correto pode excluir um produtor que deixou de preencher um atributo. A política de dead letter pertence à subscrição. O serviço pode encaminhar mensagens que o consumidor não consegue confirmar, mas o máximo de tentativas é aproximado. A contagem depende de configuração e permissões corretas e pode reiniciar em algumas condições. Não a uses como orçamento rígido de operações de negócio. A mensagem encaminhada inclui um novo envelope e informação sobre a subscrição original; o leitor da quarentena deve preservar essa correlação. Confirma as permissões da identidade gerida de Pub/Sub para publicar no destino e confirmar na subscrição. O operador conseguir publicar manualmente não prova esse acesso. O destino também precisa de um caminho de consumo, observação e tratamento. Criar um tópico não cria automaticamente uma operação de suporte completa.
3. Recuperar erros sem inventar dados
Classifica a falha antes de escolher retry. Uma indisponibilidade temporária de um serviço pode desaparecer; uma data impossível continua impossível quando repetes o mesmo parsing. Não substituas um valor rejeitado por zero ou pela data atual apenas para obter uma execução verde. Um erro de registo pode seguir uma saída separada; uma falha sistémica de autenticação ou armazenamento precisa de tratamento operacional, em vez de ser escondida como milhares de dados inválidos. O padrão de quarentena da aplicação permite guardar elementos não processados para investigação e correção. A política de dead letter do transporte não substitui automaticamente essa saída, sobretudo quando a aplicação já confirmou a entrega. Confirma que o destino de rejeições recebe realmente os registos e que uma falha nesse destino não fica invisível. Contabiliza falhas de classificação e de persistência separadamente. O backoff reduz pressão durante certas falhas transitórias, mas não é um agendador exato. Também não transforma uma operação em idempotente. Se um pedido perdeu a resposta, pode ter produzido o efeito. Consulta estado por referência quando possível e define orçamento, limite de concorrência e condição de escalamento. Replicar tentativas em várias camadas pode multiplicar chamadas e custo.
4. Executar o gate de quarentena local
O exercício Python abaixo classifica um lote fictício em memória. Não chama APIs, guarda ficheiros ou confirma mensagens. O envelope contém records, com uma a 100 entradas. Cada entrada tem recordId único e payload. Um envelope ambíguo interrompe o programa; problemas nos dados seguem para quarentena. Esta diferença evita emitir um plano aparentemente válido quando nem sequer é possível referenciar as entradas de forma única. O payload exige exatamente eventId, amount e currency. Os identificadores usam caracteres ASCII alfanuméricos, hífen ou underscore, até 64 caracteres. amount é texto com duas casas decimais, sem espaços, sinal positivo ou zeros iniciais redundantes; admite sinal negativo e até nove algarismos antes do ponto. currency aceita apenas EUR ou USD neste contrato didático. O programa converte o valor para cêntimos inteiros sem arredondar. Estes são limites locais, não limites de Google Cloud. Executa python3 content/labs/pde-quarantine-gate/run.py < content/labs/pde-quarantine-gate/case.json na raiz do projeto. Prevê o resultado primeiro. O fixture tem duas cópias equivalentes de A, dois B em conflito, um C válido com um par inválido, um identificador inválido e um D válido. Explica cada classificação antes de consultar a saída.
5. Interpretar identidades e contagens
recordId identifica uma entrada de transporte; eventId identifica um evento lógico imutável neste contrato. Duas entradas com o mesmo evento e valores normalizados iguais originam um evento aceite, mantendo ambas as referências. A normalização considera 0.00 e -0.00 equivalentes. Não normaliza vírgulas, notação científica ou moedas em minúsculas. Se os valores válidos diferem, todo o grupo fica em conflito, sem escolher o primeiro nem o maior. Existe outra regra deliberadamente conservadora: se uma entrada tem eventId válido mas outro campo inválido, os pares válidos desse evento recebem peer-invalid. O modelo não escolhe uma versão autoritativa. No fixture, C fica integralmente pendente, apesar de uma das entradas ter um montante válido. Esta política deve ser discutida com quem define o contrato; não é uma regra universal para todos os pipelines. O resultado é: dois eventos aceites representam três entradas e cinco entradas ficam em quarentena. Três mais cinco conserva as oito entradas, enquanto dois conta resultados lógicos. A lista duplicates explica a diferença sem classificar uma cópia equivalente como inválida. Altera a ordem do lote e confirma que o resultado permanece igual. Depois muda um montante e identifica que contagens e razões se alteram.
6. Conhecer os limites do exercício
A saída inclui referências e razões, sem copiar o payload rejeitado. Isto reduz a informação exposta pelo relatório, mas não prova anonimização: os próprios identificadores podem precisar de proteção num sistema real. O modelo não implementa retenção, controlo de acesso ou recuperação. rawPayloadIncluded, durablyStored, messagesAcknowledged e productionApproval permanecem falsos. Não confundas a palavra quarantined com uma confirmação de escrita num destino duradouro. A comparação só conhece o lote recebido. Se A:1.00 aparece num processo e A:2.00 noutro, cada execução pode aceitar o que vê. Para detetar o conflito global é preciso definir estado partilhado, fronteira temporal, concorrência, atualizações consistentes e tratamento de dados tardios. Aumentar o limite de entradas ou ordenar cada lote não resolve essa ausência de contexto. Constrói três contraexemplos: um campo adicional não autorizado, um recordId duplicado e um evento com duas moedas. Explica por que o primeiro pode ser uma rejeição de dados, o segundo interrompe o envelope e o terceiro é um conflito. Depois descreve uma arquitetura real que persista os resultados sem prometer atomicidade entre serviços que não a fornecem. Identifica onde uma falha exigiria reconciliação.
7. Diagnosticar o arranque e a pressão externa
Localiza a fase em que o pipeline falha. Num Flex Template, o programa que constrói o pipeline tem de terminar para permitir o lançamento. Um wait_until_finish nesse programa pode bloquear o arranque. Verifica também o entrypoint da imagem e as dependências de rede do launcher. A documentação atual descreve uma imagem de logging obtida de gcr.io, mesmo quando a imagem própria está em Artifact Registry; acesso a uma imagem não prova acesso a todas. Compara um template de referência usando a mesma identidade e rede, mas trata o resultado como evidência de diagnóstico, não como prova universal. Pré-instalar dependências testadas numa imagem pode reduzir instalações no arranque e evitar repositórios públicos indisponíveis. Continua a ser necessário validar acesso à imagem, compatibilidade de bibliotecas e configuração efetivamente promovida. Depois do arranque, mede throughput útil e rejeições. Mais workers podem aumentar 429 numa API limitada. Controla concorrência, tamanho de batches e retries de acordo com a capacidade do destino. Num enriquecimento por join, uma referência demasiado grande pode tornar side inputs ineficientes; avalia distribuição por chave e o custo de shuffle. Cada ajuste precisa de resultados comparáveis, não apenas de utilização de CPU.
8. Entregar um procedimento utilizável a RUN
O handover deve ligar sinais a decisões. Define quem observa crescimento de rejeições, como distingue erro do produtor de indisponibilidade do destino e onde encontra a versão das regras. Prepara amostras fictícias conhecidas: válida, duplicada, conflituosa e inválida. A equipa que assume suporte deve prever os resultados e repetir o exercício sem depender do autor do pipeline. Para uma recuperação real, fixa o conjunto de entradas, a versão do parser, o destino e os efeitos externos permitidos. Uma correção aprovada não deve apagar a evidência anterior. Confirma o armazenamento recuperável antes de depender da quarentena e reconcilia resultados após o replay. Se o pedido de criação perder a resposta, procura o estado existente antes de criar outro lote com uma referência diferente. Como avaliação final desta aula, prepara uma nota de decisão para o fixture: quais os eventos aceites, quais as cinco entradas pendentes, que evidência falta e quem decide a publicação? Acrescenta uma falha do destino de rejeições e explica como muda a decisão. O resumo é simples: contabilizar, classificar, preservar evidência, corrigir com autorização e verificar a entrega. Relaciona estas etapas com contratos, replay e capacidade das aulas anteriores.
"""Bounded fictional batch gate; no storage, cloud calls or acknowledgement."""
import json
import re
import sys
IDENTIFIER = re.compile(r'[A-Za-z0-9_-]{1,64}', re.ASCII)
AMOUNT = re.compile(r'-?(?:0|[1-9][0-9]{0,8})\.[0-9]{2}', re.ASCII)
def identifier(value):
return isinstance(value, str) and IDENTIFIER.fullmatch(value) is not None
def inspect(payload):
issues = []
if not isinstance(payload, dict):
return None, None, ['payload-shape']
event = payload.get('eventId')
if not identifier(event):
event = None
issues.append('event-id')
if set(payload) != {'eventId', 'amount', 'currency'}:
issues.append('payload-fields')
amount = payload.get('amount')
cents = None
if not isinstance(amount, str) or AMOUNT.fullmatch(amount) is None:
issues.append('amount-format')
else:
negative = amount.startswith('-')
major, minor = amount.lstrip('-').split('.')
cents = (int(major) * 100 + int(minor)) * (-1 if negative else 1)
currency = payload.get('currency')
if not isinstance(currency, str) or currency not in ('EUR', 'USD'):
issues.append('currency')
return event, (cents, currency) if not issues else None, sorted(issues)
def evaluate(data):
if not isinstance(data, dict) or set(data) != {'records'}:
raise ValueError('expected records envelope')
records = data['records']
if not isinstance(records, list) or not 1 <= len(records) <= 100:
raise ValueError('expected 1..100 records')
groups, rejected, ids = {}, [], set()
for record in records:
if not isinstance(record, dict) or set(record) != {'recordId', 'payload'}:
raise ValueError('invalid record envelope')
rid = record['recordId']
if not identifier(rid) or rid in ids:
raise ValueError('invalid or duplicate recordId')
ids.add(rid)
event, value, reasons = inspect(record['payload'])
row = {'recordId': rid, 'eventId': event, 'reasons': reasons}
if event is None:
rejected.append(row)
else:
groups.setdefault(event, []).append((row, value))
accepted, duplicates = [], []
for event, group in sorted(groups.items()):
invalid_peer = any(row['reasons'] for row, _ in group)
values = {value for _, value in group if value is not None}
conflict = len(values) > 1
if invalid_peer or conflict:
for row, _ in group:
reasons = set(row['reasons'])
if invalid_peer and not reasons:
reasons.add('peer-invalid')
if conflict:
reasons.add('event-conflict')
rejected.append({**row, 'reasons': sorted(reasons)})
else:
cents, currency = next(iter(values))
source_ids = sorted(row['recordId'] for row, _ in group)
accepted.append({'eventId': event, 'cents': cents,
'currency': currency, 'recordIds': source_ids})
if len(source_ids) > 1:
duplicates.append({'eventId': event, 'recordIds': source_ids,
'collapsed': len(source_ids) - 1})
rejected.sort(key=lambda row: row['recordId'])
return {'accepted': accepted, 'quarantined': rejected,
'duplicates': duplicates, 'inputRecords': len(records),
'acceptedSourceRecords': sum(len(x['recordIds']) for x in accepted),
'quarantinedRecords': len(rejected),
'rawPayloadIncluded': False, 'durablyStored': False,
'messagesAcknowledged': False, 'productionApproval': False}
def unique_object(pairs):
result = {}
for key, value in pairs:
if key in result:
raise ValueError('duplicate JSON key')
result[key] = value
return result
if __name__ == '__main__':
try:
raw = sys.stdin.read(1000001)
if len(raw) > 1000000:
raise ValueError('input too large')
print(json.dumps(evaluate(json.loads(raw, object_pairs_hook=unique_object)), sort_keys=True))
except (ValueError, TypeError, RecursionError) as error:
print('invalid input: ' + str(error), file=sys.stderr)
sys.exit(2)
Oito entradas originam dois eventos aceites a partir de três entradas e cinco rejeições; nenhuma escrita duradoura é executada.
Armadilhas comuns
Backlog vazio como validação completa; retry como correção de dados; zero como substituição; classificação como persistência; mais workers como resposta a 429.
Tópicos relacionados: Contratos e consumidores · Replay e efeitos externos · Capacidade e recuperação
Um processo verde não prova uma entrega completa: conserva contagens, razões e evidência recuperável.
Referência: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)