Aller au contenu principal

Blocage et confirmation de message

Les SDK Zerobus Ingest offrent plusieurs méthodes pour ingérer un enregistrement, qui permettent de trouver un compromis entre le throughput et le niveau de confirmation de durabilité obtenu. Cette page explique chaque méthode et quand bloquer sur la durabilité. Pour réagir aux accusés de réception de manière asynchrone au lieu de bloquer, consultez Callbacks d'accusé de réception.

Les exemples de cette page utilisent le SDK Python. Pour connaître le délai d'expiration exact et les options de configuration acceptées par chaque méthode (y compris leurs valeurs par default et leurs unités), consultez le repository du SDK Zerobus. Les SDK des autres langages exposent des options équivalentes.

Qu’est-ce qu’un décalage ?

Chaque enregistrement que vous ingérez se voit attribuer un offset : sa position dans le stream. L'offset est la manière dont vous faites référence à un enregistrement spécifique lorsque vous souhaitez confirmer qu'il a été écrit de manière durable. Zerobus Ingest fournit des garanties de livraison au moins une fois, et l'attente sur un offset est la façon dont un client confirme cette garantie pour un enregistrement donné.

La confirmation d'un offset signifie que l'enregistrement est durable, et non qu'il est déjà interrogeable dans la table Delta. Zerobus Ingest matérialise les enregistrements durables dans la table peu de temps après. Pour les chiffres de latence, consultez Latence.

Méthodes d’ingestion

Les SDK offrent deux manières d'ingérer un enregistrement. (Les noms de méthode ci-dessous proviennent du SDK Python. D'autres SDK exposent des méthodes équivalentes.)

Méthode

Renvoie

Utilisez-le lorsque

Basé sur l’offset , ingest_record_offset()

L'offset de l'enregistrement, une fois que l'enregistrement est mis en file d'attente sur le stream.

Recommended default. Vous souhaitez mettre les enregistrements en file d'attente dans l'ordre et éventuellement confirmer la durabilité plus tard en attendant un offset.

Basé sur le futur , ingest_record()

Un RecordAcknowledgment sur lequel vous pouvez attendre.

Obsolète. Préférez le mode basé sur le décalage (offset) pour de meilleures performances.

Méthode

Renvoie

Utilisez-le lorsque

Basé sur l’offset , ingest_record_offset()

L'offset de l'enregistrement, une fois que l'enregistrement est mis en file d'attente sur le stream.

Recommended default. Vous souhaitez mettre les enregistrements en file d'attente dans l'ordre et éventuellement confirmer la durabilité plus tard en attendant un offset.

Basé sur le futur , ingest_record()

Un RecordAcknowledgment sur lequel vous pouvez attendre.

Obsolète. Préférez le mode basé sur le décalage (offset) pour de meilleures performances.

Basé sur l’offset (recommandé)

ingest_record_offset() soumet l'enregistrement et renvoie son offset une fois que l'enregistrement est mis en file d'attente sur le Stream. L'appel s'exécute sur votre thread appelant, de sorte que les enregistrements sont mis en file d'attente dans l'ordre où vous appelez la méthode, et l'offset renvoyé vous permet de confirmer la durabilité ultérieurement avec wait_for_offset(). Il s'agit du « default » recommandé pour la plupart des producteurs, et c'est la méthode utilisée dans les exemples Utiliser Zerobus Ingest.

Basé sur le futur (obsolète)

ingest_record() renvoie un objet RecordAcknowledgment sur lequel vous pouvez attendre pour la durabilité. Il est déprécié au profit de la méthode basée sur l'offset, qui offre de meilleures performances. Utilisez-le uniquement pour le code existant qui n'a pas encore été migré.

Ingestion enregistrement par enregistrement ou par batch

Chaque méthode d’ingestion possède une variante par batch (par exemple, ingest_records_offset()) qui soumet une liste d’enregistrements en un seul appel. Le traitement par batch est plus efficace que les appels individuels pour l’ingestion en masse.

Pour JSON et Protocol Buffers (protobuf), un batch commit de manière atomique : soit chaque enregistrement du batch est accepté et rendu durable, soit l’intégralité du batch est rejetée. Zerobus Ingest n’effectue pas d’upload partiel ni d’accusé de réception partiel pour ces formats ; votre table ne contient donc jamais de batch partiel. Un batch qui échoue à la validation (par exemple, une incompatibilité de schéma) échoue rapidement, avant d’atteindre la table, plutôt que d’enregistrer certains éléments et d’en supprimer d’autres.

Comme un batch JSON ou protobuf est envoyé sous forme de message unique, la taille maximale de message de 10 Mo s'applique à la fois à un enregistrement unique et à un batch complet : tous les enregistrements d'un batch doivent tenir dans 10 Mo. Dimensionnez vos batches pour rester en dessous de cette limite. Consultez Taille d'enregistrement.

Les batches Arrow Flight font exception

