Verwerkingsgaranties in Lakeflow-pijpleidingen

Nieuwe pogingen en heruitvoeringen zijn onvermijdelijk in elke echte pipeline, dus op deze pagina wordt uitgelegd welke verwerkingsgaranties Lakeflow-pipelines bieden en hoe u de onderdelen die u schrijft veilig opnieuw kunt uitvoeren.

Overview

Twee gerelateerde eigenschappen bepalen of het opnieuw aanleggen van een pijpleiding veilig is:

  • Idempotentie betekent dat een pijplijn hetzelfde resultaat oplevert, hoe vaak je hem ook over dezelfde invoer draait. Opnieuw uitvoeren na een mislukking, twee keer een datumbereik invullen of een taak handmatig opnieuw activeren creëert nooit dubbele rijen of corrumpeert de status.
  • Verwerkingsgarantie beschrijft hoe vaak elk record het resultaat beïnvloedt. At-least-once-verwerking garandeert dat ieder record wordt verwerkt, maar door een fout en een nieuwe poging kunnen sommige records meer dan één keer worden verwerkt, waardoor duplicaten kunnen ontstaan. Exact-once verwerking zorgt ervoor dat elk record het resultaat beïnvloedt alsof het precies één keer is verwerkt, zelfs over herhalingen, zonder duplicaten en zonder gaten.

Lakeflow-pipelines zijn standaard idempotent voor de onderdelen die ze beheren en bieden exact-once-verwerking in hun eigen beheerde tabellen. Het belangrijkste om te begrijpen is waar die garanties ophouden automatisch te zijn, zodat je de juiste waarborgen aan de randen van je pijplijn kunt toevoegen.

Hoe werkt het?

Lakeflow-pipelines bieden exact-once-verwerking en idempotentie voor de stromen die ze beheren, en bieden je hulpmiddelen om ook de logica die je schrijft idempotent te houden.

Exact-once-verwerking voor beheerde tabellen

In beheerde tabellen krijg je standaard exact-once-verwerking. Streaming-tabellen gebruiken Structured Streaming-checkpoints in combinatie met de transactionele schrijfbewerkingen van Delta Lake: elke micro-batch commit zijn bronoffsets en zijn uitvoer samen, zodat een batch die na een fout opnieuw wordt uitgevoerd ofwel volledig slaagt, ofwel volledig wordt teruggedraaid en opnieuw geprobeerd, en nooit gedeeltelijk twee keer wordt toegepast. Dit geldt voor het inlezen van bestanden met Auto Loader, leesbewerkingen uit Kafka, Kinesis en Azure Event Hubs, en AUTO CDC upserts, zonder dat je code hoeft te schrijven.

Als een bron die minstens één keer hetzelfde record meerdere keren stuurt, verwerkt de pipeline ze als unieke records en schrijft ze allemaal naar jouw tabel. Het verwijderen van die duplicaten is jouw verantwoordelijkheid. Zie Dedupliceren van at-least-once-bronnen.

Idempotentie voor leesbewerkingen vloeit voort uit diezelfde checkpoints. Auto Loader en streaming table checkpoints garanderen dat elk bronbestand of offset eenmaal wordt verwerkt voor state-tracking doeleinden, zodat het herverwerken van een pipeline-update na een mislukking vanaf checkpoint wordt hervat in plaats van data opnieuw te verwerken of over te slaan. Je krijgt dit door streamingtabellen te gebruiken in spark.readStream plaats van handmatig gerolde batchloops. Zie Streaming-tabellen.

Gebruik AUTO CDC in plaats van handgeschreven MERGE

AUTO CDC INTO is van nature idempotent ten opzichte van haar keys en sequence_by. Door hetzelfde wijzigingsrecord twee keer toe te passen, of records uit de volgorde toe te passen, ontstaat dezelfde eindtoestand, omdat de pipeline de sequentiekolom gebruikt om te bepalen of een binnenkomende rij daadwerkelijk nieuwer is dan wat er is opgeslagen:

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;

Als je je eigen upsertlogica buiten AUTO CDC schrijft (zeldzaam, maar soms noodzakelijk bij complexe mergecondities), baseer die dan op een stabiele zakelijke sleutel en zorg ervoor dat die zonder problemen twee keer kan worden toegepast, bijvoorbeeld een MERGE ... WHEN MATCHED op basis van order_id in plaats van een blinde INSERT. Voor meer informatie, zie De AUTO CDC API's: Vereenvoudig wijzigingsdata-opname met pijplijnen.

Houd je eigen transformaties idempotent

