Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Os SDKs de ingestão do Zerobus oferecem vários métodos para fazer a ingestão de um registro, que compensam a taxa de transferência e o nível de confirmação de durabilidade recebido. Esta página explica cada método e quando bloquear a durabilidade. Para reagir a confirmações de forma assíncrona em vez de bloquear, veja Chamadas de retorno de reconhecimento.
Os exemplos nesta página usam o SDK Python. Para as opções exatas de timeout e configuração que cada método aceita (incluindo seus padrões e unidades), veja o repositório Zerobus SDK. Os outros SDKs de linguagem expõem opções equivalentes.
O que é um offset?
Cada registro que você ingere recebe um offset: sua posição no fluxo. O desvio é a forma de referenciar um registro específico quando você deseja confirmar que ele foi gravado de forma durável. O Zerobus Ingest oferece garantias de entrega pelo menos uma vez, e aguardar por um desvio é como o cliente confirma essa garantia para um registro específico.
Confirmar um desvio significa que o registro é durável, não que ele ainda possa ser consultado na tabela Delta. O Zerobus Ingest materializa registros duráveis na tabela logo depois. Para números de latência, veja Latência.
Métodos de ingestão
Os SDKs oferecem duas formas de ingerir um registro. (Os nomes dos métodos abaixo são do SDK Python. Outros SDKs expõem métodos equivalentes.)
| Método | Returns | Use-o quando |
|---|---|---|
Baseado em deslocamento, ingest_record_offset() |
O deslocamento do registro, depois que o registro é enfileirado no fluxo. | Padrão recomendado. Você quer enfileirar registros em ordem e, opcionalmente, confirmar a durabilidade depois, esperando por um deslocamento. |
Baseado no futuro, ingest_record() |
Um RecordAcknowledgment em que você pode confiar. |
Preterido. Prefiro baseado em deslocamento para melhor desempenho. |
Baseado em offset (recomendado)
ingest_record_offset() envia o registro e retorna seu deslocamento assim que o registro é enfileirado na transmissão. A chamada é executada na sua thread de chamada; portanto, os registros são enfileirados na ordem em que você chama o método, e o deslocamento retornado permite confirmar a durabilidade posteriormente com wait_for_offset(). Esse é o padrão recomendado para a maioria dos produtores, e é o método usado nos exemplos de Use Zerobus Ingest .
Baseado no futuro (obsoleto)
ingest_record() retorna um objeto RecordAcknowledgment no qual você pode aguardar a confirmação de durabilidade. Foi descontinuado em favor do método baseado em deslocamento, que oferece melhor desempenho. Use apenas para código existente que ainda não foi migrado.
Registro a registro vs. ingestão por lote
Cada método de ingestão possui uma variante em lote (por exemplo, ingest_records_offset()) que envia uma lista de registros em uma chamada. O processamento em lote é mais eficiente do que chamadas individuais para ingestão em massa.
Para JSON e Protocol Buffers (protobuf), a confirmação de um lote é atômica: ou todos os registros do lote são aceitos e persistidos de forma durável, ou o lote inteiro é rejeitado. O Zerobus Ingest não realiza uploads parciais nem reconhecimento parcial para esses formatos, então sua tabela nunca contém um lote parcial. Um lote que falha na validação (por exemplo, por uma incompatibilidade de esquema) falha imediatamente, antes de afetar a 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 juntos devem caber dentro de 10 MB. Ajuste o tamanho dos lotes para ficar abaixo desse limite. Veja Tamanho do registro.
Os lotes do Arrow Flight são exceção
A ingestão do Apache Arrow Flight não segue o modelo tudo ou nada de mensagem única acima. Um lote do Arrow pode ser muito maior do que um lote JSON ou protobuf, e o Arrow Flight divide um lote grande em mensagens de transporte menores, que são enviadas e confirmadas individualmente, em vez de serem tratadas como uma única 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 forma. Um grande lote Arrow é dividido em mensagens de transporte, em vez de ser rejeitado devido ao tamanho.
- A durabilidade é confirmada na granularidade da mensagem de transporte; então, um lote lógico muito grande poderá ser parcialmente durável se uma falha ocorrer no meio do caminho, em vez de comprometer tudo ou nada.
ingest_batch() continua retornando um único deslocamento lógico para o lote que você enviou, e wait_for_offset() nesse deslocamento só será concluído depois que cada mensagem de transporte que compõe o lote for confirmada. Para o modelo completo do Arrow Flight, as orientações sobre processamento em lote e a recuperação de dados não reconhecidos, consulte Usar o Arrow Flight com o Zerobus Ingest.
Quando você deve bloquear uma mensagem?
Bloquear um deslocamento compensa uma taxa de transferência por uma garantia de durabilidade por registro mais forte no código do cliente. Escolha com base na sua carga de trabalho:
- Não bloqueie: o padrão certo para streaming de alto volume, onde você se importa com o throughput sustentado e pode confirmar a durabilidade no total (por exemplo, no encerramento do stream ou por meio de um callback de reconhecimento). A maioria dos produtores deveria começar por aqui.
-
Bloquear em um deslocamento: considere esta opção quando seu aplicativo precisar saber que um registro específico foi gravado de forma durável antes de executar outra ação. Por exemplo:
- Você está prestes a excluir ou reconhecer a origem dos dados (uma mensagem de fila, um arquivo, um cursor upstream) e não deve perdê-los caso a ingestão falhe.
- Você está ingerindo pontos de verificação ou limites transacionais e precisa que cada ponto de verificação seja durável antes de avançar.
- Você está fazendo gravações de baixo volume e alto valor onde a confirmação por registro importa mais do que a taxa de transferência.
Não bloqueie para cada registro em um loop de alta taxa de transferência. Isso faz com que seu produtor opere de forma serial, com uma ida e volta ao servidor para cada registro, e reduz drasticamente a taxa de transferência. Em vez disso, o Azure Databricks recomenda ingerir uma grande parte dos registros e depois confirmar a durabilidade uma vez para todo o bloco. Você tem duas formas de fazer isso: esperar o último offset ou lavar o fluxo. 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.
Aguarde o deslocamento
wait_for_offset() bloqueia até que o Zerobus Ingest confirme que o registro nesse deslocamento foi gravado de forma durável, ou até ocorrer o tempo limite. Use-o para confirmar um ponto específico na transmissão, mais comumente o último registro de um bloco. Ingira a parte, mantenha o deslocamento final que o loop retorna e aguarde esse único deslocamento em vez de esperar após cada registro:
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()
Libere o fluxo
flush() bloqueia até que todos os registros que você ingeriu até agora sejam gravados de forma durável e, depois, retorna. Ao contrário de wait_for_offset(), você não rastreia um deslocamento: a liberação aguarda tudo o que está pendente no fluxo. Ela não fecha o fluxo; então, você pode continuar a ingestão depois.
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 para uma parte. Escolha com base no que você está confirmando:
- Use
wait_for_offset(offset)quando desejar confirmar até um registro específico, por exemplo, um limite de ponto de verificação, enquanto outros registros ainda podem estar em andamento atrás dele. - Use
flush()quando desejar confirmar que todos os registros pendentes são duráveis antes de avançar, por exemplo, no final de um lote, antes de avançar um cursor upstream ou antes de desligar.flush()é controlado por um tempo limite de flush configurável.
close() descarrega e fecha o fluxo, para que os registros sejam sempre gravados de forma durável em um desligamento normal. Sempre chame em um bloco finally.
Reagir aos reconhecimentos de forma assíncrona
Se, em vez de bloquear, você desejar reagir às confirmações e erros de durabilidade conforme eles chegam, enquanto seu produtor continua pressionando em velocidade máxima, registre um retorno de chamada de reconhecimento na transmissão. Os retornos de chamada são um recurso separado das chamadas de bloqueio nesta página. Confira Retornos de chamada de confirmação.
Related
- Retornos de chamada de confirmação: reaja a confirmações e erros de forma assíncrona.
- Use Zerobus Ingest: Escreva um cliente.
- Streams: Fluxos e deslocamentos.
- Tratamento de erros Zerobus Ingest: Tratamento de erros.