L’ingestion Apache Arrow Flight ne suit pas le modèle « tout ou rien » à message unique ci-dessus. Un batch Arrow peut être beaucoup plus volumineux qu’un batch JSON ou protobuf, et le chemin Arrow Flight divise un batch important en messages de transport plus petits qui sont envoyés et accusés réception individuellement plutôt que comme une unité atomique. En conséquence :

  • La limite de 10 Mo par message qui s’applique aux batchs JSON et protobuf ne s’applique pas de la même manière à un batch Arrow. Un batch Arrow volumineux est divisé en messages de transport au lieu d’être rejeté pour sa taille.
  • La durabilité est confirmée au niveau de la granularité des messages de transport ; ainsi, un batch logique très volumineux peut être partiellement durable si une défaillance survient en cours de route, plutôt que de commit selon une logique de tout ou rien.

ingest_batch() renvoie toujours un offset logique unique pour le batch que vous avez soumis, et wait_for_offset() sur cet offset ne se termine qu'après que chaque message de transport constituant le batch a été acquitté. Pour le modèle Arrow Flight complet, les conseils sur le traitement par batch et la récupération des données non acquittées, consultez Utiliser Arrow Flight avec Zerobus Ingest.

Quand devez-vous bloquer un message ?

Le blocage sur un offset sacrifie le throughput au profit d'une garantie de durabilité par enregistrement plus forte dans votre code client. Faites votre choix en fonction de votre workload :

  • Ne pas bloquer : la valeur default appropriée pour le streaming à haut volume, où vous vous souciez d’un throughput soutenu et pouvez confirmer la durabilité de manière agrégée (par exemple, à la fermeture du Stream ou via un rappel d’accusé de réception). La plupart des producteurs devraient start ici.
  • Bloquer sur un décalage : envisagez cette option lorsque votre application doit savoir qu'un enregistrement spécifique est durable avant d'effectuer une autre action. Par exemple :
    • Vous êtes sur le point de supprimer ou d'acquitter la source des données (un message de file d'attente, un fichier, un curseur amont) et ne devez pas la perdre en cas d'échec de l'ingestion.
    • Vous effectuez l'ingestion dans des points de contrôle ou des limites transactionnelles et avez besoin que chaque point de contrôle soit durable avant de continuer.
    • Vous effectuez des écritures à faible volume et à haute valeur où la confirmation par enregistrement importe plus que le throughput.

Ne bloquez pas sur chaque enregistrement dans une boucle à throughput élevé. Cela sérialise votre producteur lors d'un aller-retour vers le serveur pour chaque enregistrement et réduit considérablement le throughput. Au lieu de cela, Databricks recommande d'ingérer un grand bloc d'enregistrements, puis de confirmer la durabilité une seule fois pour l'ensemble du bloc. Vous avez deux façons de procéder : attendre le dernier offset ou vider le Stream. Le blocage par enregistrement individuel doit être réservé aux cas spécifiques ci-dessus où un seul enregistrement doit être confirmé avant l'action suivante.

Attendre un offset

wait_for_offset() bloque jusqu'à ce que Zerobus Ingest confirme que l'enregistrement à ce décalage est écrit de manière durable, ou jusqu'à l'expiration du délai. Utilisez-le pour confirmer un point spécifique dans le Stream, généralement le dernier enregistrement d'un bloc. Ingérez le bloc, conservez le décalage final renvoyé par la boucle et attendez sur ce décalage unique au lieu d'attendre après chaque enregistrement :

Python
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()

Vider le Stream

flush() bloque jusqu'à ce que tous les enregistrements que vous avez ingérés jusqu'à présent soient écrits de manière durable, puis renvoie. Contrairement à wait_for_offset(), vous ne suivez pas de décalage : flush attend que tout ce qui est en attente sur le stream soit traité. Cela ne ferme pas le stream, vous pouvez donc continuer l'ingestion par la suite.

Python
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. flush

Les deux confirment la durabilité pour un segment. Choisissez en fonction de ce que vous confirmez :

  • Utilisez wait_for_offset(offset) lorsque vous souhaitez confirmer jusqu’à un enregistrement spécifique, par exemple une limite de point de contrôle, alors que d’autres enregistrements peuvent encore être en cours de traitement derrière celui-ci.
  • Utilisez flush() lorsque vous souhaitez confirmer que tous les enregistrements en attente sont durables avant de continuer, par exemple à la fin d'un batch, avant d'avancer un curseur amont ou avant l'arrêt. flush() est régi par un délai d'expiration de vidage configurable.

close() vide et ferme le stream, afin que les enregistrements soient toujours rendus durables lors d'un arrêt propre. Appelez-le toujours dans un bloc finally.

Réagir aux confirmations de manière asynchrone

Si, au lieu de bloquer, vous souhaitez réagir aux confirmations de durabilité et aux erreurs au fur et à mesure de leur arrivée, tout en laissant votre producteur continuer à envoyer des données à pleine vitesse, enregistrez un callback d'accusé de réception sur le stream. Les callbacks sont une fonctionnalité distincte des appels bloquants sur cette page. Voir Callbacks d'accusé de réception.

Connexes