Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Spouštějte úlohy strukturovaného streamování v produkčním prostředí jako naplánované úlohy Lakeflow na Azure Databricks. Podívejte se na Úlohy Lakeflow.
Databricks doporučuje, abyste vždy nakonfigurovali následující:
- Odeberte nepotřebný kód z poznámkových bloků, který by vrátil výsledky, například
displayacount. - Nespouštějte úlohy strukturovaného streamování pomocí univerzálních výpočetních prostředků. Vždy naplánujte streamy jako úlohy Lakeflow pomocí výpočetních prostředků úloh.
- Plánování úloh Lakeflow pomocí
Continuousrežimu Toto se týká funkce plánování úloh Azure Databricks, nikoli intervalu spuštění strukturovaného streamování. - Nepovolujte automatické škálování pro výpočetní prostředky pro úlohy strukturovaného streamování.
Některé úlohy můžou využívat tyto výhody:
- Konfigurovat úložiště stavů RocksDB na Azure Databricks
- Asynchronní vytváření kontrolních bodů stavu pro stavové dotazy
- Asynchronní sledování průběhu
Databricks zavedla kanály Lakeflow, které snižují složitost správy produkční infrastruktury pro úlohy strukturovaného streamování. Databricks doporučuje používat kanály Lakeflow pro nové kanály strukturovaného streamování. Viz deklarativní kanály Sparku.
Poznámka:
Automatické škálování má omezení při zmenšování velikosti clusteru pro práce se strukturovaným streamováním. Databricks doporučuje používat deklarativní kanály Sparku ve službě Lakeflow s vylepšeným automatickým škálováním pro úlohy streamování. Viz Optimalizace využití clusteru kanálu Lakeflow pomocí automatického škálování.
:::Poznámka: Výpočetní prostředky bez serveru
Na serverless výpočtech se podporují jenom Trigger.AvailableNow() a Trigger.Once(). Databricks doporučuje Trigger.AvailableNow().
Pro průběžné streamování na bezserverové výpočetní technice použijte režim spuštěný událostmi vs. průběžný režim pipeline v průběžném režimu.
Viz omezení streamování.
:::
Snížení latence pro operační streamování
Provozní streamové úlohy přijímají a transformují data a provádějí nad nimi akce téměř v reálném čase. Běžné příklady zahrnují detekci podvodů, detekci anomálií, personalizaci a monitorování a upozornění v reálném čase, kde zpožděné zpracování přímo ovlivňuje obchodní výsledky. Nízká latence u těchto pracovních zátěží obvykle znamená desítky až stovky milisekund, i když mnoho týmů nastavuje dohody o úrovni služeb (SLA) v rozmezí sekund, aby zohlednily variabilitu ve vyšších percentilech.
Pro nejnižší latenci mezi koncovými body použijte režim v reálném čase, který dosahuje latence mezi koncovými body pod jednu sekundu v krajních případech a kolem 300 milisekund v běžných případech. Viz koncepty režimu v reálném čase.
Když režim v reálném čase nevyhovuje vaší pracovní zátěži, následující osvědčené postupy snižují latenci u mikrodávkového strukturovaného streamování:
- Výstupní režim: Použijte režim aktualizace, kde je podporují dotazovací operátory a spotřebič. Režim aktualizace po každém triggeru vydává aktualizované řádky a aktualizuje je, dokud vodoznak nevyprší, takže nastavte svůj downstream sink idempotent pro zpracování aktualizovaných výsledků. Používejte režim append pro úlohy, které režim update nepodporuje, například spojení streamů, nebo když můžete zahodit opožděně přicházející data. Nepoužívejte režim Complete pro nízkou latenci. Viz Výběr výstupního režimu prostrukturovaného streamování .
-
Trigger: Použijte
processingTimetrigger s intervalem0, který spustí další mikro-batch hned, jakmile skončí předchozí a jsou k dispozici nová data. To zajišťuje nejnižší mikro-batch latenci, ale zvyšuje náklady na API cloudového úložiště. NepoužívejteAvailableNow,Once, aniContinuouspro provozní pracovní zátěže. Viz Konfigurace intervalů triggeru strukturovaného streamování. - Vodoznak: Nastavte vodoznak tak, aby byl dostatečně dlouhý a zahrnoval i opožděně doručená data, která vaše úloha nesmí zahodit. Vodoznak určuje, jak dlouho dotaz přijímá data času události přicházející mimo pořadí, než je zahodí a odstraní stav, takže příliš krátký vodoznak bez upozornění zahazuje platné opožděné záznamy. V rámci tohoto omezení kratší vodoznak snižuje latenci a zachovává menší stav, zatímco delší vodoznak toleruje více pozdních dat na úkor latence a stavu. Malý násobek vašeho SLA latence, například 2x, je rozumný výchozí bod pro ladění. Viz Použití vodoznaků pro řízení prahových hodnot zpracování dat.
-
Zdroje a odběry: Čtěte z nízkolatenčních zdrojů, jako jsou sběrnice zpráv (Apache Kafka, Amazon Kinesis, Apache Pulsar nebo Google Cloud Pub/Sub), nebo měňte datové kanály z tabulek Delta Lake a Apache Iceberg. Zapisujte do nízkolatenčních sinků s vysokou propustností, jako jsou sběrnice zpráv, provozní databáze nebo
foreachsinky. Navrhněte výstupní operace tak, aby byly idempotentní, aby navazující systémy dokázaly zpracovat duplicitní i opožděně doručená data. - Stav a vytváření kontrolních bodů: Pro dotazy se stavem použijte úložiště stavu RocksDB, které je nezbytné jak pro vytváření kontrolních bodů protokolu změn, tak pro asynchronní vytváření kontrolních bodů stavu. Povolte checkpointování changelogu pro ukládání pouze přírůstkových změn stavu. Je-li ukládání stavu do kontrolních bodů úzkým hrdlem ovlivňujícím dobu trvání dávky, povolte asynchronní checkpointování stavu, aby se zápisy checkpointů překrývaly se zpracováním další mikrodávky, a to po seznámení se s omezeními týkajícími se obnovy po selhání a změny velikosti clusteru. Každému dotazu dejte vlastní adresář kontrolních bodů v odolném cloudovém úložišti. Viz Configure RocksDB state store on Azure Databricks, Asynchronous state checkpointing for stateful queries a Structured Streaming checkpoints.
-
Správa offsetů: Chcete-li snížit latenci způsobenou ukládáním kontrolních bodů offsetů v nepřetržitých streamech, povolte asynchronní sledování průběhu, které aktualizuje offsety a protokoly commitů bez blokování zpracování dat. Není kompatibilní se spouštěči
AvailableNowneboOnce. Podívejte se na sledování asynchronního průběhu. - Přeskoky úložiště: Pokud je to možné, udržujte výpočet v rámci jediné streamovací pipeline. Rozdělení logiky mezi více úloh nebo kanálů vede k dalším mezikrokům při ukládání dat, které zvyšují latenci.
Návrh úloh streamování tak, aby očekával selhání
Databricks doporučuje, abyste vždy nakonfigurovali úlohy streamování tak, aby se automaticky restartovala při selhání. Některé funkce, včetně evoluce schématu, vyžadují, aby se úlohy Structured Streaming automaticky opakovaly. Viz Konfigurace úloh strukturovaného streamování pro restartování dotazů streamování při selhání.
Některé operace jako foreachBatch poskytují záruky alespoň jednou namísto přesně jednou. U těchto operací se ujistěte, že váš zpracovatelský řetězec je idempotentní. Viz Použití příkazu foreachBatch k zápisu do libovolných datových jímek.
Poznámka:
Když se dotaz restartuje, zpracuje se mikrodávka naplánovaná během předchozího spuštění. Pokud vaše úloha selhala kvůli chybě nedostatku paměti nebo jste ručně zrušili úlohu kvůli nadměrné mikrodávce, možná budete muset vertikálně navýšit kapacitu výpočetních prostředků, aby bylo možné úspěšně zpracovat mikrodávku.
Pokud změníte konfigurace mezi spuštěními, tyto konfigurace se vztahují na první plánovanou dávku. Viz Obnovení po změnách v dotazu strukturovaného streamingu.
Když se úloha opakuje
Jako součást úlohy v Azure Databricks můžete naplánovat více úkolů. Když konfigurujete úlohu pomocí průběžného triggeru, nemůžete nastavit závislosti mezi úkoly.
Pomocí jednoho z následujících přístupů můžete naplánovat více datových proudů v jedné úloze:
- Více úkolů: Definujte úlohu s více úlohami, které spouštějí úlohy streamování pomocí průběžného triggeru.
- Více dotazů: Definujte více streamovaných dotazů ve zdrojovém kódu pro jeden úkol.
Tyto strategie můžete také kombinovat. Následující tabulka porovnává tyto přístupy.
| Strategie | Více úkolů | Více dotazů |
|---|---|---|
| Jak se výpočetní funkce sdílí? | Databricks doporučuje nasazení výpočetních prostředků adekvátní velikosti pro každý streamovací úkol. Volitelně můžete sdílet výpočetní prostředky napříč úkoly. | Všechny dotazy sdílejí stejný výpočetní výkon. Dotazy můžete volitelně přiřadit fondům plánovače. |
| Jak se zpracovávají opakované pokusy? | Všechny úkoly musí selhat, než dojde k opakování úkolu. | Úloha se opakuje, pokud některý dotaz selže. |
Další podrobnosti o práci s více úlohami nebo dotazy najdete v tématu Spouštění více dotazů strukturovaného streamování ve stejném clusteru.
Konfigurace úloh strukturovaného streamování pro restartování dotazů streamování při selhání
Databricks doporučuje nakonfigurovat všechny streamovací úlohy pomocí průběžného spouštěče. Viz Průběžné spouštění úloh.
Průběžná aktivační událost má ve výchozím nastavení následující chování:
- Zabraňuje souběžnému spuštění úlohy více než jednou.
- Spustí nové spuštění v případě, že předchozí spuštění selže.
- Pro opakování používá exponenciální ústup.
Databricks doporučuje při plánování pracovních postupů vždy používat výpočetní prostředky úloh místo výpočetních prostředků pro všechny účely. Při selhání úlohy a opakování nasadíte nové výpočetní prostředky.
Poznámka:
Databricks doporučuje, abyste je nepoužíli streamingQuery.awaitTermination() nebo spark.streams.awaitAnyTermination(). Viz Kdy použít awaitTermination().
Kdy použít awaitTermination()
streamingQuery.awaitTermination() a spark.streams.awaitAnyTermination() zablokujte aktuální vlákno, dokud se dotaz streamování neukončil. To, jestli se mají tyto funkce používat, závisí na vašem spouštěcím prostředí.
Pro úlohy Lakeflow nepoužívejte streamingQuery.awaitTermination() nebo spark.streams.awaitAnyTermination(). Tyto funkce nejsou nezbytné, protože služba pro úlohy automaticky zabraňuje dokončení procesu, když je streamovací dotaz aktivní. Obě funkce blokují buňky poznámkového bloku od dokončení a brání službě úloh ve sledování toku streamovacích dotazů, což narušuje metriky backlogu a oznámení úloh.
Použijte awaitTermination() v následujících případech:
| Případ použití | Chování |
|---|---|
| Interaktivní poznámkové bloky na výpočetních prostředcích pro všechny účely |
awaitTermination() udržuje buňku spuštěnou, umožňuje sledovat stav dotazu a zajistit, že se v výstupu poznámkového bloku zobrazí selhání. |
| Místní a vývojová prostředí | Při místním spuštění programu Spark se proces po dokončení hlavního vlákna ukončí. Zavolejte awaitTermination(), aby program zůstal aktivní, dokud se dotaz ke streamování nedokončí nebo se nezdaří. |
| Šíření selhání do ovladače | Bez awaitTermination() toho se selhání streamovacího dotazu v nekontextu úlohy nemusí rozšířit do volajícího vlákna. Dotaz může bezobslužně selhat, což znesnadňuje zjišťování a diagnostiku selhání. Volání awaitTermination() v ovladači opět vyvolá výjimku při dotazu. |