Modèles de récupération et de nouvelle tentative
Cette analyse approfondie explique comment créer un client Zerobus Ingest résilient dans Lakeflow Connect. Il couvre la récupération intégrée du SDK, la manière dont les erreurs apparaissent et comment récupérer les enregistrements non acquittés lorsqu’un stream échoue de manière permanente. Les noms des méthodes et des options ci-dessous proviennent du SDK Python. D’autres SDK proposent des équivalents.
Rétablissement intégré
Les SDK Zerobus récupèrent automatiquement des défaillances transitoires. La récupération est déclenchée lorsque le stream rencontre une erreur récupérable, généralement un délai d'attente ou une coupure réseau, ou lorsque le stream reçoit un signal d'arrêt normal. La récupération est activée par default et vous pouvez la régler via les options de configuration du Stream :
Option | Description |
|---|---|
| Activez la récupération automatique de stream. |
| Délai d'expiration pour une opération de récupération. |
| Délai entre les tentatives de récupération. |
| Nombre maximal de tentatives de récupération. |
Pour la plupart des workloads, les valeurs par défaut constituent un bon point de départ, et vous n'avez pas besoin d'écrire votre propre boucle de reconnexion pour les problèmes transitoires. Lorsque vous utilisez un SDK, les jetons OAuth sont également actualisés automatiquement lors de la création et de la récupération du stream ; votre client n'a donc rien à gérer. L'exception est l'API REST, où votre client récupère et refresh lui-même le jeton OAuth. Voir Utiliser Zerobus Ingest. Pour les valeurs et unités par default de ces options, consultez le repository du SDK Zerobus.
Erreurs et nouvelles tentatives
Le SDK réessaie automatiquement les erreurs temporaires, telles que les problèmes réseau ou les erreurs de serveur temporaires, grâce à sa récupération intégrée. Les échecs dont il ne peut pas se remettre, tels que des identifiants invalides ou une table manquante, apparaissent sous la forme ZerobusException. Interceptez ZerobusException pour gérer une défaillance, puis décidez s'il faut corriger la cause sous-jacente, effectuer une récupération sur un nouveau Stream ou arrêter.
from zerobus.sdk.shared import ZerobusException
try:
stream.ingest_record_offset(record)
except ZerobusException as e:
# Handle the failure: log it, fix the cause, recover on a new stream, or stop.
...
Pour la liste complète des codes d’erreur, voir gestion des erreurs Zerobus Ingest.
Récupération des enregistrements non acquittés
Lorsqu’un stream échoue de manière permanente, une fois la récupération automatique du SDK épuisée, les enregistrements soumis mais non encore accusés de réception par le serveur sont toujours conservés par le client dans le tampon en cours de transfert décrit dans Communication asynchrone. Récupérez-les afin de ne pas perdre de données :
get_unacked_records()renvoie les enregistrements non acquittés sous forme d'octets bruts.get_unacked_batches()renvoie les batchs non acquittés (chacun étant une liste d'enregistrements) pour la logique de nouvelle tentative de batch.
Les enregistrements reviennent sous leur forme sérialisée : décodez le JSON avec json.loads(record.decode('utf-8')), ou désérialisez les Protocol Buffers (protobuf) avec votre type de message. Enregistrez-les ou rejouez-les sur un nouveau stream.
Récupérer un stream après une défaillance permanente
Le SDK gère automatiquement les nouvelles tentatives pour les erreurs temporaires. Les échecs de mise en file d’attente (enqueue), de vidage (flush) et de fermeture (close) apparaissent tous sous la forme ZerobusException. get_unacked_records() et recreate_stream() ne réussissent qu’une fois le stream déjà fermé, ce qu’un échec terminal provoque. Un échec de mise en file d’attente laisse le stream actif, ces appels échouent donc ; dans ce cas, levez l’erreur initiale et conservez le stream. recreate_stream() remet en file d’attente les enregistrements qui ont déjà été acceptés ; il ne relance pas une charge utile qui n’a pas pu être mise en file d’attente.
from zerobus.sdk.shared import ZerobusException
try:
for i in range(10000):
stream.ingest_record_offset(record)
stream.flush()
except ZerobusException as e:
print(f"Ingestion failed: {e}")
try:
unacked = list(stream.get_unacked_records())
except ZerobusException:
raise e
print(f"{len(unacked)} previously queued records were unacknowledged.")
try:
new_stream = sdk.recreate_stream(stream)
try:
new_stream.flush()
finally:
new_stream.close()
except ZerobusException:
raise e
else:
stream.close()
Utilisez get_unacked_batches() pour inspecter le regroupement par batch d'origine après la fermeture du stream :
unacked_batches = list(stream.get_unacked_batches())
print(f"{len(unacked_batches)} batches remain unacknowledged")
Gestion des doublons lors de la relecture
Zerobus Ingest fournit une distribution au moins une fois, et non exactement une fois ; par conséquent, la relecture d’enregistrements récupérés, ou toute nouvelle tentative, peut entraîner l’écriture d’un enregistrement plus d’une fois. Si votre charge de travail ne tolère pas les doublons, dédupliquez les données dans le lakehouse :
- Incluez un identifiant unique stable sur chaque enregistrement (par exemple, un ID d’événement assigné à la source ou une clé naturelle).
- Dédupliquez lors de la lecture ou pendant le traitement en aval. Par exemple, utilisez un
MERGE INTOqui correspond à l’identifiant, ou unROW_NUMBER()fenêtré dans une transformation.
Comme les enregistrements sur un stream sont commités dans l’ordre, un numéro de séquence augmentant de façon monotone fonctionne également bien comme clé de déduplication.
Vidage et fermeture propre
flush()attend que le serveur confirme que les enregistrements que vous avez soumis sont durables, sans fermer le stream. Appelez-le lorsque vous avez besoin d'un point de contrôle de durabilité au milieu d'un stream.close()vide et ferme le stream en douceur, en attendant que les enregistrements en attente soient acquittés comme durables avant de retourner. Utilisez-le pour un arrêt en douceur, et non pour récupérer d'une défaillance de stream. Lorsqu'un stream a échoué, récupérez plutôt les enregistrements non acquittés, comme indiqué dans Récupérer un stream après une défaillance permanente.
Pour confirmer la durabilité d’un enregistrement spécifique plutôt que de l’ensemble du stream, consultez Message blocking and acknowledgment.
Récupération des données depuis l’emplacement de fallback durable
Si une modification incompatible est apportée à votre table cible après que Zerobus Ingest a rendu vos données durables mais avant qu’il ne puisse les publier, Zerobus Ingest écrit ces données sous forme de fichiers Parquet dans un répertoire de fallback situé sous la racine de stockage de votre table, au lieu de les supprimer. Voir Emplacement de fallback durable.
Comment savoir si des données y ont été écrites : le répertoire de fallback est _zerobus/table_rejected_parquets/, relatif à l'emplacement de stockage physique racine de la table. Si l'ingestion s'est poursuivie mais qu'il manque des lignes dans la table après une modification de table, vérifiez la présence de fichiers Parquet dans ce répertoire.
Retraiter les données de fallback dans la table
Une fois la cause corrigée (généralement en alignant le schéma de la table sur ce que vos producteurs envoient), retraitez les fichiers Parquet de fallback dans la table cible. Les fichiers sont au format Parquet standard dans l'emplacement de stockage de la table, vous pouvez donc les charger avec COPY INTO:
-
Résoudre l'incohérence de schéma. Faites évoluer la table cible (ou le schéma de votre producteur) afin que les enregistrements de fallback correspondent. Consultez la gestion de schémas.
-
Inspectez les données de fallback avant le chargement. Pointez une query vers le chemin de fallback pour confirmer ce qui s'y trouve et vérifier qu'il correspond désormais à la table :
SQLSELECT * FROM parquet.`<table-storage-root>/_zerobus/table_rejected_parquets/` LIMIT 10; -
Chargez les fichiers avec
COPY INTO.COPY INTOest idempotent : il suit les fichiers qu’il a déjà chargés, de sorte que son exécution ne chargera pas deux fois les mêmes fichiers de fallback :SQLCOPY INTO <catalog>.<schema>.<table>
FROM '<table-storage-root>/_zerobus/table_rejected_parquets/'
FILEFORMAT = PARQUET
COPY_OPTIONS ('mergeSchema' = 'false'); -
Vérifiez que le nombre de lignes attendu a bien été atteint, puis, une fois que vous avez confirmé que les données sont dans la table, nettoyez le répertoire fallback si vous n’en avez plus besoin.
Comme Zerobus Ingest fonctionne en mode at-least-once, les enregistrements qui ont été à la fois publiés dans la table et écrits dans l'emplacement de fallback pourraient être chargés deux fois. Si les doublons sont importants, dédupliquez comme décrit dans Gestion des doublons lors de la relecture. Pour un retraitement continu ou automatisé, vous pouvez pointer Auto Loader vers le chemin de fallback au lieu d’exécuter COPY INTO manuellement.
Cette procédure est une approche générale de premier passage utilisant les outils Delta standard sur les fichiers Parquet de fallback. Validez-le par rapport à votre configuration de table et de stockage avant de vous y fier en production.
Connexes
- Blocage et accusé de réception des messages: méthodes d’ingestion et confirmation de la durabilité.
- Gestion des erreurs Zerobus Ingest: référence complète des codes d'erreur.
- Emplacement de fallback durable: l'emplacement de fallback durable.