Pular para o conteúdo principal

Padrões de recuperação e nova tentativa

Esta análise detalhada explica como criar um cliente Zerobus Ingest resiliente no Lakeflow Connect. Ele aborda a recuperação integrada do SDK, como os erros aparecem e como resgatar registros não confirmados quando uma transmissão falha permanentemente. Os nomes de métodos e opções abaixo são do SDK do Python. Outros SDKs expõem equivalentes.

Recuperação integrada

Os SDKs do Zerobus se recuperam automaticamente de falhas transitórias. A recuperação é acionada quando a transmissão encontra um erro recuperável, normalmente um tempo limite ou uma interrupção de rede, ou quando a transmissão recebe um sinal de desligamento normal. A recuperação está ativada por default, e você a ajusta por meio das opções de configuração de transmissão:

Opção

Descrição

recovery

Ativar a recuperação automática de transmissão.

recovery_timeout_ms

Tempo esgotado para uma operação de recuperação.

recovery_backoff_ms

Atraso entre tentativas de recuperação.

recovery_retries

Número máximo de tentativas de recuperação.

Opção

Descrição

recovery

Ativar a recuperação automática de transmissão.

recovery_timeout_ms

Tempo esgotado para uma operação de recuperação.

recovery_backoff_ms

Atraso entre tentativas de recuperação.

recovery_retries

Número máximo de tentativas de recuperação.

Para a maioria das cargas de trabalho, os padrões são um bom ponto de partida, e você não precisa escrever seu próprio loop de reconexão para problemas transitórios. Ao usar um SDK, os tokens OAuth também são atualizados automaticamente na criação e recuperação da transmissão, portanto, não há nada para seu cliente gerenciar. A exceção é a API REST, na qual seu cliente busca e refresh o token OAuth por conta própria. Consulte Use o Zerobus Ingest. Para os valores e unidades default dessas opções, consulte o repository do Zerobus SDK.

Erros e tentativas de repetição

O SDK tenta repetir automaticamente erros transitórios, como problemas de rede ou erros temporários de servidor, por meio de sua recuperação integrada. Falhas das quais não é possível se recuperar, como credenciais inválidas ou uma tabela ausente, aparecem como ZerobusException. Capture ZerobusException para tratar uma falha e, em seguida, decida se deve corrigir a causa subjacente, recuperar em uma nova transmissão ou parar.

Python
from zerobus.sdk.shared import ZerobusException

try:
stream.ingest_record_offset(record)
except ZerobusException as e:
# Handle the failure: log it, fix the cause, recover on a new stream, or stop.
...

Para a lista completa de códigos de erro, veja tratamento de erros do Zerobus Ingest.

Recuperando registros não confirmados

Quando uma transmissão falha permanentemente, após a recuperação automática do SDK ser esgotada, os registros que foram enviados, mas ainda não confirmados pelo servidor, permanecem retidos pelo cliente no buffer em trânsito descrito em Comunicação assíncrona. Recupere-os para não perder dados:

  • get_unacked_records() retorna os registros não confirmados como bytes brutos.
  • get_unacked_batches() retorna lotes não confirmados (cada um uma lista de registros) para a lógica de repetição de lotes.

Os registros retornam em sua forma serializada: decodifique JSON com json.loads(record.decode('utf-8')) ou desserialize Protocol Buffers (protobuf) com seu tipo de mensagem. Salve-os ou reproduza-os em uma nova transmissão.

Recuperar uma transmissão após falha permanente

O SDK lida automaticamente com novas tentativas para erros transitórios. Falhas de enfileiramento, liberação e fechamento aparecem como ZerobusException. get_unacked_records() e recreate_stream() são bem-sucedidos apenas após a transmissão ter sido fechada, o que uma falha terminal faz. Uma falha de enfileiramento deixa a transmissão ativa, portanto essas chamadas falham; nesse caso, gere o erro original e mantenha a transmissão. recreate_stream() re-enfileira os registros que já foram aceitos; ele não tenta novamente um payload que falhou ao ser enfileirado.

Python
from zerobus.sdk.shared import ZerobusException

try:
for i in range(10000):
stream.ingest_record_offset(record)
stream.flush()
except ZerobusException as e:
print(f"Ingestion failed: {e}")
try:
unacked = list(stream.get_unacked_records())
except ZerobusException:
raise e
print(f"{len(unacked)} previously queued records were unacknowledged.")
try:
new_stream = sdk.recreate_stream(stream)
try:
new_stream.flush()
finally:
new_stream.close()
except ZerobusException:
raise e
else:
stream.close()

