Strukturált streamelési ellenőrzőpontok

Az ellenőrzőpontok és az előre írt naplók együttműködnek a strukturált streamelési számítási feladatok feldolgozási garanciáinak biztosításához. Az ellenőrzőpont nyomon követi a lekérdezést azonosító információkat, beleértve az állapotadatokat és a feldolgozott rekordokat. Amikor törli a fájlokat egy ellenőrzőpont-könyvtárban, vagy új ellenőrzőpont-helyre vált, a lekérdezés következő futtatása újraindul.

Az ellenőrzőpont-címtár a következőket tartalmazza:

  • Eltolások: Az egyes mikrofolyamatokban feldolgozott forráseltolások. Ez lehetővé teszi, hogy a lekérdezés pontosan onnan folytatódjon, ahol az adatok feldolgozása nélkül abbahagyta.
  • Véglegesítések: Egy rekord arról, hogy mely mikro-batch-ek lettek commitálva a célba, lehetővé téve az exactly-once szemantikát.
  • Állapot: Állapotalapú lekérdezések (aggregációk, stream-stream-illesztések, deduplikációk és egyéni állapotalapú operátorok, például transformWithState) esetén az ellenőrzőpont az állapotalapú operátor, az állapotséma és az állapottár-szolgáltató által kezelt ellenőrzőponttal ellátott állapottár-tartalom metaadatait tárolja.
  • Metaadatok: A lekérdezés azonosításához használt egyedi lekérdezésazonosító. A konfigurációs beállítások az eltolásnapló részeként vannak tárolva.

Minden lekérdezésnek más ellenőrzőponttal kell rendelkeznie. Több lekérdezésnek soha nem szabad ugyanazt a helyet megosztania.

Megjegyzés

Ez a cikk a streamelési lekérdezések strukturált streamelési ellenőrzőpontjait ismerteti. A Unity Catalog-kötetek nem streaming adattábláinak végrehajtási terveinek csonkításával kapcsolatos információért tekintse meg a DataFrame.checkpoint() című részt.

Ellenőrzőpontok engedélyezése strukturált streamelési lekérdezésekhez

A streamelési lekérdezés futtatása előtt meg kell adnia a checkpointLocation beállítást, ahogyan az alábbi példában is látható:

Python

(df.writeStream
  .option("checkpointLocation", "/Volumes/catalog/schema/volume/path")
  .toTable("catalog.schema.table")
)

Scala

df.writeStream
  .option("checkpointLocation", "/Volumes/catalog/schema/volume/path")
  .toTable("catalog.schema.table")

Megjegyzés

Bizonyos kimeneti végpontok, mint például a jegyzetfüzetek display() kimenete és a memory végpont, automatikusan létrehoznak egy ideiglenes ellenőrzőpont helyét, ha kihagyja ezt a beállítást. Ezek az ideiglenes ellenőrzőpontok nem biztosítják a hibatűrést vagy az adatkonzisztencia garanciáját, és előfordulhat, hogy nem lesznek megfelelően megtisztítva. A Databricks azt javasolja, hogy mindig adja meg az ellenőrzőpont-helyet ezekhez a célpontokhoz.

A helyreállítás a strukturált streamelési lekérdezés módosításai után

Az ugyanabból az ellenőrzőpont-helyről történő újraindítások között korlátozva van, hogy a streamelési lekérdezések milyen módosításokat hajthatnak végre.

Az általánosan új ellenőrzőpontot igénylő módosítások közé tartozik a bemeneti források száma vagy típusa, az előfizetett Kafka-témakörök vagy az Automatikus betöltő elérési útjai, az állapotalapú művelettípusok, az állapotséma és a kimeneti fogadó típusa.

Az általánosan biztonságos módosítások közé tartoznak a szűrők hozzáadása vagy eltávolítása, a sebességkorlátok módosítása, a triggerintervallumok és a felhasználó által definiált függvénylogika mapGroupsWithState frissítése (bár a szemantika változhat).

