Bloqueo y acuse de recibo de mensajes

Los SDK de ingestión de Zerobus ofrecen varios métodos para incorporar un registro, que equilibran el rendimiento con el grado de confirmación de durabilidad que se obtiene. Esta página explica cada método y cuándo bloquear la durabilidad. Para reaccionar a los acuses de recibo de forma asíncrona en lugar de bloquear, véase Llamadas de reconocimiento.

Los ejemplos de esta página utilizan el SDK de Python. Para las opciones exactas de tiempo de espera y configuración que acepta cada método (incluyendo sus valores predeterminados y unidades), consulta el repositorio Zerobus SDK. Los otros SDKs del lenguaje exponen opciones equivalentes.

¿Qué es un desplazamiento?

A cada registro que ingieres se le asigna un desplazamiento: su posición en el flujo. El offset es la forma de referirse a un registro específico cuando quieres confirmar que se escribió de forma persistente. Zerobus Ingest ofrece garantías de entrega al menos una vez, y esperar a un offset es la forma en que un cliente confirma esa garantía para un registro determinado.

Confirmar un offset significa que el registro es persistente, no que todavía se pueda consultar en la tabla Delta. Zerobus Ingest materializa registros duraderos en la tabla poco después. Para las cifras de latencia, véase Latencia.

Métodos de ingesta

Los SDKs ofrecen dos formas de ingerir un registro. (Los nombres de los métodos a continuación provienen del SDK de Python. Otros SDKs exponen métodos equivalentes.)

Método Returns Úselo cuando
Basado en desplazamiento, ingest_record_offset() El desplazamiento del disco es después de que el disco está en cola en la transmisión. Predeterminado (recomendado). Quieres poner los registros en cola en orden y, opcionalmente, confirmar la durabilidad más adelante esperando un desplazamiento.
Basado en el futuro, ingest_record() Alguien RecordAcknowledgment en quien puedes confiar. Deprecated. Prefiero basado en offset para mejorar el rendimiento.

Basado en desplazamiento (recomendado)

ingest_record_offset() envía el registro y devuelve su offset una vez que el registro está en cola en la transmisión. La llamada se ejecuta en el hilo desde el que llamas, por lo que los registros se encolan en el orden en que llamas al método, y el offset devuelto te permite confirmar la durabilidad más adelante con wait_for_offset(). Esta es la opción predeterminada recomendada para la mayoría de los productores, y es el método utilizado en los ejemplos de Uso de Zerobus Ingest.

Basado en el futuro (obsoleto)

ingest_record() Devuelve un RecordAcknowledgment objeto que puedes esperar para que sea resistente. Está obsoleto en favor del método basado en desplazamientos, que ofrece mejor rendimiento. Úsalo solo para código existente que aún no se ha migrado.

Registro a registro vs. ingestión por lotes

Cada método de ingestión tiene una variante por lotes (por ejemplo, ingest_records_offset()) que envía una lista de registros en una llamada. El procesamiento por lotes es más eficiente que las llamadas individuales para la ingesta masiva.

Para JSON y Protocol Buffers (protobuf), un lote se compromete atómicamente: o bien se acepta y se hace duradero cada registro del lote, o se rechaza todo el lote. Zerobus Ingest no realiza subidas parciales ni acuses de recibo parciales para estos formatos, por lo que tu tabla nunca contiene un lote parcial. Un lote que no supera la validación (por ejemplo, por una incompatibilidad de esquema) falla de inmediato, antes de afectar a la tabla, en lugar de cargar algunos registros y descartar otros.

Como un lote JSON o protobuf se envía como un solo mensaje, el tamaño máximo de mensaje de 10 MB se aplica tanto a un solo registro como a un lote completo: todos los registros de un lote juntos deben caber dentro de 10 MB. Dimensiona tus lotes para mantenerte por debajo de ese límite. Ver Tamaño de registro.

Los lotes de Arrow Flight son la excepción

La ingestión de Apache Arrow Flight no sigue el modelo de todo o nada y de mensaje único mencionado anteriormente. Un lote Arrow puede ser mucho más grande que un lote JSON o protobuf, y la trayectoria de vuelo Arrow divide un lote grande en mensajes de transporte más pequeños que se envían y reconocen individualmente en lugar de como una unidad atómica única. Como resultado:

  • El límite de 10 MB por mensaje que se aplica a lotes JSON y protobuf no se aplica a un lote Arrow de la misma manera. Un lote grande de Arrow se divide en mensajes de transporte en lugar de ser rechazado debido a su tamaño.
  • La durabilidad se confirma en la granularidad del mensaje de transporte, por lo que un lote lógico muy grande puede ser parcialmente duradero si ocurre un fallo a mitad de camino, en lugar de comprometer todo o nada.

