Garanties de traitement dans les Lakeflow pipelines
Les essais et rediffusions sont inévitables dans tout pipeline réel, donc cette page explique les garanties de traitement que les Lakeflow pipelines offrent et comment garder les parties que vous écrivez en sécurité pour les relancer.
Présentation
Deux propriétés connexes déterminent si la réexécution d'un pipeline est sûre :
- L'idempotence signifie qu'un pipeline produit le même résultat, quel que soit le nombre de fois où vous l'exécutez sur la même entrée. La réexécution après un échec, le remplissage d'une plage de dates deux fois ou le redéclenchement manuel d'un Job ne crée jamais de lignes en double et ne corrompt jamais l'état.
- La garantie de traitement décrit combien de fois chaque enregistrement affecte le résultat. Le traitement At-least-once (au moins une fois) garantit que chaque enregistrement est traité, mais une défaillance suivie d'une nouvelle tentative peut entraîner le traitement de certains enregistrements plus d'une fois, ce qui présente un risque de doublons. Le traitement Exactly-once (exactement une fois) garantit que chaque enregistrement affecte le résultat comme s'il avait été traité précisément une fois, même en cas de nouvelles tentatives, sans doublons ni lacunes.
LakeFlow pipelines sont default pour les éléments qu'ils gèrent et vous offrent un traitement « exactly-once » au sein de leurs propres tables gérées. L'important est de comprendre où ces garanties cessent d'être automatiques, afin que vous puissiez ajouter les protections appropriées aux extrémités de votre pipeline.
Comment ça fonctionne
Les Lakeflow pipelines assurent un traitement « exactly-once » (exactement une fois) et l'idempotence pour les flux qu'ils gèrent, et vous fournissent des outils pour rendre également idempotente la logique que vous écrivez.
Traitement « exactement une fois » pour les tables gérées
Au sein des tables gérées, vous bénéficiez par default d'un traitement exactement une fois. Les tables de streaming utilisent des points de contrôle Structured Streaming combinés aux écritures transactionnelles de Delta Lake : chaque micro-batch commit ses offsets sources et sa sortie ensemble, de sorte qu'un batch relancé après un échec réussit entièrement ou est entièrement annulé et relancé, sans jamais être appliqué partiellement deux fois. Cela s'applique à l'ingestion de fichiers Auto Loader, aux lectures Kafka, Kinesis et Azure Event Hubs, ainsi qu'aux upserts AUTO CDC, sans aucun code de votre part.
Si une source « au moins une fois » envoie le même enregistrement plusieurs fois, le pipeline les traite comme des enregistrements uniques et les écrit tous dans votre table. Il vous incombe de supprimer ces doublons. Voir Dédupliquer les sources « au moins une fois ».
L'idempotence pour les lectures découle de ces mêmes points de contrôle. Les points de contrôle d'Auto Loader et des tables de streaming garantissent que chaque fichier source ou décalage est traité une seule fois à des fins de suivi d'état ; ainsi, le retraitement d'une mise à jour de pipeline après un échec reprend à partir du point de contrôle plutôt que de retraiter ou de sauter des données. Vous obtenez cela en utilisant des tables de streaming sur spark.readStream au lieu de boucles de batch créées manuellement. Voir Tables de streaming.
Utilisez AUTO CDC au lieu d'un MERGE écrit manuellement
AUTO CDC INTO est intrinsèquement idempotent par rapport à ses keys et sequence_by. L'application deux fois du même enregistrement de modification, ou l'application d'enregistrements dans le désordre, 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 ce qui est stocké :
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 (rare, mais parfois nécessaire pour des conditions de Merge complexes), basez-la sur une clé métier stable et assurez-vous 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 APIs AUTO CDC : simplifiez la capture de données modifiées avec des pipelines.
Assurez l'idempotence de vos propres transformations
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, n’utilisez pas
current_timestamp()pour compute une valeur métier qui devrait rester fixe une fois écrite ; Prenez le Timestamp de l’événement source ou passez-le comme parameter afin que le recalcul produise une sortie identique. - Concevez des full refresh sécurisées. Une full refresh supprime et recalcule une table à partir de zéro, ce qui n'est sûr que si chaque source en amont peut toujours produire l'historique complet. Si une source en amont n'expose qu'une fenêtre glissante de changements, un full refresh d'une table
AUTO CDCen aval peut entraîner une perte silencieuse de l'historique ; concevez donc la rétention des sources et des rubriques en gardant cela à l'esprit.
Obtenir un traitement « exactly-once » aux extrémités
Le traitement « exactly-once » cesse d'être automatique aux limites de ce que le pipeline contrôle directement, comme les écritures vers des systèmes externes. Lorsque vous effectuez une diffusion vers un système externe, rendez l'écriture elle-même idempotente, par exemple en effectuant un upsert par clé du côté récepteur, car un micro-batch relancé pourrait sinon écrire le même batch deux fois. Le récepteur suivant écrit chaque partition du batch à partir des exécuteurs et utilise une clé d'idempotence afin qu'un batch relancé ne soit pas écrit deux fois :
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 vers des systèmes externes, consultez Sinks in Lakeflow pipelines.
Dédupliquer les sources « at-least-once »
Lorsqu’une source peut transmettre un enregistrement plus d’une fois, effectuez une déduplication en aval. Combinez un filigrane avec dropDuplicatesWithinWatermark, qui est compatible avec les filigranes et ne nécessite pas d’état illimité pour détecter les doublons. Dédupliquez sur les colonnes qui identifient 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 clic n’est unique qu’au sein de sa session ; les deux colonnes identifient donc l’événement ensemble :
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 à partir du contrat d'unicité de la source, et non à partir de ce qui semble distinct 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 traitez comme l'identité. Un utilisateur cliquant deux fois sur la même publicité est un exemple courant : la déduplication sur l'utilisateur et la publicité supprime silencieusement le second clic.
AUTO CDCLes sémantiques d’upsert basées sur les clés suppriment également les doublons naturellement. Ainsi, le routage des données « at-least-once » via un flux AUTO CDC basé sur une clé métier stable est un autre moyen de converger vers un état « exactly-once ».
Limitations
Le traitement « exactly-once » s'applique aux flux Delta-to-Delta managés. Traitez les segments suivants comme au moins une fois et ajoutez-y une logique explicite de déduplication ou d’écriture idempotente :
foreach_batch_sinket écritures externes personnalisées. Spark garantit qu'une tentative de traitement de batch est effectuée au moins une fois, mais un batch relancé après une écriture partielle peut rendre 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 batch sur lequel le récepteur peut effectuer une déduplication.- Kafka comme un évier. Les sujets Kafka ne supportent pas exactement les écritures transactionnelles une fois comme le fait Delta, donc une écriture micro-batch réessayée sur Kafka peut produire des messages doublés. Si les consommateurs en aval sont sensibles aux doublons, déduppez du côté consommateur, par exemple par l’identifiant d’événement.
- Sources de données Python personnalisées utilisées comme sources. Le fait que les lectures soient effectuées exactement une fois dépend de la capacité de votre implémentation de source à signaler et à reprendre correctement à partir des offsets. S’il ne suit pas les offsets, traitez-le comme « at-least-once » et dédupliquez en aval avec
dropDuplicatessur un ID d’événement ou en vous appuyant sur les sémantiques d’upsert basées sur les clés deAUTO CDC.
En règle générale, si tout votre pipeline est Delta-à-Delta (streaming tables et vues matérialisées lisant et écrivant des tables Delta via des flux gérés), vous avez déjà exactement une seule fois. Au moment où vous ajoutez un foreach_batch_sink, un puits non-Delta ou une source personnalisée non vérifiée, considérez cette arête spécifique comme au moins une fois et ajoutez une logique d’écriture idempotente ou de déduplication à cet endroit.