A következő szakasz azokat a módosításokat ismerteti, amelyek vagy nem engedélyezettek, vagy a módosítás hatása nem megfelelően van meghatározva, ahol:

  • Az engedélyezett kifejezés azt jelenti, hogy elvégezheti a megadott módosítást, de a lekérdezéstől és a változástól függ, hogy a hatás szemantikája megfelelően van-e definiálva.
  • A nem engedélyezett kifejezés azt jelenti, hogy nem szabad elvégeznie a megadott módosítást, mivel az újraindított lekérdezés valószínűleg kiszámíthatatlan hibákkal fog meghiúsulni.
  • sdf streamelt DataFrame/Dataset-et jelöl, amely a sparkSession.readStream-val/vel jön létre.

A strukturált streamelési lekérdezések változásainak típusai

  • A bemeneti források számának vagy típusának módosítása: Ez alapértelmezés szerint nem engedélyezett, mert a strukturált streamelés a források helyét azonosítja a lekérdezési tervben. Ha bekapcsolja a forráselnevezést, átrendezheti a meglévő forrásokat, és új forrásokat vehet fel új ellenőrzőpont nélkül. Lásd: Streamforrások módosítása a forrásfejlődéssel.

  • A bemeneti források paramétereinek változásai: Az, hogy ez engedélyezett-e, és hogy a változás szemantikája megfelelően van-e definiálva, a forrástól és a lekérdezéstől függ, beleértve az olyan belépési vezérlőket is, mint a maxFilesPerTrigger vagy maxOffsetsPerTrigger. Íme néhány példa:

    • A sebességkorlátok hozzáadása, törlése és módosítása engedélyezett:

      spark.readStream.format("kafka").option("subscribe", "article")
      

      felhasználóként a(z)

      spark.readStream.format("kafka").option("subscribe", "article").option("maxOffsetsPerTrigger", ...)
      

      További információ: Strukturált streamelési kötegméret konfigurálása az Azure Databricksben

    • Az előfizetett cikkek és fájlok módosítása általában nem engedélyezett, mivel az eredmények kiszámíthatatlanok: spark.readStream.format("kafka").option("subscribe", "article")spark.readStream.format("kafka").option("subscribe", "newarticle")

  • Az eseményindító időközének változásai: Módosíthatja az eseményindítókat a növekményes kötegek és az időintervallumok között. Lásd: Az eseményindító futások közötti időközeinek módosítása.

  • A kimeneti fogadó típusának változásai: A fogadók egyes kombinációi közötti változások engedélyezettek. Ezt eseti alapon kell ellenőrizni. Szeretnénk ismertetni néhány példát.

    • A fájl-adatvégpont Kafka adatvégpontra történő kapcsolása engedélyezett. A Kafka csak az új adatokat fogja látni.
    • Kafka és fájl szinkronizálása nem engedélyezett.
    • A Kafka sink "foreach"-re történő módosítása, vagy fordítva, megengedett.
  • A kimeneti fogadó paramétereinek változásai: A fogadótól és a lekérdezéstől függ, hogy ez engedélyezett-e, és hogy a változás szemantikája megfelelően van-e definiálva. Szeretnénk ismertetni néhány példát.

    • A fájlelső kimeneti könyvtárának módosítása nem engedélyezett: sdf.writeStream.format("parquet").option("path", "/somePath")sdf.writeStream.format("parquet").option("path", "/anotherPath")
    • A kimeneti témakör módosítása engedélyezett: sdf.writeStream.format("kafka").option("topic", "topic1")sdf.writeStream.format("kafka").option("topic", "topic2")
    • A felhasználó által definiált foreach-fogadó (azaz a ForeachWriter kód) módosítása engedélyezett, de a módosítás szemantikája a kódtól függ.
  • A vetítési/szűrési/ térképszerű műveletek változásai: Bizonyos esetekben engedélyezettek. Például:

    • A szűrők hozzáadása/törlése engedélyezett: sdf.selectExpr("a") és sdf.where(...).selectExpr("a").filter(...) között.
    • Az azonos kimeneti sémával rendelkező projekciók változásai engedélyezettek: sdf.selectExpr("stringColumn AS json").writeStreamsdf.select(to_json(...).as("json")).writeStream.
    • A különböző kimeneti sémákkal rendelkező vetítések változásai feltételesen engedélyezettek: a sdf.selectExpr("a").writeStream átalakítás sdf.selectExpr("b").writeStream-re csak akkor lehetséges, ha a kimeneti csatorna engedélyezi a séma változtatását "a"-ről "b"-ra.
  • Állapotalapú műveletek változásai: A streamelési lekérdezések egyes műveleteinek állapotadatokat kell fenntartaniuk az eredmény folyamatos frissítéséhez. A strukturált streamelés automatikusan ellenőrzi az állapotadatokat a hibatűrő tárolóra (például DBFS, Azure Blob Storage) és visszaállítja azokat az újraindítás után. Ez azonban feltételezi, hogy az állapotadatok sémája az újraindítások során változatlan marad. Ez azt jelenti, hogy a streamelési lekérdezés állapotalapú műveleteinek módosításai (vagyis hozzáadásai, törlései vagy sémamódosításai) nem engedélyezettek az újraindítások között. Íme azoknak az állapotalapú műveleteknek a listája, amelyek sémája nem módosítható az újraindítások között az állapot helyreállítása érdekében:

    • Stream-összesítés: Például sdf.groupBy("a").agg(...). A csoportosítási kulcsok vagy összesítések számának vagy típusának módosítása nem engedélyezett.
    • Stream-deduplikáció: Például sdf.dropDuplicates("a"). A csoportosítási kulcsok vagy összesítések számának vagy típusának módosítása nem engedélyezett.
    • Stream-stream illesztés: Például sdf1.join(sdf2, ...) (azaz mindkét bemenet sparkSession.readStream-kel van létrehozva). A séma vagy az egyenrangú oszlopok módosítása nem engedélyezett. Az illesztés típusának (külső vagy belső) módosítása nem engedélyezett. Az illesztés feltételének egyéb változásai nem definiálva vannak.
    • Tetszőleges állapotalapú művelet: Például. sdf.groupByKey(...).mapGroupsWithState(...)sdf.groupByKey(...).flatMapGroupsWithState(...) A felhasználó által definiált állapot sémájának és az időtúllépés típusának módosítása nem engedélyezett. A felhasználó által definiált állapotleképezési függvényben bármilyen változás megengedett, de a változás szemantikai hatása a felhasználó által definiált logikától függ. Ha valóban támogatni szeretné az állapotséma módosításait, akkor a séma áttelepítését támogató kódolási/dekódolási sémával explicit módon bájtokra kódolhatja/dekódolhatja az összetett állapotadat-struktúrákat. Ha például Avro kódolású bájtként menti az állapotot, módosíthatja az Avro-state-sémát a lekérdezés újraindítása között, mivel ez visszaállítja a bináris állapotot.

