Pular para o conteúdo principal

Bloqueio e reconhecimento de mensagens

Os SDKs do Zerobus Ingest oferecem vários métodos para ingerir um registro, que equilibram o throughput com a quantidade de confirmação de durabilidade que você recebe de volta. Esta página explica cada método e quando bloquear quanto à durabilidade. Para reagir a confirmações de forma assíncrona em vez de bloquear, consulte Callbacks de confirmação.

Os exemplos nesta página usam o SDK do Python. Para obter o tempo limite exato e as opções de configuração que cada método aceita (incluindo seus defaults e unidades), consulte o repository do Zerobus SDK. Os SDKs de outras linguagens expõem opções equivalentes.

O que é um deslocamento?

Cada registro que você ingere recebe um offset: sua posição na transmissão. O offset é como você se refere a um registro específico quando deseja confirmar que ele foi gravado de forma durável. O Zerobus Ingest oferece garantias de entrega pelo menos uma vez, e aguardar um offset é a forma como um cliente confirma essa garantia para um determinado registro.

Confirmar um offset significa que o registro é durável, não que ele já possa ser consultado na tabela Delta. O Zerobus Ingest materializa registros duráveis na tabela logo em seguida. Para números de latência, consulte Latência.

Métodos de ingestão

Os SDKs oferecem duas maneiras de ingerir um registro. (Os nomes dos métodos abaixo são do SDK para Python. Outros SDKs expõem métodos equivalentes.)

Método

Devolve

Use-o quando

Baseado em deslocamento , ingest_record_offset()

O offset do registro, após o registro ser enfileirado na transmissão.

Default recomendado. Você deseja enfileirar registros em ordem e, opcionalmente, confirmar a durabilidade mais tarde aguardando em um offset.

Baseado em futuro , ingest_record()

Um(a) RecordAcknowledgment no(a) qual você pode aguardar.

Descontinuado. Prefira baseado em offset para melhor desempenho.

Método

Devolve

Use-o quando

Baseado em deslocamento , ingest_record_offset()

O offset do registro, após o registro ser enfileirado na transmissão.

Default recomendado. Você deseja enfileirar registros em ordem e, opcionalmente, confirmar a durabilidade mais tarde aguardando em um offset.

Baseado em futuro , ingest_record()

Um(a) RecordAcknowledgment no(a) qual você pode aguardar.

Descontinuado. Prefira baseado em offset para melhor desempenho.

Baseado em deslocamento (recomendado)

ingest_record_offset() envia o registro e retorna seu deslocamento assim que o registro é enfileirado na transmissão. A chamada é executada em sua thread de chamada, portanto, os registros são enfileirados na ordem em que você chama o método, e o deslocamento retornado permite que você confirme a durabilidade posteriormente com wait_for_offset(). Este é o default recomendado para a maioria dos produtores, e é o método usado nos exemplos Use Zerobus Ingest.

Baseado em futuro (descontinuado)

ingest_record() retorna um objeto RecordAcknowledgment no qual você pode aguardar a durabilidade. Ele está obsoleto em favor do método baseado em deslocamento, que apresenta melhor desempenho. Use-o apenas para código existente que ainda não foi migrado.

Ingestão registro a registro vs. em lotes

Cada método de ingestão tem uma variante em lote (por exemplo, ingest_records_offset()) que envia uma lista de registros em uma única chamada. O processamento em lote é mais eficiente do que chamadas individuais para ingestão em massa.

Para JSON e Protocol Buffers (protobuf), um lote faz commit atomicamente: ou todos os registros no lote são aceitos e tornados duráveis, ou o lote inteiro é rejeitado. O Zerobus Ingest não realiza uploads parciais ou confirmações parciais para esses formatos, portanto, sua tabela nunca contém um lote parcial. Um lote que falha na validação (por exemplo, uma incompatibilidade de esquema) falha rapidamente, antes de tocar na tabela, em vez de gravar alguns registros e descartar outros.

Como um lote JSON ou protobuf é enviado como uma única mensagem, o tamanho máximo de mensagem de 10 MB se aplica tanto a um único registro quanto a um lote inteiro: todos os registros em um lote devem caber juntos dentro de 10 MB. Dimensione seus lotes para permanecer abaixo desse limite. Consulte Tamanho do registro.

Os lotes do Arrow Flight são a exceção

