Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Futtasson éles Structured Streaming számítási feladatokat ütemezett Lakeflow Jobsként az Azure Databricksben. Lásd Lakeflow Jobs.
A Databricks azt javasolja, hogy mindig konfigurálja a következőket:
- Távolítsa el a szükségtelen kódot olyan jegyzetfüzetekből, amelyek eredményeket adnak vissza, például
displayéscount. - Ne futtasson strukturált adatfolyam-feldolgozási munkaterheléseket általános célú számítási erőforrásokon. Mindig a feladatokhoz tartozó számítási erőforrás használatával ütemezze az adatfolyamokat Lakeflow Jobsként.
- A Lakeflow-feladatok ütemezése
Continuousmódban. Ez a Azure Databricks feladatok ütemezési funkcióra vonatkozik, nem a strukturált streamelési próbálkozás időközére. - Ne engedélyezze az automatikus skálázást a strukturált streamelési feladatok számítási feladataihoz.
Egyes számítási feladatok a következőkből profitálnak:
- Konfigurálása a RocksDB-állapottárnak az Azure Databricks-en
- Állapotalapú lekérdezések aszinkron állapot-ellenőrzőpontja
- Aszinkron folyamatkövetés
A Databricks lakeflow-folyamatokat vezetett be, hogy csökkentse a strukturált streamelési számítási feladatok éles infrastruktúrájának kezelésének összetettségét. A Databricks a Lakeflow-folyamatokat javasolja az új strukturált streamelési folyamatokhoz. Lásd: Spark deklaratív adatfeldolgozási folyamatok.
Megjegyzés
A számítási erőforrások automatikus skálázása korlátozott a strukturált streamelési számítási feladatok fürt méretének csökkentésében. A Databricks a Lakeflowban elérhető, továbbfejlesztett automatikus skálázással rendelkező Spark deklaratív adatfeldolgozási folyamatait ajánlja streaming számítási feladatokhoz. Lásd A Lakeflow-adatcsatorna fürt kihasználásának optimalizálása automatikus skálázással.
:::megjegyzés Kiszolgáló nélküli számítás
Kiszolgáló nélküli számításon csak Trigger.AvailableNow() és Trigger.Once() támogatott. A Databricks javasolja Trigger.AvailableNow().
A kiszolgáló nélküli számítás folyamatos streameléséhez használja a Trigger vagy folyamatos folyamat módot folyamatos módban.
Lásd a streamelési korlátozásokat.
:::
Csökkentse a késleltetést az operatív streaminghez
Az operatív streaming munkaterhelések közel valós időben olvassák be, alakítják át és használják fel az adatokat. Gyakori példák a csalásészlelés, az anomáliák észlelése, a személyre szabás, valamint a valós idejű megfigyelés és riasztás, ahol a késleltetett feldolgozás közvetlenül befolyásolja az üzleti eredményeket. Az alacsony késleltetés ezeknél a feladatterheléseknél jellemzően néhány tíz–néhány száz milliszekundumot jelent, bár sok csapat a szolgáltatási szintre vonatkozó megállapodásokat (SLA-kat) másodperces tartományban határozza meg, hogy figyelembe vegye a magasabb percentiliseknél tapasztalható ingadozást.
A legalacsonyabb végponttól végpontig tartó késleltetéshez valós idejű módot használjunk, amely végponttól végpontig tartó késleltetést ér el egy másodperc alatt, míg a gyakori esetekben körülbelül 300 milliszekundot ér el. Lásd a valós idejű mód koncepciókat.
Ha a valós idejű mód nem illik a munkaterheléshez, a következő legjobb gyakorlatok csökkentik a késleltetést a mikro-batch Structured Streaming esetén:
- Kimeneti mód: Használjon frissítési módot, ha a lekérdezési operátorok és a fogadó támogatják azt. Az Update mód minden trigger után frissített sorokat küld, és folyamatosan frissíti őket, amíg a watermark érvényessége le nem jár, ezért az utófeldolgozó nyelőt idempotenssé kell tenni, hogy kezelni tudja a frissített eredményeket. Használj append módot olyan munkaterhelésekhez, amelyeket a frissítési mód nem támogat, például stream-stream csatlakozások, vagy amikor el lehet dobni a későn érkező adatokat. Ne használd a teljes módot alacsony késleltetés miatt. Lásd: Kimeneti mód kiválasztása strukturált streameléshez.
-
Trigger: Használj egy
processingTimeintervallummal rendelkező triggert0, amely azonnal elindítja a következő mikro-adagot, amint az előző befejezi és új adatok elérhetők. Ez a legalacsonyabb mikro-tételes késleltetést biztosítja, de növeli a felhőalapú tárolási API költségeit. Ne használdAvailableNow,Once, vagyContinuousaz üzemeltetési terhelésekhez. Lásd: Strukturált streamelési eseményindítók időközeinek konfigurálása. - Vízjel: Állítsa be a vízjelet elég hosszúra ahhoz, hogy lefedje a késve érkező adatokat, amelyeket a munkafolyamat nem veszíthet el. A vízjel szabályozza, hogy a lekérdezés mennyi ideig fogadja el a nem sorrendben érkező, eseményidő szerinti adatokat, mielőtt eldobná őket és eltávolítaná a tárolt állapotot, ezért a túl rövid vízjel észrevétlenül eldobja a késve érkező érvényes rekordokat. Ebben a korláton belül a rövidebb vízjel csökkenti a késleltetést és kevesebb állapotot tart fenn, míg a hosszabb vízjel több késői adatot tűr, de a késleltetés és állapot árán. A késleltetési SLA-d kis többszöröse, például a 2-szerese, ésszerű kiindulópont a finomhangoláshoz. Lásd: Vízjelek alkalmazása az adatfeldolgozási küszöbértékek szabályozásához.
-
Források és elnyelők: Olvass alacsony késleltetésű forrásokból, például üzenetbuszokról (Apache Kafka, Amazon Kinesis, Apache Pulsar vagy Google Cloud Pub/Sub), vagy változtatásadatfolyamokról a Delta Lake és Apache Iceberg táblákból. Írj alacsony késleltetésű, nagy áteresztőképességű elnyelőkbe, mint például üzenetbuszokat, működési adatbázisokat vagy
foreachelsigőket. Tervezze úgy a nyelő oldali műveleteket, hogy azok idempotensek legyenek, így a downstream fogyasztók képesek kezelni a duplikátumokat és a késve érkező adatokat. - Állapot és ellenőrzőpont: Állapotos lekérdezésekhez használd a RocksDB állapot tárolót, amely mind a changelog ellenőrzőponthoz, mind aszinkron állapotellenőrzéshez szükséges. Engedélyezze a változásnapló-ellenőrzőpontozást, hogy csak a növekményes állapotváltozások kerüljenek megőrzésre. Ha az állapot-ellenőrzőpontok készítése jelenti a szűk keresztmetszetet a batch feldolgozási idejében, a hibából való helyreállítással és a klaszter átméretezésével kapcsolatos korlátok áttekintése után engedélyezd az aszinkron állapot-ellenőrzőpontkészítést, hogy az ellenőrzőpontok írása átfedésbe kerülhessen a következő mikroköteggel. Adj minden lekérdezésnek saját ellenőrzőpont-könyvtárat tartós felhőtárolóban. Lásd: A RocksDB állapottár konfigurálása az Azure Databricksben, Aszinkron állapot-ellenőrzőpont-készítés állapottal rendelkező lekérdezésekhez és A Structured Streaming ellenőrzőpontjai.
-
Offset menedzsment: Az offset ellenőrzőpontok folyamatos áramlási késleltetésének csökkentése érdekében engedélyezzük az aszinkron előrehaladás követését, amely frissíti az offset és commit naplókat anélkül, hogy akadályozná az adatfeldolgozást. Nem kompatibilis az
AvailableNoworOncetriggerekkel. Lásd az aszinkron folyamatkövetést. - Tárolóugrók: A számítást egy egyetlen streaming pipeline-en belül tartsd, ahol lehetséges. A logika több feladat vagy pipeline közötti megosztása tárolóugrásokat eredményez, amelyek növelik a késleltetést.
Streamelési számítási feladatok megtervezése, amelyek számolnak a hibákkal
A Databricks azt javasolja, hogy a streamingfeladatokat mindig úgy állítsa be, hogy hiba esetén automatikusan újrainduljanak. Bizonyos képességek, például a séma fejlődése megkövetelik, hogy a strukturált streamelési számítási feladatok automatikusan újrapróbálkoznak. Lásd: Strukturált streamelési feladatok konfigurálása a streamelési lekérdezések sikertelen újraindításához.
Egyes műveletek, például foreachBatch a pontos garancia helyett legalább egyszer nyújtanak garanciát. Ezeknek a műveleteknek a végrehajtása során győződjön meg arról, hogy a feldolgozási folyamat idempotens. Lásd: A foreachBatch használata tetszőleges adatelnyelőkbe való íráshoz című témakört.
Megjegyzés
Amikor egy lekérdezés újraindul, az előző futtatás során tervezett mikrorészlet kerül feldolgozásra. Ha a feladat memóriahiba miatt meghiúsult, vagy egy túlméretezett mikroköteg miatt manuálisan megszakított egy feladatot, előfordulhat, hogy fel kell skáláznia a számítást a mikro köteg sikeres feldolgozásához.
Ha a futtatások közötti konfigurációkat módosítja, ezek a konfigurációk az első tervezett új kötegre vonatkoznak. Lásd: Helyreállítás a strukturált streamelési lekérdezés változásai után.
Amikor egy feladat újrapróbálkozik
Több tevékenységet is ütemezhet egy Azure Databricks feladat részeként. Ha egy feladatot a folyamatos eseményindítóval konfigurál, nem állíthat be függőségeket a tevékenységek között.
Az alábbi módszerek egyikével több streamet ütemezhet egy feladatba:
- Több tevékenység: Több feladattal rendelkező feladat definiálása, amely folyamatos eseményindítóval futtat streamelési számítási feladatokat.
- Több lekérdezés: Több streamelési lekérdezés definiálása egyetlen tevékenység forráskódjában.
Ezeket a stratégiákat kombinálhatja is. Az alábbi táblázat ezeket a megközelítéseket hasonlítja össze.
| Stratégia | Több feladat | Több lekérdezés |
|---|---|---|
| Hogyan osztják meg a számítást? | A Databricks azt javasolja, hogy az egyes streamelési feladatokhoz megfelelő méretű számítási feladatokat helyezzen üzembe. Igény szerint megoszthatja a számítást a tevékenységek között. | Minden lekérdezés ugyanazt a számítást használja. Igény szerint hozzárendelhet lekérdezéseket az ütemezőkészletekhez. |
| Hogyan kezelik az újrapróbálkozásokat? | A feladat újrapróbálkozása előtt minden tevékenységnek sikertelennek kell lennie. | A feladat újrapróbálkozza, ha valamelyik lekérdezés meghiúsul. |
A több tevékenység vagy lekérdezés használatával kapcsolatos további részletekért lásd: Több strukturált streamelési lekérdezés futtatása ugyanazon a fürtön.
Strukturált streamelési feladatok konfigurálása a streamelési lekérdezések sikertelen újraindításához
A Databricks azt javasolja, hogy az összes streamelési számítási feladatot konfigurálja a folyamatos eseményindítóval. Tekintse meg a feladatok folyamatos futtatása című témakört.
A folyamatos eseményindító alapértelmezés szerint a következő viselkedéssel rendelkezik:
- Megakadályozza, hogy egy feladat egyszerre többször is fusson.
- Új futtatás indítása, ha egy korábbi futtatás meghiúsul.
- Exponenciális visszalépést használ az újrapróbálkozáshoz.
A Databricks azt javasolja, hogy a munkafolyamatok ütemezése során mindig a feladatok számítását használja a teljes körű számítás helyett. Feladathiba és újrapróbálkozás esetén új számítási erőforrások kerülnek telepítésre.
Megjegyzés
A Databricks azt javasolja, hogy ne használja streamingQuery.awaitTermination() vagy spark.streams.awaitAnyTermination(). Lásd: Mikor érdemes használniawaitTermination().
Mikor érdemes használni? awaitTermination()
streamingQuery.awaitTermination() és spark.streams.awaitAnyTermination() tiltsa le az aktuális szálat, amíg egy streamelési lekérdezés le nem fejeződik. A függvények használata a végrehajtási környezettől függ.
A Lakeflow-feladatokhoz ne használja streamingQuery.awaitTermination() vagy spark.streams.awaitAnyTermination(). Ezek a függvények nem szükségesek, mert a Feladatok szolgáltatás automatikusan megakadályozza a futtatás végrehajtását, amikor egy streamelési lekérdezés aktív. Mindkét függvény megakadályozza a jegyzetfüzetcellák befejeződését, és megakadályozzák a Jobs szolgáltatást abban, hogy nyomon kövesse a streaming lekérdezést, ami megzavarja a hátralékmetrikákat és a feladatértesítéseket.
Használja awaitTermination() a következő esetekben:
| Felhasználási eset | Magatartás |
|---|---|
| Interaktív jegyzetfüzetek a teljes körű számításhoz |
awaitTermination() Működteti a cellát, lehetővé teszi a lekérdezési állapot megfigyelését, és biztosítja, hogy a hibák láthatóvá váljanak a jegyzetfüzet kimenetében. |
| Helyi és fejlesztési környezetek | Spark-program helyi futtatásakor a folyamat a fő szál befejeződésekor lép ki. Hívja meg awaitTermination() , hogy tartsa életben a programot, amíg a streamelési lekérdezés be nem fejeződik vagy meghiúsul. |
| Hiba propagálása az illesztőprogramra | E nélkül awaitTermination() előfordulhat, hogy a nem feladathoz kötött környezetben egy streamelési lekérdezési hiba nem terjed át a hívó szálra. A lekérdezés csendesen meghiúsulhat, ami megnehezíti a hibák észlelését és diagnosztizálását. A awaitTermination() meghívása ismét lekérdezési kivételt vált ki az illesztőprogramban. |