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 nouvelles tentatives et les réexécutions sont inévitables dans tout pipeline réel. Cette page explique donc les garanties de traitement fournies par les pipelines Lakeflow et comment faire en sorte que les parties que vous écrivez puissent être réexécutées en toute sécurité.
Overview
Deux propriétés connexes déterminent si la reprise d’un pipeline est sûre :
- L’idempotence signifie qu’un pipeline produit le même résultat, peu importe combien de fois on l’exécute sur la même entrée. Réexécuter après une défaillance, remplir deux fois une plage de dates ou relancer manuellement un travail ne crée jamais de lignes dupliquées ni ne corrompt l’état.
- La garantie de traitement décrit le nombre de fois que chaque enregistrement affecte le résultat. Au moins une fois, le traitement garantit que chaque enregistrement est traité, mais une défaillance et une nouvelle tentative peuvent traiter certains enregistrements plus d’une fois, ce qui risque de dupliques. Le traitement exact-once garantit que chaque enregistrement affecte le résultat comme s’il avait été traité précisément une fois, même entre les tentatives, sans doublons ni interruptions.
Les pipelines Lakeflow sont idempotents par défaut pour les éléments qu’ils gèrent et vous garantissent un traitement « exactement une fois » dans leurs propres tables gérées. L’essentiel à comprendre est où ces garanties cessent d’être automatiques, afin de pouvoir ajouter les bonnes protections aux bords de votre pipeline.
Fonctionnement
Les pipelines Lakeflow fournissent un traitement « exactement une fois » et une idempotence pour les flux qu’ils gèrent, et vous donnent des outils pour que la logique que vous écrivez reste idempotente.
Traitement « exactement une fois » pour les tables gérées
Avec les tables gérées, vous bénéficiez par défaut d’un traitement exactement une fois. Les tables de streaming utilisent les points de contrôle de Structured Streaming combinés aux écritures transactionnelles de Delta Lake : chaque micro-lot valide ses offsets de source et ses données de sortie ensemble, de sorte qu’un lot relancé après un échec soit réussit complètement, soit est entièrement annulé, puis relancé, sans jamais être appliqué partiellement deux fois. Cela s’applique à l’ingestion de fichiers avec le chargeur automatique, à la lecture depuis Kafka, Kinesis et Azure Event Hubs, ainsi qu’aux opérations d’upsert AUTO CDC, sans aucun code de votre part.
Si une source de type « au moins une fois » envoie le même enregistrement plusieurs fois, le pipeline les traite comme des enregistrements distincts et les écrit tous dans votre table. Supprimer ces doublons est votre responsabilité. Consultez Dédupliquer les sources au moins une fois.
L’idempotence en lecture découle de ces mêmes points de contrôle. Les points de contrôle du chargeur automatique et des tables de streaming garantissent que chaque fichier source ou décalage n’est traité qu’une seule fois à des fins de suivi de l’état. Ainsi, le retraitement d’une mise à jour du pipeline après une défaillance reprend à partir du point de contrôle au lieu de retraiter ou d’ignorer des données. Vous obtenez cela en utilisant des tables de streaming sur spark.readStream plutôt que des boucles de traitement par lots développées manuellement. Consultez Streaming tables.
Utilisez AUTO CDC au lieu de MERGE écrit manuellement
AUTO CDC INTO est intrinsèquement idempotente par rapport à ses keys et sequence_by. Appliquer deux fois le même enregistrement de modification, ou appliquer des enregistrements dans le mauvais ordre, produit le même état final, car le pipeline utilise la colonne de séquence pour déterminer si une ligne entrante est réellement plus récente que celle stockée :
CREATE FLOW customers_cdc_flow AS AUTO CDC INTO customers_silver
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY sequence_num
STORED AS SCD TYPE 1;
Si vous écrivez votre propre logique d’upsert en dehors de AUTO CDC (c’est rare, mais parfois nécessaire pour des règles de fusion complexes), basez-la sur une clé métier stable et faites en sorte qu’elle puisse être appliquée deux fois sans risque, par exemple un MERGE ... WHEN MATCHED basé sur order_id plutôt qu’un INSERT aveugle. Pour plus d’informations, voir Les API AUTO CDC : Simplifier la capture des données de modifications avec des pipelines.
Gardez vos propres transformations idempotentes
Pour garder la logique idempotente lors de la reprise des opérations d’écriture, suivez ces deux directives :
- Évitez les transformations non déterministes dans les vues matérialisées. Comme une vue matérialisée peut être recalculée entièrement ou de manière incrémentale, évitez les fonctions dont la sortie dépend du moment où elles s’exécutent plutôt que de l’entrée . Par exemple, ne pas utiliser
current_timestamp()pour calculer une valeur métier qui devrait rester fixe une fois écrite ; prenez l’horodatage de l’événement source ou passez-le comme paramètre pour que le recalcul produise une sortie identique. - Concevez des rafraîchissements complets en toute sécurité. Une actualisation complète supprime et recalcule entièrement une table à partir de zéro, ce qui n’est sûr que si chaque source en amont peut encore produire l’historique complet. Si une source amont n’expose qu’une fenêtre glissante des modifications, un rafraîchissement complet d’une table
AUTO CDCen aval peut entraîner une perte silencieuse de l’historique ; prévoyez donc la rétention de la source et des topics en conséquence.
Obtenez « exactement une fois » aux bords
L’exécution « exactement une fois » cesse d’être automatique aux bords de ce que le pipeline contrôle directement, comme les écritures vers des systèmes externes. Lorsque vous effectuez une distribution (fan-out) vers un système externe, faites en sorte que l’écriture elle-même soit idempotente, par exemple, avec un upsert par clé du côté récepteur. Sinon, une relance de micro-lot pourrait écrire le même lot deux fois. Le récepteur suivant écrit chaque partition du lot depuis les exécuteurs et utilise une clé d’idempotence pour qu’un lot relancé n’écrive pas les données en double :
from pyspark import pipelines as dp
@dp.foreach_batch_sink(name="orders_to_external_api")
def write_orders_to_api(batch_df, batch_id):
def write_partition(rows):
# Open one client per partition.
for row in rows:
# Use an idempotency key (order_id) so a retried batch doesn't double-write.
upsert_to_external_system(key=row.order_id, payload=row.asDict())
batch_df.select("order_id", "amount").foreachPartition(write_partition)
Pour plus d’informations sur l’écriture dans des systèmes externes, consultez Récepteurs dans les pipelines Lakeflow.
Dédupliquez au moins une fois les sources
Lorsqu’une source peut livrer un enregistrement plus d’une fois, dédupliquez en aval. Associez un filigrane avec dropDuplicatesWithinWatermark, qui est compatible avec les filigranes et ne nécessite pas un état illimité pour détecter les doublons. Supprimez les doublons selon les colonnes qui permettent d’identifier de manière unique un événement. L’identité peut s’étendre sur plusieurs colonnes lorsqu’aucune colonne n’est unique en soi. Dans l’exemple suivant, un numéro de séquence de clics n’est unique qu’au sein de sa session, donc les deux colonnes ensemble identifient l’événement :
from pyspark import pipelines as dp
@dp.table(name="clicks_deduped")
def clicks_deduped():
return (
spark.readStream.table("clicks_bronze")
.withWatermark("click_ts", "5 minutes")
.dropDuplicatesWithinWatermark(["session_id", "click_seq_num"])
)
Choisissez ces colonnes du contrat d’unicité de la source, et non celles qui semblent distinctes dans les données d’échantillon. Les colonnes qui peuvent légitimement se répéter ignorent les événements réels lorsque vous les utilisez comme identifiants. Un utilisateur cliquant deux fois sur la même publicité est un exemple courant : déduplication de l’utilisateur et la publicité supprime silencieusement le deuxième clic.
La sémantique d’upsert basé sur une clé de AUTO CDC fusionne aussi naturellement les doublons. Ainsi, le routage de données « au moins une fois » via un flux AUTO CDC utilisant une clé professionnelle stable comme clé constitue une autre façon de converger vers un état de type « exactement une fois ».
Limitations
Le traitement « exactement une fois » s’applique aux flux Delta vers Delta gérés. Traitez les bords suivants comme « au moins une fois » et ajoutez une logique explicite de déduplication ou d’écriture idempotente :
-
foreach_batch_sinket écritures externes personnalisées. Spark garantit qu’un batch est tenté au moins une fois, mais un batch réessayé après une écriture partielle peut laisser certaines lignes visibles deux fois dans le système externe. Rendez l’écriture externe idempotente, par exemple en effectuant un upsert sur une clé naturelle ou en écrivant un ID de lot sur lequel le récepteur peut dédupliquer. - Kafka comme récepteur. Les rubriques Kafka ne prennent pas en charge les écritures transactionnelles « exactement une fois » comme le fait Delta, donc une nouvelle tentative de micro-lot écrivant sur Kafka peut produire des messages en double. Si les consommateurs en aval sont sensibles aux doublons, déduppez du côté consommateur, par exemple par l’identifiant d’événement.
- Des sources de données Python personnalisées utilisées comme sources. La possibilité d’effectuer des lectures exactement une fois dépend de la capacité de votre implémentation source à signaler correctement les décalages et à reprendre correctement à partir de ceux-ci. S’il ne suit pas les décalages, considérez-le comme fonctionnant en mode « au moins une fois » et dédupliquez en aval avec
dropDuplicatessur un ID d’événement ou en vous appuyant sur la sémantique d’upsert basée sur une clé deAUTO CDC.
En règle générale, si tout votre pipeline est Delta vers Delta (tables de streaming et vues matérialisées lisant et écrivant des tables Delta via des flux gérés), vous avez déjà le mode « exactement une fois ». Dès que vous ajoutez un foreach_batch_sink, un récepteur autre que Delta ou une source personnalisée non vérifiée, traitez ce bord spécifique en mode « au moins une fois » et ajoutez à ce niveau une logique d’écriture idempotente ou de déduplication.