A ingestão do Apache Arrow Flight não segue o modelo de mensagem única tudo ou nada acima. Um lote do Arrow pode ser muito maior do que um lote JSON ou Protobuf, e o caminho do Arrow Flight divide um lote grande em mensagens de transporte menores que são enviadas e confirmadas individualmente, em vez de como uma unidade atômica. Como resultado:

  • O limite de 10 MB por mensagem que se aplica a lotes JSON e protobuf não se aplica a um lote Arrow da mesma maneira. Um lote Arrow grande é dividido em mensagens de transporte em vez de ser rejeitado por tamanho.
  • A durabilidade é confirmada na granularidade da mensagem de transporte, portanto, um lote lógico muito grande pode ser parcialmente durável se ocorrer uma falha durante o processo, em vez de commit tudo ou nada.

ingest_batch() ainda retorna um único offset lógico para o lote que você enviou, e wait_for_offset() nesse offset só é concluído após cada mensagem de transporte que compõe o lote ter sido confirmada. Para o modelo completo do Arrow Flight, orientações sobre lotes e recuperação de dados não confirmados, consulte Use o Arrow Flight com o Zerobus Ingest.

Quando você deve bloquear em uma mensagem?

O bloqueio em um offset troca throughput por uma garantia de durabilidade por registro mais forte em seu código cliente. Escolha com base na sua carga de trabalho:

  • Não bloquear : o default correto para transmissão de alto volume, onde você se preocupa com o throughput sustentado e pode confirmar a durabilidade em agregado (por exemplo, no fechamento da transmissão ou por meio de um callback de confirmação). A maioria dos produtores deve começar aqui.
  • Bloquear em um offset : considere isso quando sua aplicação precisar saber que um registro específico é durável antes de realizar outra ação. Por exemplo:
    • Você está prestes a excluir ou confirmar a fonte dos dados (uma mensagem de fila, um arquivo, um cursor upstream) e não deve perdê-la se a ingestão falhar.
    • Você está ingerindo em pontos de verificação ou limites transacionais e precisa que cada ponto de verificação seja durável antes de avançar.
    • Trata-se de gravações de baixo volume e alto valor, nas quais a confirmação por registro é mais importante do que o throughput.

Não bloqueie a cada registro em um loop de alto throughput. Isso serializa seu produtor em um round-trip para o servidor para cada registro e reduz drasticamente o throughput. Em vez disso, o Databricks recomenda ingerir uma grande parte de registros e, em seguida, confirmar a durabilidade uma vez para toda a parte. Você tem duas maneiras de fazer isso: aguardar o último deslocamento ou liberar a transmissão. O bloqueio por registro individual deve ser reservado para os casos específicos acima, onde um único registro deve ser confirmado antes da próxima ação.

Aguardar um offset

wait_for_offset() bloqueia até que o Zerobus Ingest confirme que o registro naquele offset foi gravado de forma durável, ou até que o tempo limite seja atingido. Use-o para confirmar um ponto específico na transmissão, mais comumente o último registro de um bloco. Ingira o bloco, mantenha o offset final que o loop retorna e aguarde nesse offset em vez de aguardar após cada registro:

Python
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

table_properties = TableProperties("main.default.air_quality")
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

try:
last_offset = 0
for row in records:
last_offset = stream.ingest_record_offset(row)

# Block until everything up to the last record of the chunk is durable
stream.wait_for_offset(last_offset)
print("Chunk durably written.")
finally:
stream.close()

Liberar a transmissão

flush() bloqueia até que todos os registros que você ingeriu até o momento sejam gravados de forma durável e, em seguida, retorna. Ao contrário de wait_for_offset(), você não rastreia um deslocamento: o flush aguarda tudo o que estiver pendente na transmissão. Isso não fecha a transmissão, portanto, você pode continuar a ingestão posteriormente.

Python
try:
for row in records:
stream.ingest_record_offset(row)

# Block until every pending record is durable
stream.flush()
print("All ingested records durably written.")
finally:
stream.close()

wait_for_offset vs. flush

Ambos confirmam a durabilidade de um fragmento. Escolha com base no que você está confirmando:

  • Use wait_for_offset(offset) quando quiser confirmar até um registro específico, por exemplo, um limite de checkpoint, enquanto outros registros ainda podem estar em trânsito atrás dele.
  • Use flush() quando quiser confirmar que todos os registros pendentes estão duráveis antes de prosseguir, por exemplo, ao final de um lote, antes de avançar um cursor upstream ou antes de desligar. flush() é regido por um limite de tempo de liberação configurável.

close() libera e fecha a transmissão, para que os registros sejam sempre tornados duráveis em um desligamento normal. Sempre chame-o em um bloco finally.

Reagir a confirmações de forma assíncrona

Se, em vez de bloquear, você desejar reagir às confirmações de durabilidade e aos erros à medida que chegam, enquanto seu produtor continua enviando na velocidade máxima, registre um callback de reconhecimento na transmissão. Os callbacks são um recurso separado das chamadas de bloqueio nesta página. Veja Callbacks de reconhecimento.

Relacionado