Bloqueio e confirmação de mensagens

Os SDKs de ingestão do Zerobus oferecem vários métodos para ingerir um registo, que estabelecem um compromisso entre a taxa de processamento e o nível de confirmação de durabilidade devolvido. Esta página explica cada método e quando bloquear na durabilidade. Para reagir a confirmações de receção de forma assíncrona em vez de bloquear, consulte funções de chamada de retorno para confirmações de receção.

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 os seus valores predefinidos e unidades), consulte o repositório Zerobus SDK. Os outros SDKs de linguagem expõem opções equivalentes.

O que é um desfasamento?

Cada registo que introduz recebe um offset: a sua posição no stream. O offset é a forma de referir um registo específico quando quer confirmar que foi gravado de forma persistente. O Zerobus Ingest oferece garantias de entrega pelo menos uma vez, e esperar por um offset é a forma como um cliente confirma essa garantia para um dado registo.

Confirmar um offset significa que o registo é durável, não que ainda seja consultável na tabela Delta. O Zerobus Ingest materializa registos duráveis na tabela pouco depois. Para valores de latência, veja Latência.

Métodos de ingestão

Os SDKs oferecem duas formas de ingerir um registo. (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 offset, ingest_record_offset() O deslocamento do disco, depois de o disco ser colocado na fila na stream. Padrão recomendado. Queres enfileirar registos por ordem e, opcionalmente, confirmar a durabilidade mais tarde, esperando por um offset.
Orientado para o futuro, ingest_record() Um(a) RecordAcknowledgment em que pode confiar. Deprecated. Prefira a opção baseada em offset para melhor desempenho.

Baseado em offset (recomendado)

ingest_record_offset() envia o registo e devolve o respetivo offset assim que o registo é colocado em fila no fluxo. A chamada é executada no seu encadeamento de chamada, pelo que os registos são colocados em fila pela ordem em que chama o método, e o offset devolvido permite-lhe confirmar posteriormente a durabilidade com wait_for_offset(). Esta é a opção predefinida recomendada para a maioria dos produtores e é o método utilizado nos exemplos de Usar a ingestão do Zerobus.

Baseado no futuro (obsoleto)

ingest_record() devolve um objeto RecordAcknowledgment sobre o qual pode esperar para garantir a durabilidade. Foi descontinuado em favor do método baseado em deslocamento, que oferece melhor desempenho. Usa-o apenas para código existente que ainda não foi migrado.

registo a registo vs. ingestão em lote

Cada método de ingestão tem uma variante em lote (por exemplo, ingest_records_offset()) que submete uma lista de registos numa chamada. O processamento em lotes é 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 registos do lote são aceites e persistidos de forma durável, ou o lote inteiro é rejeitado. O Zerobus Ingest não realiza uploads parciais nem reconhecimentos parciais para estes formatos, pelo que a sua tabela nunca contém um lote parcial. Um lote que não passa na validação (por exemplo, devido a uma incompatibilidade de esquema) falha imediatamente, antes de chegar a afetar a tabela, em vez de gravar alguns registos e descartar outros.

Como um lote JSON ou protobuf é enviado como uma única mensagem, o tamanho máximo de mensagem de 10 MB aplica-se tanto a um único registo como a um lote inteiro: todos os registos de um lote juntos devem caber dentro de 10 MB. Dimensiona os teus lotes para ficares abaixo desse limite. Ver Tamanho do registo.

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

A ingestão do Apache Arrow Flight não segue o modelo tudo ou nada de mensagem única acima. Um lote Arrow pode ser muito maior do que um lote JSON ou protobuf, e o percurso de voo Arrow divide um lote grande em mensagens de transporte menores que são enviadas e reconhecidas individualmente em vez de como uma unidade atómica única. 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 quanto ao tamanho.
  • A durabilidade é confirmada na granularidade da mensagem de transporte, pelo que um lote lógico muito grande pode ser parcialmente durável se ocorrer uma falha a meio, em vez de comprometer tudo ou nada.

ingest_batch() continua a devolver um único offset lógico para o lote que submeteu, e wait_for_offset() nesse offset só fica concluído depois de todas as mensagens de transporte que compõem o lote serem reconhecidas. Para o modelo completo do Arrow Flight, orientações sobre processamento em lotes e recuperação de dados não reconhecidos, consulte Utilizar o Arrow Flight com o Zerobus Ingest.

Quando deve bloquear uma mensagem?

Bloquear o throughput de uma troca offset para garantir uma durabilidade por registo 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 se preocupa com a taxa de transferência sustentada e pode confirmar a durabilidade no agregado (por exemplo, no encerramento do stream ou através de um callback de confirmação). A maioria dos produtores deve começar por aqui.
  • Bloquear num offset: considere esta opção quando a sua aplicação precisa de saber que um registo específico foi persistido de forma durável antes de executar outra ação. Por exemplo:
    • Está prestes a apagar ou reconhecer a origem dos dados (uma mensagem de fila, um ficheiro, um cursor a montante) e não deve perdê-la se a ingestão falhar.
    • Está a ingerir pontos de controlo ou limites transacionais e precisa que cada ponto de controlo seja durável antes de avançar.
    • Está a fazer escritas de baixo volume e alto valor, onde a confirmação por registo importa mais do que o rendimento.

Não bloqueie em todos os registos num ciclo de alta produtividade. Isso serializa o seu produtor numa viagem de ida e volta ao servidor para cada disco e reduz drasticamente o rendimento. Em vez disso, o Azure Databricks recomenda ingerir um grande conjunto de registos e, em seguida, confirmar a persistência uma única vez para todo o conjunto. Tens duas formas de o fazer: esperar pelo último deslocamento ou lavar o fluxo. O bloqueio por registo individual deve ser reservado para os casos específicos acima em que um único registo deve ser confirmado antes da ação seguinte.

Espera por um deslocamento

wait_for_offset() bloqueia até que o Zerobus Ingest confirme que o registo nesse offset foi escrito de forma duradoura, ou até que seja atingido o tempo limite. Utilize-o para confirmar um ponto específico no stream, mais frequentemente o último registo de um segmento. Processe o bloco, mantenha o offset final que o ciclo devolve e aguarde esse offset em vez de aguardar após cada registo:

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()

Purgar o fluxo

flush() bloqueia até que todos os registos ingeridos até ao momento sejam gravados de forma persistente e, em seguida, retorna. Ao contrário de wait_for_offset(), não rastreias um desfasamento: o flush aguarda tudo o que estiver pendente no fluxo. Não fecha o fluxo, por isso pode continuar a ingerir 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 atestam a durabilidade de um bloco. Escolha com base no que está a confirmar:

  • Use wait_for_offset(offset) quando quiser confirmar até um registo específico, por exemplo, um limite de ponto de controlo, enquanto outros registos ainda podem estar em voo atrás dele.
  • Use flush() quando quiser confirmar que todos os registos pendentes são duráveis antes de avançar, por exemplo no final de um lote, antes de avançar um cursor a montante, ou antes de desligar. flush() é controlado por um tempo limite de esvaziamento configurável.

close() esvazia e fecha o stream, pelo que os registos ficam sempre persistidos num encerramento normal. Chame-o sempre num bloco finally.

Responder a confirmações de forma assíncrona

Se, em vez de bloquear, quiseres reagir às confirmações e erros de durabilidade à medida que chegam, enquanto o teu produtor continua a pressionar a toda a velocidade, regista uma chamada de confirmação na transmissão. Os callbacks são uma funcionalidade separada das chamadas de bloqueio nesta página. Ver chamadas de reconhecimento.