Produktionsöverväganden för strukturerad direktuppspelning

Kör Structured Streaming-arbetsbelastningar i produktion som schemalagda Lakeflow-jobb på Azure Databricks. Se Lakeflow Jobs.

Databricks rekommenderar att du alltid konfigurerar följande:

  • Ta bort onödig kod från notebook-filer som returnerar resultat, till exempel display och count.
  • Kör inte Structured Streaming-arbetsbelastningar med allmän beräkning. Schemalägg alltid strömmar som Lakeflow-jobb med jobbberäkning.
  • Schemalägg Lakeflow Jobs i Continuousläge. Detta avser schemaläggningsfunktionen för Azure Databricks Jobs, inte Structured Streaming triggerintervallet.
  • Aktivera inte automatisk skalning för beräkning för strukturerade direktuppspelningsjobb.

Vissa arbetsbelastningar drar nytta av följande:

Databricks introducerade Lakeflow-pipelines för att minska komplexiteten i hanteringen av produktionsinfrastruktur för strukturerade strömningsarbetsbelastningar. Databricks rekommenderar att du använder Lakeflow-pipelines för nya pipelines för strukturerad direktuppspelning. Se Spark deklarativa datapipelines.

Kommentar

Automatisk skalning av beräkningskraft har begränsningar för att skala ned klusterstorleken för strukturerad strömmande arbetsbelastning. Databricks rekommenderar att du använder deklarativa Spark-pipelines i Lakeflow med utökad automatisk skalning för strömningsarbetslaster. Se Optimera användning av Lakeflow-pipelinekluster med automatisk skalning.

:::obs Serverlös beräkning

På serverlös beräkning stöds endast Trigger.AvailableNow() och Trigger.Once() . Databricks rekommenderar Trigger.AvailableNow().

För kontinuerlig direktuppspelning vid serverlös beräkning använder du Utlöst kontra kontinuerligt pipelineläge i kontinuerligt läge.

Se Begränsningar för direktuppspelning.

:::

Minska latensen för operativ streaming

Operativa strömningsarbetsbelastningar tar in, transformerar och agerar på data i nästan realtid. Vanliga exempel är bedrägeridetektion, avvikelsedetektion, personalisering samt realtidsövervakning och varningar, där fördröjd hantering direkt påverkar affärsresultaten. Låg latens för dessa arbetsbelastningar innebär vanligtvis tiotals till hundratals millisekunder, även om många team sätter service level agreements (SLA) i sekunderintervallet för att ta hänsyn till variation vid högre procentiler.

För lägsta end-to-end-latens bör du använda realtidsläge, som ger en end-to-end-latens på under en sekund som mest och omkring 300 millisekunder i typiska fall. Se koncept för realtidsläge.

När realtidsläge inte passar din arbetsbelastning minskar följande bästa praxis latensen för mikrobatch-strukturerad streaming:

  • Output-läge: Använd update-läge där dina frågeoperatorer och sink stödjer det. Uppdateringsläget skickar ut uppdaterade rader efter varje utlösare och fortsätter att uppdatera dem tills vattenstämpeln upphör att gälla, så se till att din nedströmsmottagare är idempotent för att hantera uppdaterade resultat. Använd append-läge för arbetsbelastningar som uppdateringsläget inte stöder, till exempel ström-ström-sammanslagningar, eller när du kan utelämna sent inkommande data. Använd inte komplett läge för låg latens. Se Välj ett utdataläge för Structured Streaming.
  • Trigger: Använd en processingTime trigger med intervall 0 , som startar nästa mikrobatch så snart den föregående är klar och ny data är tillgänglig. Detta ger den lägsta mikrobatchlatensen, men ökar kostnaderna för molnlagrings-API:er. Använd inte AvailableNow, Once, eller Continuous för operativa arbetsbelastningar. Se Konfigurera utlösarintervall för strukturerad direktuppspelning.
  • Vattenstämpel: Ställ in vattenstämpeln tillräckligt länge för att inkludera den data som anländer sent så att din arbetsbelastning inte får minska. Vattenmärket styr hur länge frågan accepterar händelsetidsdata som kommer i fel ordning innan den kasserar dem och rensar tillståndet, så ett vattenmärke som är för kort kasserar i tysthet giltiga försenade poster. Inom den begränsningen ger en kortare watermark lägre latens och kräver mindre tillståndsinformation, medan en längre watermark kan hantera mer försenade data på bekostnad av latens och tillståndsinformation. En liten multipel av din latens-SLA, som 2x, är en rimlig utgångspunkt för tuning. Se Tillämpa vattenmärken för att styra databehandlingens tröskelvärden.
  • Källor och mottagare: Läs från datakällor med låg latens, som meddelandebussar (Apache Kafka, Amazon Kinesis, Apache Pulsar eller Google Cloud Pub/Sub) eller ändringsdataflöden från Delta Lake- och Apache Iceberg-tabeller. Skriv till sänkor med låg latens och hög genomströmning, såsom meddelandebussar, operativa databaser eller foreach sänkor. Designa sänkoperationer så att de är idempotenta så att nedströms konsumenter hanterar dubbletter och sent ankommande data.
  • Tillståndshantering och checkpointing: För tillståndsfulla frågor använder du RocksDB:s tillståndslager, som krävs för både checkpointing av ändringslogg och asynkron checkpointing av tillstånd. Aktivera kontrollpunkter för ändringsloggen för att endast bevara inkrementella tillståndsändringar. När tillståndskontroll är flaskhalsen i din batchvaraktighet aktiverar du asynkron tillståndskontroll så att skrivningar av kontrollpunkter kan överlappa med nästa mikrobatch, när du har gått igenom begränsningarna för felåterställning och klusterändring. Ge varje fråga en egen checkpoint-katalog i hållbar molnlagring. Se Konfigurera RocksDB-tillståndslager på Azure Databricks, Asynkron kontrollpunktslagring för tillståndsfulla frågor och Kontrollpunkter i Structured Streaming.
  • Offsethantering: För att minska den latens som uppstår vid skapande av kontrollpunkter för offset i kontinuerliga strömmar aktiverar du asynkron förloppsspårning, som uppdaterar loggar för offset och commit utan att blockera databehandlingen. Det är inte kompatibelt med utlösarna AvailableNow eller Once. Se Asynkron förloppsspårning.
  • Storage hops: Håll beräkningarna inom en enda strömningspipeline där det är möjligt. Att dela upp logiken över flera jobb eller pipelines lägger till lagringshopp som ökar latensen.

