Bearbetningsgarantier i Lakeflow-pipelines

Nya försök och omkörningar är oundvikliga i alla verkliga pipelines, så den här sidan förklarar vilka bearbetningsgarantier Lakeflow-pipelines ger dig och hur du ser till att de delar du skriver är säkra att köra om.

Overview

Två relaterade egenskaper avgör om det är säkert att köra om en pipeline:

  • Idempotens betyder att en pipeline ger samma resultat oavsett hur många gånger du kör den över samma input. Att köra om efter ett fel, bakåtfylla ett datumintervall två gånger eller manuellt utlösa ett jobb skapar aldrig dubbla rader eller ett korrupt tillstånd.
  • Bearbetningsgaranti beskriver hur många gånger varje post påverkar resultatet. Bearbetning minst en gång garanterar att varje post bearbetas, men ett fel och ett nytt försök kan innebära att vissa poster bearbetas mer än en gång, vilket kan leda till dubbletter. Behandling exakt en gång garanterar att varje post påverkar resultatet som om den hade bearbetats precis en gång, även vid omförsök, utan vare sig dubbletter eller bortfall.

Lakeflow pipelines är som standard idempotenta för de delar de hanterar och ger dig bearbetning exakt en gång i de tabeller de själva hanterar. Det viktiga att förstå är var dessa garantier slutar vara automatiska, så att du kan lägga till rätt skydd vid kanterna av din pipeline.

Så här fungerar det

Lakeflow Pipelines ger bearbetning exakt en gång och idempotens i de flöden de hanterar, och ger dig verktyg för att se till att logiken du skriver också är idempotent.

Exakt en-gångsbehandling för hanterade tabeller

I hanterade tabeller får du som standard bearbetning exakt en gång. Strömningstabeller använder strukturerade strömningskontroller kombinerade med Delta Lakes transaktionella skrivningar: varje mikrobatch committar sina källoffsets och sin utdata tillsammans, så en omprövad batch efter ett misslyckande lyckas antingen fullt ut eller rullas tillbaka och prövas igen, aldrig delvis applicerad två gånger. Detta gäller för filinläsning med Auto Loader samt läsningar från Kafka, Kinesis och Azure Event Hubs, och AUTO CDC upserts, utan att du behöver skriva någon kod.

Om en källa som minst en gång skickar samma post flera gånger, bearbetar pipelinen dem som unika poster och skriver alla till din tabell. Att ta bort dessa dubbletter är ditt ansvar. Se Deduplicera minst en gång.

Idempotens för läsningar följer från samma kontrollpunkter. Auto Loader och streamingtabellskontrollpunkter garanterar att varje källfil eller offset behandlas en gång för tillståndsspårning, så att ombearbeta en pipelineuppdatering efter ett fel återupptas från kontrollpunkten istället för att bearbeta eller hoppa över data. Du får det här genom att använda streamingtabeller över spark.readStream i stället för egenskrivna batchloopar. Se Strömmande tabeller.

Använd AUTO CDC istället för manuellt skriven MERGE

AUTO CDC INTO är inneboende idempotent med avseende på dess keys och sequence_by. Att applicera samma ändringspost två gånger, eller att tillämpa poster i fel ordning, ger samma slutstatus, eftersom pipelinen använder sekvenskolumnen för att avgöra om en inkommande rad faktiskt är nyare än den som lagras:

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;

Om du skriver din egen upsert-logik utanför AUTO CDC (sällsynt, men ibland nödvändigt för komplexa villkor för sammanslagning), basera den på en stabil affärsnyckel och se till att den kan tillämpas två gånger utan problem, till exempel en MERGE ... WHEN MATCHED med order_id som nyckel i stället för en okontrollerad INSERT. För mer information, se AUTO CDC API:erna: Förenkla ändringsdatafångst med hjälp av pipelines.

Se till att dina egna transformationer är idempotenta