Om de logica idempotent te houden bij het opnieuw uitvoeren van schrijfoperaties, volg deze twee richtlijnen:

  • Vermijd niet-deterministische transformaties in gematerialiseerde weergaven. Omdat een gematerialiseerde weergave volledig of incrementeel kan herberekenen, vermijd functies waarvan de output afhangt van wanneer ze draaien in plaats van wat de invoer is. Gebruik bijvoorbeeld niet current_timestamp() om een bedrijfswaarde te berekenen die vast moet blijven zodra deze is geschreven; neem de tijdstempel van het bronevent of geef deze door als parameter zodat herberekening identieke output oplevert.
  • Ontwerp volledige vernieuwingen zo dat ze veilig zijn. Een volledige verversing verwijdert een tabel en bouwt die vanaf nul opnieuw op, wat alleen veilig is als elke bovenliggende bron nog steeds de volledige geschiedenis kan aanleveren. Als een upstream-bron alleen een rollend venster met wijzigingen beschikbaar stelt, kan een volledige verversing van een downstream-AUTO CDC-tabel ongemerkt historische gegevens verliezen, dus stem bron- en topicretentie hierop af.

Ga precies één keer aan de randen

Waar exactly-once niet langer automatisch is, is aan de randen van wat de pijplijn direct controleert, zoals schrijfacties naar externe systemen. Wanneer je naar een extern systeem schrijft, maak de schrijfbewerking zelf dan idempotent, bijvoorbeeld door aan de ontvangende kant een upsert op basis van een sleutel uit te voeren, omdat een opnieuw uitgevoerde micro-batch anders dezelfde batch twee keer kan schrijven. De volgende sink schrijft elke partitie van de batch vanuit de executors weg en gebruikt een idempotentiesleutel zodat een opnieuw uitgevoerde batch niet dubbel wordt weggeschreven:

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)

Raadpleeg Sinks in Lakeflow-pijpleidingen voor meer informatie over het wegschrijven naar externe systemen.

Dedupliceer ten minste één bronnen

Wanneer een bron een record meer dan eens kan leveren, deduplicleer dan downstream. Combineer een watermerk met dropDuplicatesWithinWatermark, dat watermerk-bewust is en geen onbegrensde toestand vereist om duplicaten te detecteren. Verwijder duplicaten op basis van de kolommen die een gebeurtenis uniek identificeren. De identiteit kan meerdere kolommen beslaan wanneer geen enkele kolom uniek op zichzelf is. In het volgende voorbeeld is een kliksequentienummer uniek alleen binnen zijn sessie, dus identificeren de twee kolommen samen het evenement:

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"])
    )

Kies die kolommen uit het uniciteitscontract van de bron, niet uit wat er verschillend uitziet in de voorbeeldgegevens. Kolommen die legitiem kunnen herhalen, laten echte gebeurtenissen vallen als je ze als identiteit behandelt. Een gebruiker die twee keer op dezelfde advertentie klikt is een veelvoorkomend voorbeeld: dedupliceren bij de gebruiker en de advertentie laat stilletjes de tweede klik vallen.

AUTO CDCDe op sleutels gebaseerde Upsert-semantiek laat ook duplicaten natuurlijk instorten, dus het routen van ten minste één data door een AUTO CDC flow die is gebaseerd op een stabiele bedrijfssleutel is een andere manier om te convergeren naar exact één toestand.

Limitations

Exactly-once-verwerking is van toepassing op beheerde Delta-naar-Delta-gegevensstromen. Behandel de volgende randen als ten minste één en voeg daar expliciete deduplicatie- of idempotent-write-logica toe:

  • foreach_batch_sink en aangepaste externe schrijfbewerkingen. Spark garandeert dat een batch minstens één keer wordt geprobeerd , maar een batch die na een gedeeltelijk schrijven opnieuw wordt geprobeerd, kan sommige rijen twee keer zichtbaar laten in het externe systeem. Maak de externe write idempotent, bijvoorbeeld door op een natuurlijke sleutel te upserten of een batch-ID te schrijven waarop de ontvanger kan dedupliceren.
  • Kafka als gootsteen. Kafka-topics ondersteunen geen transactionele exactly-once-schrijfbewerkingen zoals Delta dat wel doet, dus een opnieuw uitgevoerde micro-batch die naar Kafka schrijft kan dubbele berichten produceren. Als downstream consumenten gevoelig zijn voor duplicaten, dedupéer dan aan de consumentenkant, bijvoorbeeld op basis van event ID.
  • Aangepaste Python-databronnen gebruikt als bronnen. Of leesbewerkingen precies één keer plaatsvinden, hangt ervan af of je bronimplementatie offsets correct rapporteert en vanaf die offsets wordt hervat. Als het geen offsets bijhoudt, behandel het dan als at-least-once en dedupliceer downstream met dropDuplicates op basis van een gebeurtenis-ID of door te vertrouwen op de sleutelgebaseerde upsert-semantiek van AUTO CDC.

Als vuistregel geldt: als je volledige pijplijn Delta-naar-Delta is (streamingtabellen en gematerialiseerde weergaven die Delta-tabellen lezen en schrijven via beheerde stromen), dan heb je al exact-once-semantiek. Zodra je een foreach_batch_sink toevoegt, een niet-Delta-sink of een niet-geverifieerde aangepaste bron, beschouw die specifieke verbinding dan als at-least-once en voeg daar logica voor idempotente schrijfbewerkingen of deduplicatie toe.

Aanvullende bronnen