Fontos

Az állapotalapú operátorok dropDuplicates()dropDuplicatesWithinWatermark() nem indulhatnak újra az állapotséma kompatibilitási ellenőrzése miatt a számítási hozzáférési módok közötti váltáskor.

A dedikált és az izoláció nélküli hozzáférési módok közötti váltás engedélyezett. A standard és a kiszolgáló nélküli hozzáférési módok közötti váltás engedélyezett. Ne próbáljon meg más hozzáférési mód kombinációk között váltani.

A hiba elkerülése érdekében ne módosítsa az ezeket az operátorokat tartalmazó streamelési lekérdezések számítási hozzáférési módját.

Streamelési források módosítása a forrásfejlődéssel

A strukturált streamelés alapértelmezés szerint a lekérdezési tervben elfoglalt pozíciójuk alapján azonosítja a forrásokat, például 0, 1, 2stb. A bemeneti források számának vagy sorrendjének bármilyen módosítása megszakítja az ellenőrzőpontok kompatibilitását, és új ellenőrzőpontot igényel. A forrásfejlődés lehetővé teszi, hogy stabil, felhasználó által definiált neveket rendeljen minden streamforráshoz, hogy az ellenőrzőpont állapotának elvesztése nélkül átrendezhesse, hozzáadhassa vagy eltávolíthassa a lekérdezések forrásait.