Utforma strömmande arbetsflöden med beredskap för fel

Databricks rekommenderar att du alltid konfigurerar direktuppspelningsjobb för att automatiskt starta om vid fel. Vissa funktioner, inklusive schemaevolution, kräver att Structured Streaming-arbetsbelastningar gör automatiska återförsök. Se Konfigurera strukturerade strömmande jobb för att starta om strömmande frågor vid fel.

Vissa åtgärder som foreachBatch erbjuder garantier för minst en gång i stället för exakt en gång. Säkerställ att din bearbetningspipeline är idempotent för dessa åtgärder. Se även Använda foreachBatch för att skriva till godtyckliga datamottagare.

Kommentar

När en fråga startas om, bearbetas mikrobatchen som planerades under föregående körning. Om jobbet misslyckades på grund av ett minnesfel eller om du avbröt ett jobb manuellt på grund av en överdimensionerad mikrobatch kan du behöva skala upp beräkningen för att kunna bearbeta mikrobatchen.

Om du ändrar konfigurationer mellan körningar gäller dessa konfigurationer för den första nya batchen som planeras. Se Återställning efter ändringar i en Structured Streaming-fråga.

När ett jobb försöker igen

Du kan schemalägga flera aktiviteter som en del av ett Azure Databricks jobb. När du konfigurerar ett jobb med den kontinuerliga utlösaren kan du inte ange beroenden mellan aktiviteter.

Du kan välja att schemalägga flera strömmar i ett enda jobb med någon av följande metoder:

  • Flera uppgifter: Definiera ett jobb med flera aktiviteter som kör strömmande arbetsbelastningar med hjälp av den kontinuerliga utlösaren.
  • Flera frågor: Definiera flera strömmande frågor i källkoden för en enda uppgift.

Du kan också kombinera dessa strategier. I följande tabell jämförs dessa metoder.

Strategi Flera uppgifter Flera frågor
Hur delas datorkapacitet? Databricks rekommenderar att du distribuerar jobb med lämplig storlek för varje direktuppspelningsaktivitet. Du kan också dela beräkning mellan aktiviteter. Alla frågor delar samma beräkning. Du kan också tilldela frågor till scheduler-pooler.
Hur hanteras återförsök? Alla uppgifter måste misslyckas innan jobbet omstartar. Uppgiften körs igen om någon sökfråga misslyckas.

Mer information om hur du arbetar med flera uppgifter eller frågor finns i Köra flera frågor för strukturerad direktuppspelning i samma kluster.

Konfigurera strukturerade strömningsjobb för att återstarta strömmande frågor vid fel

Databricks rekommenderar att du konfigurerar alla strömmande arbetsbelastningar med hjälp av den kontinuerliga utlösaren. Se Kör jobb kontinuerligt.

Den kontinuerliga utlösaren har följande beteende som standard:

  • Förhindrar mer än en samtidig körning av jobbet.
  • Startar en ny körning när en tidigare körning misslyckas.
  • Använder den exponentiella "backoff"-metoden för återförsök.

Databricks rekommenderar att du alltid använder jobbberäkning i stället för all-purpose compute när du schemalägger arbetsflöden. Vid jobbfel och återförsök distribueras nya beräkningsresurser.

Kommentar

Databricks rekommenderar att du inte använder streamingQuery.awaitTermination() eller spark.streams.awaitAnyTermination(). Se När du ska använda awaitTermination().

När du ska använda awaitTermination()

streamingQuery.awaitTermination() och spark.streams.awaitAnyTermination() blockerar den aktuella tråden tills en strömmande fråga avslutas. Om du vill använda dessa funktioner beror på din körningsmiljö.

För Lakeflow-jobb ska du inte använda streamingQuery.awaitTermination() eller spark.streams.awaitAnyTermination(). Dessa funktioner är inte nödvändiga eftersom Jobs service automatiskt förhindrar en körning från att slutföras när en streamingfråga är aktiv. Båda funktionerna blockerar notebook-celler från att slutföras och förhindrar att jobbtjänsten spårar strömningsfrågan, vilket stör kvarvarande mått och jobbmeddelanden.

Använd awaitTermination() i följande fall:

Användningsfall Beteende
Interaktiva anteckningsböcker för generell databehandling awaitTermination() ser till att cellen fortsätter köra, gör att du kan observera frågetillståndet och ser till att problem visas i anteckningsboksresultaten.
Lokala miljöer och utvecklingsmiljöer När du kör ett Spark-program lokalt avslutas processen när huvudtråden är klar. Anropa awaitTermination() för att hålla programmet vid liv tills strömningsfrågan har slutförts eller misslyckats.
Felöverföring till drivrutinen Utan awaitTermination() kan det hända att ett fel i en strömmande frågeprocess i en icke-jobbkontext inte sprids till den anropande tråden. Frågan kan misslyckas tyst, vilket gör fel svårare att identifiera och diagnostisera. Att anropa awaitTermination() utlöser frågeundantaget i drivrutinen igen.