Use get_unacked_batches() para inspecionar o agrupamento de lotes original após o fechamento da transmissão:

Python
unacked_batches = list(stream.get_unacked_batches())
print(f"{len(unacked_batches)} batches remain unacknowledged")

Tratamento de duplicatas na reprodução

O Zerobus Ingest fornece entrega pelo menos uma vez, não exatamente uma vez, portanto, a reprodução de registros resgatados, ou qualquer repetição, pode gravar um registro mais de uma vez. Se sua carga de trabalho não puder tolerar duplicatas, deduplique os dados no lakehouse:

  • Inclua um identificador exclusivo estável em cada registro (por exemplo, um ID de evento atribuído pela origem ou uma key natural).
  • Deduplique na leitura ou durante o processamento downstream. Por exemplo, use um MERGE INTO que corresponda ao identificador, ou um ROW_NUMBER() com janela em uma transformação.

Como os registros em uma transmissão são confirmados em ordem, um número de sequência crescente de forma monótona também funciona bem como a key de desduplicação.

Liberação e fechamento normal

  • flush() aguarda que o servidor reconheça os registros que você enviou como duráveis, sem fechar a transmissão. Chame-o quando precisar de um ponto de verificação de durabilidade no meio da transmissão.
  • close() libera e fecha a transmissão de forma elegante, aguardando que os registros pendentes sejam reconhecidos como duráveis antes de retornar. Utilize-o para um desligamento elegante, não para se recuperar de uma falha na transmissão. Quando uma transmissão falhar, resgate os registros não reconhecidos, conforme mostrado em Recuperar uma transmissão após falha permanente.

Para confirmar a durabilidade de um registro específico em vez de toda a transmissão, consulte Bloqueio e confirmação de mensagens.

Recuperando dados do local de fallback durável

Se uma alteração interruptiva for feita na sua tabela de destino após o Zerobus Ingest ter tornado seus dados duráveis, mas antes que ele possa publicar, o Zerobus Ingest gravará esses dados como arquivos Parquet em um diretório de fallback na raiz de armazenamento da sua tabela, em vez de descartá-los. Consulte Local de fallback durável.

Como saber se os dados foram gravados lá: o diretório de fallback é _zerobus/table_rejected_parquets/, relativo ao local de armazenamento raiz físico da tabela. Se a ingestão continuou, mas faltam linhas na tabela após uma alteração na tabela, verifique esse diretório em busca de arquivos Parquet.

Reprocessar dados de fallback na tabela

Assim que você corrigir a causa (normalmente alinhando o esquema da tabela com o que seus produtores enviam), reprocesse os arquivos Parquet de fallback na tabela de destino. Os arquivos são Parquet padrão no local de armazenamento da tabela, portanto, você pode carregá-los com COPY INTO:

  1. Resolva a incompatibilidade de esquema. Evolua a tabela de destino (ou o esquema do seu produtor) para que os registros de fallback caibam. Consulte Gerenciamento de esquema.

  2. Inspecione os dados de fallback antes de carregar. Aponte uma query para o caminho de fallback para confirmar o que está lá e se ele agora corresponde à tabela:

    SQL
    SELECT * FROM parquet.`<table-storage-root>/_zerobus/table_rejected_parquets/` LIMIT 10;
  3. Carregue os arquivos com COPY INTO. COPY INTO é idempotente: ele rastreia os arquivos que já carregou, portanto, executá-lo novamente não carregará em duplicidade os mesmos arquivos de fallback:

    SQL
    COPY INTO <catalog>.<schema>.<table>
    FROM '<table-storage-root>/_zerobus/table_rejected_parquets/'
    FILEFORMAT = PARQUET
    COPY_OPTIONS ('mergeSchema' = 'false');
  4. Verifique se as contagens de linhas esperadas foram gravadas e, após confirmar que os dados estão na tabela, limpe o diretório de fallback se não precisar mais dele.

Como o Zerobus Ingest é pelo menos uma vez, os registros que foram publicados na tabela e gravados no local de fallback podem ser carregados duas vezes. Se as duplicatas forem importantes, faça a desduplicação conforme descrito em Handling duplicates on replay. Para reprocessamento contínuo ou automatizado, você pode apontar o Auto Loader para o caminho de fallback em vez de executar COPY INTO manualmente.

nota

Este procedimento é uma abordagem geral de primeira passagem usando ferramentas padrão do Delta nos arquivos Parquet de fallback. Valide-o em relação à sua configuração de tabela e armazenamento antes de confiar nele em produção.

Relacionado