A strukturált streamelés gyártási megfontolásai

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 és count.
  • 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 Continuous mó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:

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.

:::

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.