Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Une fonction de rappel de confirmation permet à votre client de réagir de façon asynchrone aux confirmations d’enregistrement et aux erreurs, sans bloquer votre boucle de production. Lorsque les enregistrements deviennent durables ou tombent en panne, Zerobus Ingest invoque votre rappel en arrière-plan, afin que vous puissiez suivre vos progrès et mettre à jour les métriques sans ralentir votre producteur, et apprendre les échecs dès qu’ils surviennent.
Cela diffère de l’attente d’un décalage ou de l’effacement : ce sont des appels de blocage où votre code attend la durabilité en ligne. Un rappel n’est pas un appel de blocage. C’est 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 flux Arrow Flight ne prennent pas en charge les fonctions de rappel; pour confirmer la persistance sur un flux Arrow, utilisez wait_for_offset() ou flush(). Voir Use Arrow Flight avec Zerobus Ingest.
Les noms des méthodes et types ci-dessous proviennent du SDK Python. D’autres SDK Zerobus proposent des fonctions de rappel d’accusé de réception lorsque cette fonctionnalité est prise en charge, en utilisant des mécanismes équivalents dans chaque langage de programmation.
Comment fonctionnent les rappels
Vous définissez un callback en sous-classant AckCallback et en implémentant deux méthodes :
-
on_ack(offset: int): appelé lorsqu’une soumission (un enregistrement ou un lot) est reconnue avec succès comme durable par le serveur.offsetidentifie la soumission ayant fait l’objet d’un accusé de réception. -
on_error(offset: int, error_message: str): appelé lorsqu’une soumission rencontre une erreur.on_errorest facultatif. Implémentez-le pour gérer ou enregistrer les pannes.
La fonction de rappel est appelée une fois pour chaque enregistrement ou lot soumis lorsque son offset logique est confirmé ou que le traitement échoue ; elle constitue ainsi un indicateur en continu de la progression de l’ingestion sur l’ensemble du flux.
Vos méthodes de rappel s’exécutent sur les threads en arrière-plan du SDK, donc les invoquer ne bloque pas votre producteur. Gardez-les performants et non bloquants. Ce qu’il faut faire face à une défaillance est la responsabilité de votre client : enregistrer, alerter, réessayer ou arrêter. Certaines erreurs sont terminales et, si on_error signale que le flux a échoué de façon permanente, vous devez reprendre sur un nouveau flux. Voir les schémas de récupération et de réessayage.
Configurer un rappel
Vous associez un callback à un flux en passant une instance de votre sous-classe AckCallback comme option ack_callback dans StreamConfigurationOptions, lors de la création du flux. Le rappel s’applique alors à chaque enregistrement ingéré sur ce flux.
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 la fonction de rappel enregistrée, vous n’attendez pas de manière synchrone.
on_ack s’enflamme à mesure que chaque enregistrement est confirmé durable, et on_error s’enflamme si un enregistrement échoue.
Fonctions de rappel vs. blocage
Les rappels et les appels de blocage résolvent différents problèmes, et vous pouvez les utiliser ensemble :
- Utilisez une fonction de rappel d’accusé de réception pour réagir aux confirmations de durabilité et aux erreurs lorsqu’elles surviennent, de façon asynchrone, tout en maintenant un débit élevé. C’est bon pour le suivi des progrès, les métriques et l’enregistrement des erreurs.
- Utilisez
wait_for_offset()ouflush()lorsque votre code doit être bloqué jusqu’à ce qu’un enregistrement spécifique, ou tous les documents en attente, soient durables avant de continuer.
Related
-
Blocage et accusé de réception des messages : Blocage sur 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 confirmés.
- Gestion des erreurs Zerobus Ingest : référence au code d’erreur.