Callbacks de reconnaissance
Un rappel d'accusé de réception permet à votre client de réagir aux accusés de réception d'enregistrements et aux erreurs de manière asynchrone, sans bloquer votre boucle de production. À mesure que les enregistrements deviennent durables, ou échouent, Zerobus Ingest invoque votre rappel en arrière-plan, afin que vous puissiez suivre la progression et mettre à jour les métriques sans ralentir votre producteur, et prendre connaissance des échecs dès qu'ils surviennent.
Ceci diffère de l’attente d’un offset ou d’un vidage: il s’agit d’appels bloquants où votre code attend la durabilité en ligne. Un rappel (callback) n’est pas un appel bloquant. Il s’agit d’un gestionnaire que le SDK invoque pour vous lorsque les accusés de réception arrivent.
Les rappels d’accusé de réception sont pris en charge pour les flux SDK JSON et Protocol Buffers (protobuf). Les Stream Arrow Flight ne prennent pas en charge les rappels ; pour confirmer la durabilité sur un Stream Arrow, utilisez wait_for_offset() ou flush(). Consultez Utiliser Arrow Flight avec Zerobus Ingest.
Les noms de méthode et de type ci-dessous proviennent du SDK Python. D'autres SDK Zerobus exposent des callbacks d'acquittement là où ils sont pris en charge, en utilisant des constructions équivalentes dans chaque langage.
Fonctionnement des rappels
Vous définissez un callback en créant une sous-classe de AckCallback et en implémentant deux méthodes :
on_ack(offset: int): appelé lorsqu'une soumission (un enregistrement ou un batch) est confirmée comme durable par le serveur. Leoffsetidentifie la soumission confirmée.on_error(offset: int, error_message: str): appelé lorsqu'une soumission rencontre une erreur.on_errorest facultatif. Implémentez-le pour gérer ou log les échecs.
Le callback est appelé une fois pour chaque enregistrement ou batch soumis lorsque son offset logique est acquitté ou échoue ; il s'agit donc d'un signal continu de la progression de l'ingestion à travers le stream.
Vos méthodes de rappel s'exécutent sur les threads d'arrière-plan du SDK ; leur invocation ne bloque donc pas votre producteur. Maintenez-les rapides et non bloquants. La gestion d'une défaillance relève de la responsabilité de votre client : log, alerter, réessayer ou arrêter. Certaines erreurs sont fatales et, si on_error signale que le stream a échoué de manière permanente, vous devez effectuer une récupération sur un nouveau stream. Voir Modèles de récupération et de nouvelle tentative.
Configurer un callback
Vous attachez un callback à un Stream en passant une instance de votre sous-classe AckCallback en tant qu'option ack_callback dans StreamConfigurationOptions lors de la création du Stream. Le callback s'applique ensuite à chaque enregistrement ingéré sur ce 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()
Une fois le callback enregistré, vous n’attendez pas en ligne. on_ack se déclenche à mesure que chaque enregistrement est confirmé comme durable, et on_error se déclenche si un enregistrement échoue.
Callbacks vs. blocage
Les callbacks et les appels bloquants résolvent des problèmes différents, et vous pouvez les utiliser ensemble :
- Utilisez un callback d’accusé de réception pour réagir aux confirmations de durabilité et aux erreurs au fur et à mesure qu’elles surviennent, de manière asynchrone, tout en maintenant un throughput élevé. Utile pour le suivi de la progression, les métriques et la journalisation des erreurs.
- Utilisez
wait_for_offset()ouflush()lorsque votre code doit se bloquer jusqu’à ce qu’un enregistrement spécifique, ou tous les enregistrements en attente, soient durables avant de poursuivre.
Connexes
- Blocage et accusé de réception des messages: blocage sur la durabilité avec
wait_for_offsetetflush. - Modèles de récupération et de nouvelle tentative: gestion des erreurs et récupération des enregistrements non acquittés.
- Gestion des erreurs Zerobus Ingest: référence des codes d'erreur.