För att säkerställa att logiken förblir idempotent när skrivåtgärder körs om, följ dessa två riktlinjer:

  • Undvik icke-deterministiska transformationer i materialiserade vyer. Eftersom en materialiserad vy kan beräknas om helt eller inkrementiellt, undvik funktioner vars utdata beror på när de körs snarare än vad indata är. Till exempel, använd current_timestamp() inte för att beräkna ett affärsvärde som ska förbli fast när det är skrivet; ta tidsstämpeln från källhändelsen eller skicka in den som en parameter så att omberäkningen ger identisk utdata.
  • Utforma fullständiga omladdningar så att de är säkra. En fullständig uppdatering tar bort och beräknar om en tabell från grunden, vilket bara är säkert om varje uppströmskälla fortfarande kan producera hela historiken. Om en uppströmskälla bara exponerar ett rullande ändringsfönster kan en fullständig omläsning av en nedströms AUTO CDC-tabell leda till att historik går förlorad utan att det märks, så planera lagringstiden för källan och ämnet med detta i åtanke.

Kom exakt en gång vid kanterna

Där exacty-once slutar vara automatisk ligger i gränsen för vad pipelinen direkt kontrollerar, såsom skrivningar till externa system. När du skriver till ett externt system bör själva skrivningen vara idempotent, till exempel genom att använda upsert baserat på en nyckel på mottagarsidan, eftersom en mikrobatch som körs om annars kan skriva samma batch två gånger. Följande sink skriver varje partition i batchen från exekutörerna och använder en idempotensnyckel så att en omprövad batch inte dubbelskriver:

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)

Mer information om att skriva till externa system finns i Sinks in Lakeflow pipelines.

Deduplicera minst en gång källor

När en källa kan leverera en post mer än en gång, deduplicera nedströms. Kombinera en vattenstämpel med dropDuplicatesWithinWatermark, som är vattenmärkesmedveten och inte kräver obegränsat tillstånd för att upptäcka dubbletter. Deduplicera utifrån de kolumner som identifierar en händelse unikt. Identiteten kan sträcka sig över flera kolumner när ingen enskild kolumn är unik i sig själv. I följande exempel är ett klicksekvensnummer unikt endast inom sin session, så de två kolumnerna tillsammans identifierar händelsen:

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

Välj dessa kolumner från källans unikhetskontrakt, inte från vad som ser ut att vara unikt i exempeldata. Kolumner som med rätta kan innehålla upprepade värden gör att verkliga händelser försvinner när du behandlar dem som en unik identifierare. En användare som klickar på samma annons två gånger är ett vanligt exempel: att deduplicera utifrån användaren och annonsen filtrerar utan att märkas bort det andra klicket.

AUTO CDCs nyckelbaserade Upsert-semantik kollapsar också dubbletter naturligt, så att routa minst en gång data genom ett AUTO CDC flöde som är nyckelat på en stabil affärsnyckel är ett annat sätt att konvergera mot exakt ett tillstånd.

Limitations

Exakt-en-gång-bearbetning gäller för hanterade Delta-till-Delta-flöden. Behandla följande kanter som minst en gång och lägg till explicit deduplicerings- eller idempotent-skrivlogik där:

  • foreach_batch_sink och anpassade externa skrivningar. Spark garanterar att en batch försöks minst en gång, men en batch som görs om efter en partiell skrivning kan lämna vissa rader synliga två gånger i det externa systemet. Gör den externa skrivningen idempotent, till exempel genom att göra en upsert på en naturlig nyckel eller skriva ett batch-ID som mottagaren kan använda för deduplicering.
  • Kafka som sink. Kafka-topicer stöder inte transaktionella skrivningar med exakt en gång på samma sätt som Delta gör, så en mikrobatch som körs om och skriver till Kafka kan ge upphov till duplicerade meddelanden. Om nedströmskonsumenter är känsliga för dubbletter, dedupeera på konsumentsidan, till exempel via händelse-ID.
  • Anpassade Python-datakällor som används som källor. Huruvida läsoperationer sker exakt en gång beror på om din källimplementering korrekt rapporterar offsetvärden och återupptas från dem. Om den inte spårar offsetar, behandla det som leverans med minst en gång och eliminera dubbletter längre ned i kedjan med dropDuplicates på ett händelse-ID eller genom att förlita dig på den nyckelbaserade upsert-semantiken i AUTO CDC.

Som tumregel, om hela din pipeline är Delta-till-Delta (strömmande tabeller och materialiserade vyer som läser och skriver Delta-tabeller genom hanterade flöden), har du redan exakt en gång. Så snart du lägger till en foreach_batch_sink, en icke-Delta-sänka eller en icke verifierad anpassad källa ska du betrakta just den anslutningen som at-least-once och lägga till logik för idempotenta skrivningar eller deduplicering där.

Ytterligare resurser