ingest_batch() sigue devolviendo un único desplazamiento lógico para el lote que has enviado, y wait_for_offset() sobre ese desplazamiento solo se completa después de que se haya reconocido cada mensaje de transporte que compone el lote. Para consultar el modelo completo de Arrow Flight, la orientación sobre el procesamiento por lotes y la recuperación de datos no confirmados, consulta Uso de Arrow Flight con Zerobus Ingest.

¿Cuándo deberías bloquear un mensaje?

Bloquearse en un offset supone sacrificar rendimiento a cambio de una garantía de durabilidad por registro más sólida en el código cliente. Elige según tu carga de trabajo:

  • No bloquees: el valor predeterminado correcto para streaming de alto volumen, donde te importa el rendimiento sostenido y puedes confirmar la durabilidad en conjunto (por ejemplo, al cerrar la transmisión o mediante una llamada de confirmación). La mayoría de los productores deberían empezar por aquí.
  • Bloquear hasta un offset: ten esto en cuenta cuando tu aplicación deba tener la certeza de que un registro concreto se ha almacenado de forma duradera antes de realizar otra acción. Por ejemplo:
    • Estás a punto de eliminar o reconocer la fuente de los datos (un mensaje de cola, un archivo, un cursor de arriba) y no debes perderla si falla la ingestión.
    • Estás consumiendo puntos de control o límites transaccionales y necesitas que cada punto de control sea resistente antes de avanzar.
    • Realizas escrituras de bajo volumen y alto valor, donde la confirmación por registro importa más que el rendimiento.

No te bloquees con cada registro en un bucle de alto rendimiento. Eso serializa a tu productor en un viaje de ida y vuelta al servidor para cada disco y reduce drásticamente el rendimiento. En su lugar, Azure Databricks recomienda ingerir una gran cantidad de registros y luego confirmar la durabilidad una vez para todo el bloque. Tienes dos formas de hacerlo: esperar al desplazamiento más reciente, o vaciar el flujo. El bloqueo por registro individual debe reservarse para los casos específicos mencionados anteriormente donde debe confirmarse un registro antes de la siguiente acción.

Espera a que se despliegue

wait_for_offset() se bloquea hasta que Zerobus Ingest confirma que el registro en ese desplazamiento se haya escrito de forma duradera, o hasta que se agote el tiempo de espera. Úsalo para confirmar un punto específico del flujo, normalmente el último registro de un fragmento. Ingiera el bloque, conserve el offset final que devuelve el bucle y espere a ese único offset en lugar de esperar tras 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()

Vaciar el flujo

flush() bloquea hasta que todos los registros que has incorporado hasta ahora se hayan escrito de forma persistente y, a continuación, devuelve el control. A diferencia de wait_for_offset(), no se registra un desplazamiento: el flush espera en todo lo que está pendiente en el stream. No cierra el chorro, así que puedes seguir ingiriendo después.

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. descarga

Ambos confirman durabilidad para un trozo. Elige en función de lo que vayas a confirmar:

  • Úsalo wait_for_offset(offset) cuando quieras confirmar hasta un registro específico, por ejemplo un límite de punto de control, mientras que otros registros pueden seguir en vuelo detrás de él.
  • Usa flush() cuando quieras confirmar que todos los registros pendientes son persistentes antes de seguir, por ejemplo al final de un lote, antes de avanzar un cursor ascendente o antes de apagar el sistema. flush() se rige por un tiempo de espera configurable para el vaciado.

close() vacía y cierra el flujo, por lo que los registros siempre quedan almacenados de forma persistente en un apagado controlado. Llámalo siempre en un bloque finally.

Responder a los acuses de recibo de forma asíncrona

Si en lugar de bloquear quieres reaccionar a confirmaciones y errores de durabilidad a medida que llegan, mientras tu productor sigue empujando a toda velocidad, registra una llamada de reconocimiento en la transmisión. Las llamadas de regreso son una función independiente de las llamadas de bloqueo en esta página. Consulta las devoluciones de llamada de confirmación.