Llamadas de reconocimiento

Una llamada de reconocimiento permite a tu cliente reaccionar de forma asíncrona a los acuses de recibo y errores de grabación, sin bloquear el bucle de producción. A medida que los registros se vuelven duraderos o fallan, Zerobus Ingest invoca tu callback en segundo plano, para que puedas seguir el progreso y actualizar métricas sin ralentizar a tu productor, y enterarte de los fallos tan pronto como ocurren.

Esto difiere de esperar un desplazamiento o un flushing: son llamadas de bloqueo donde tu código espera durabilidad en línea. Una llamada de devolución no es una llamada de bloqueo. Es un manejador que el SDK invoca por ti cuando llegan las confirmaciones de recepción.

Se soportan callbacks de acuse de recibo para flujos de SDK JSON y Protocol Buffers (protobuf). Los streams de Arrow Flight no soportan callbacks; para confirmar la durabilidad de un chorro Arrow, usa wait_for_offset() o flush(). Consulte Usar Arrow Flight con Zerobus Ingest.

Los nombres de los métodos y tipos que aparecen a continuación provienen del SDK de Python. Otros SDK de Zerobus exponen funciones de devolución de llamada de confirmación de recepción, en los casos en que están disponibles, utilizando mecanismos equivalentes en cada lenguaje.

Cómo funcionan los callbacks

Defines un callback subclasificando AckCallback e implementando dos métodos:

  • on_ack(offset: int): llamada cuando una entrega (un registro o un lote) es reconocida con éxito como duradera por el servidor. El offset identifica la entrega confirmada.
  • on_error(offset: int, error_message: str): llamado cuando una publicación encuentra un error. on_error es opcional. Implementarlo para gestionar o registrar fallos.

La función de retorno se invoca una vez por cada registro o lote enviado cuando se confirma su desplazamiento lógico o cuando dicha confirmación falla, por lo que es un indicador continuo del progreso de la ingestión en todo el flujo.

Tus métodos de callback se ejecutan en los hilos de segundo plano del SDK, así que invocarlos no bloquea a tu productor. Mantenlos rápidos y sin bloqueos. Qué hacer ante un fallo es responsabilidad de tu cliente: registrar, alertar, intentar de nuevo o detener. Algunos errores son terminales y, si on_error informa de que el flujo ha fallado permanentemente, debes recuperarte con un flujo nuevo. Consulta patrones de recuperación y reintentos.

Configurar una devolución de llamada

Asocias una función de devolución de llamada a un flujo pasando una instancia de tu subclase AckCallback como opción ack_callback en StreamConfigurationOptions cuando creas el flujo. La callback se aplica entonces a todos los registros ingeridos en ese stream.

from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import AckCallback, StreamConfigurationOptions, TableProperties

class MyAckCallback(AckCallback):
    def on_ack(self, offset: int) -> None:
        print(f"Record acknowledged at offset: {offset}")

    def on_error(self, offset: int, error_message: str) -> None:
        print(f"Error at offset {offset}: {error_message}")

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

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

try:
    for row in records:
        stream.ingest_record_offset(row)
finally:
    stream.close()

Con la devolución de llamada registrada, no esperas en la cola. on_ack se activa cuando se confirma la persistencia de cada registro, y on_error se activa si un registro falla.

Callbacks frente al bloqueo

Las callbacks y las llamadas de bloqueo resuelven diferentes problemas, y puedes usarlas juntas:

  • Utiliza una llamada de reconocimiento para reaccionar a confirmaciones y errores de durabilidad a medida que ocurren, de forma asíncrona, manteniendo un alto rendimiento. Bueno para el seguimiento del progreso, métricas y registro de errores.
  • Usa wait_for_offset() o flush() cuando tu código deba bloquearse hasta que un registro específico, o todos los registros pendientes, sean duraderos antes de continuar.