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.
Les SDK Zerobus Ingest offrent plusieurs méthodes pour ingérer un enregistrement, qui compensent le débit contre la confirmation de durabilité que vous obtenez. Cette page explique chaque méthode et quand bloquer la durabilité. Pour réagir aux accusés de réception de manière asynchrone sans bloquer, voir les fonctions de rappel d’accusé de réception.
Les exemples sur cette page utilisent le SDK Python. Pour connaître les délais exacts et les options de configuration acceptées par chaque méthode (y compris leurs valeurs par défaut et leurs unités), voir le dépôt Zerobus SDK. Les autres SDK de langage exposent des options équivalentes.
Qu’est-ce qu’un offset ?
Chaque enregistrement que vous importez se voit assigner un offset : sa position dans le flux. Le décalage sert à désigner un enregistrement précis lorsque vous voulez confirmer qu’il a bien été écrit de façon durable. Zerobus Ingest garantit une livraison au moins une fois, et l’attente d’un décalage est le moyen par lequel un client confirme cette garantie pour un enregistrement précis.
Confirmer un décalage signifie que l’enregistrement est stocké de manière durable, et non qu’il est déjà interrogeable dans la table Delta. Zerobus Ingest matérialise des enregistrements durables dans la table peu de temps après. Pour les chiffres de latence, voir Latence.
Méthodes d’ingestion
Les SDK offrent deux façons d’ingérer un enregistrement. (Les noms des méthodes ci-dessous proviennent du SDK Python. D’autres SDK exposent des méthodes équivalentes.)
| Méthode | Returns | Utilisez-le quand |
|---|---|---|
Basé sur le décalage, ingest_record_offset() |
Le décalage de l’enregistrement, après que l’enregistrement est mis en file d’attente sur le stream. | Valeur par défaut recommandée Vous souhaitez placer les enregistrements dans l’ordre en file d’attente et, éventuellement, confirmer ensuite leur persistance en attendant qu’un décalage soit atteint. |
Axé sur l’avenir, ingest_record() |
Un RecordAcknowledgment sur lequel vous pouvez compter. |
Deprecated. Préférez une approche basée sur le décalage pour de meilleures performances. |
Basé sur le décalage (recommandé)
ingest_record_offset() soumet l’enregistrement et renvoie son décalage une fois que l’enregistrement est placé en file d’attente sur le flux. L’appel s’exécute sur votre fil d’appel, donc les enregistrements sont mis en file d'attente dans l’ordre dans lequel vous appelez la méthode, et le décalage retourné vous permet de confirmer la durabilité plus tard avec wait_for_offset(). Il s’agit de la valeur par défaut recommandée pour la plupart des producteurs, et c’est la méthode utilisée dans les exemples Utiliser Zerobus Ingest.
Axé sur le futur (déprécié)
ingest_record() renvoie un RecordAcknowledgment objet sur lequel vous pouvez attendre la confirmation de la durabilité. Elle est déconseillée au profit de la méthode fondée sur le décalage, qui offre de meilleures performances. Utilisez-le uniquement pour le code existant qui n’a pas encore migré.
Enregistrement par dossier vs. ingestion par lot
Chaque méthode d’ingestion possède une variante batch (par exemple, ingest_records_offset()) qui soumet une liste d’enregistrements lors d’un appel. Le traitement par lots est plus efficace que les appels individuels pour une ingestion en masse.
Pour JSON et les tampons de protocole (protobuf), un lot est validé atomiquement : soit chaque enregistrement du lot est accepté et rendu durable, soit l’ensemble du lot est rejeté. Zerobus Ingest n’effectue pas de chargements partiels ni d’accusés de réception partiels pour ces formats, de sorte que votre table ne contient jamais de lot partiel. Un lot qui ne passe pas la validation (par exemple, en raison d’une incompatibilité de schéma) est rejeté immédiatement, avant même d’atteindre la table, au lieu de charger certains enregistrements et d’en abandonner d’autres.
Comme un lot JSON ou protobuf est envoyé en message unique, la taille maximale de 10 Mo s’applique à la fois à un seul enregistrement et à un lot entier : tous les enregistrements d’un lot doivent tenir dans un rayon de 10 Mo. Dimensionnez vos lots afin de ne pas dépasser cette limite. Voir Taille du disque.
Les lots Arrow Flight constituent l’exception
L’ingestion de Apache Arrow Flight ne suit pas le modèle tout ou rien à message unique ci-dessus. Un lot Arrow peut être beaucoup plus grand qu’un lot JSON ou protobuf, et la trajectoire Arrow Flight divise un grand lot en messages de transport plus petits qui sont envoyés et accusés de réception individuellement plutôt que comme une seule unité atomique. Par conséquent :
- La limite de 10 Mo par message qui s’applique aux lots JSON et protobuf ne s’applique pas à un lot Arrow de la même manière. Un lot Arrow volumineux est divisé en messages de transport au lieu d’être rejeté parce qu’il est trop volumineux.
- La durabilité est confirmée à la granularité du message de transport, donc un très gros lot logique peut être partiellement durable si une panne survient en cours de route, plutôt que de compromettre tout ou rien.
ingest_batch() retourne toujours un décalage logique unique pour le lot que vous avez soumis, et wait_for_offset() sur ce décalage n’est terminé qu’après que chaque message de transport constituant le lot a fait l’objet d’un accusé de réception. Pour le modèle complet d’Arrow Flight, les recommandations relatives au traitement par lots et la récupération des données dont la réception n’a pas été confirmée, voir Utiliser Arrow Flight avec Zerobus Ingest.
Quand faut-il bloquer un message ?
Bloquer sur un offset sacrifie le débit au profit d’une garantie de durabilité plus forte pour chaque enregistrement dans votre code client. Choisissez en fonction de votre charge de travail :
- Ne bloquez pas : le bon par défaut pour le streaming à haut volume, où vous vous souciez d’un débit soutenu et pouvez confirmer la durabilité en général (par exemple, à la fermeture du stream ou via un rappel d’accusé de réception). La plupart des producteurs devraient commencer par ici.
-
Blocage sur un décalage : tenez-en compte lorsque votre application doit s’assurer qu’un enregistrement précis a bien été écrit de manière durable avant d’effectuer une autre action. Par exemple :
- Vous êtes sur le point de supprimer ou de reconnaître la source des données (un message de file d’attente, un fichier, un curseur en amont) et ne devez pas la perdre si l’ingestion échoue.
- Vous effectuez l’ingestion par points de contrôle ou par limites transactionnelles et vous devez vous assurer que chaque point de contrôle est durable avant de passer au suivant.
- Vous effectuez des écritures de faible volume et de grande valeur, pour lesquelles la confirmation pour chaque enregistrement importe davantage que le débit.
Ne bloquez pas sur chaque enregistrement dans une boucle à haut débit. Cela sérialise votre producteur lors d’un aller-retour au serveur pour chaque enregistrement et réduit fortement le débit. Au lieu de cela, Azure Databricks recommande d’ingérer un lot important d’enregistrements, puis de confirmer la persistance une seule fois pour l’ensemble du lot. Vous avez deux façons de faire cela : attendre le dernier décalage, ou rincer le flux. 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 décalage
wait_for_offset() bloque jusqu’à ce que Zerobus Ingest confirme que l’enregistrement à cette position a été écrit de manière durable, ou jusqu’à ce que le délai d’attente expire. Utilisez-le pour confirmer un point précis du flux, le plus souvent le dernier enregistrement d’un bloc. Traitez le bloc, conservez l’offset final renvoyé par la boucle, et attendez cet unique décalage au lieu d’attendre après chaque enregistrement :
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()
Purgez le flux
flush() bloque jusqu’à ce que tous les enregistrements que vous avez ingérés jusqu’à présent soient durablement écrits, puis revient. Contrairement à wait_for_offset(), vous ne suivez pas de décalage : le flush attend que toutes les opérations en attente sur le flux soient terminées. Il ne ferme pas le flux, vous pouvez donc continuer l’ingestion par la suite.
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é d’un bloc. Choisissez en fonction de ce que vous confirmez :
- Utilisez
wait_for_offset(offset)lorsque vous souhaitez valider jusqu’à un enregistrement spécifique, par exemple la frontière d’un 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 vérifier que tous les enregistrements en attente sont durables avant de passer à autre chose, par exemple à la fin d’un lot, avant d’avancer un curseur en amont, ou avant de l’arrêter.flush()est régi par un délai de vidage configurable.
close() efface et ferme le flux, donc les disques sont toujours rendus durables lors d’un arrêt approprié. Appelez-le toujours dans un bloc finally.
Réagir aux accusés de réception de manière asynchrone
Si, au lieu de bloquer, vous souhaitez réagir aux confirmations de durabilité et aux erreurs au fur et à mesure qu’elles arrivent, pendant que votre producteur continue de pousser à pleine vitesse, enregistrez un rappel d’accusé de réception sur le flux. Les rappels sont une fonctionnalité distincte des appels de blocage sur cette page. Voir Rappels d’accusé de réception.
Related
- Rappels d’accusés de réception : Réagissez de manière asynchrone aux accusés de réception et aux erreurs.
- Utilisez Zerobus Ingest : Développez un client.
- Flux : Flux et décalages.
- Gestion des erreurs de Zerobus Ingest : Gestion des erreurs.