A forrásséma evolúciójához a Databricks Runtime 18.2-es vagy újabb verziója szükséges.

Szükséges konfiguráció

A forrásfejlődés engedélyezéséhez állítsa be a következő Spark-konfigurációt:

  • spark.sql.streaming.queryEvolution.enableSourceEvolution: Amikor truea lekérdezésben lévő összes streamforrást explicit módon el kell nevezni az .name() API használatával. Az alapértelmezett érték false.

Állítsa be a konfigurációt a streamelési lekérdezés definiálása előtt:

spark.conf.set("spark.sql.streaming.queryEvolution.enableSourceEvolution", "true")

Elnevezési szabályok

  • A nevek csak alfanumerikus karaktereket és aláhúzásjeleket ([a-zA-Z0-9_]+) tartalmazhatnak.
  • Minden forrásnévnek egyedinek kell lennie egy lekérdezésen belül.
  • Ha a forrásfejlődés engedélyezve van, minden streamelő forrásnak rendelkeznie kell egy névvel. A névtelen források hibát okoznak UNNAMED_STREAMING_SOURCES_WITH_ENFORCEMENT .

Források átrendezés, hozzáadás és eltávolítás

Az alábbi módosítások biztonságosak a lekérdezések újraindítása során ugyanazzal az ellenőrzőponttal:

  • Források átrendezése: Indítsa újra a lekérdezést egy másik forrásrenddel. A neve alapján minden forrás az utoljára véglegesített eltolásától folytatódik, és nem módosítja az ellenőrzőpont állapotát.
  • Új források hozzáadása: Indítsa újra a lekérdezést egy új forrással. Az új források az elejétől kezdik a feldolgozást, a meglévő források pedig az utolsó eltolási pozíciójuktól folytatják.
  • Források eltávolítása: Indítsa újra a lekérdezést forrás nélkül. A forrás véglegesen törlődik az ellenőrzőpontról. Az eltávolított forrás nem adható hozzá újra ugyanazzal a névvel.

Example

Használja a(z) .name() elemet a(z) DataStreamReader elemen, mielőtt meghívja a(z) .load() vagy a(z) .table() elemet:

Python

orders_us = (spark.readStream
  .name("orders_us")
  .table("catalog.schema.orders_us")
)

orders_eu = (spark.readStream
  .name("orders_eu")
  .table("catalog.schema.orders_eu")
)

all_orders = orders_us.union(orders_eu)

Scala

val ordersUS = spark.readStream
  .name("orders_us")
  .table("catalog.schema.orders_us")

val ordersEU = spark.readStream
  .name("orders_eu")
  .table("catalog.schema.orders_eu")

val allOrders = ordersUS.union(ordersEU)

Limitations

  • A forrás elnevezéséhez új ellenőrzőpont szükséges. A forrásfejlődés nem engedélyezhető a nélküle létrehozott meglévő ellenőrzőponton.
  • A forrásfejlődés az ellenőrzőpontok visszafordíthatatlan változása. Miután bekapcsolta a lekérdezés forrásfejlődését, nem kapcsolhatja ki a forrásfejlődést, és nem használhatja újra ugyanazt az ellenőrzőpontot.
  • A forrásnevek állandóak. Egy forrás átnevezéséhez először törölje, majd adja hozzá új néven. Az átnevezett forrásfolyamatok az elejétől kezdve.