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.
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.
-
sdfstreamelt DataFrame/Dataset-et jelöl, amely asparkSession.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
maxFilesPerTriggervagymaxOffsetsPerTrigger. Í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
ForeachWriterkód) módosítása engedélyezett, de a módosítás szemantikája a kódtól függ.
- A fájlelső kimeneti könyvtárának módosítása nem engedélyezett:
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")éssdf.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ássdf.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.
- A szűrők hozzáadása/törlése engedélyezett:
Á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 bemenetsparkSession.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.
-
Stream-összesítés: Például
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: Amikortruea lekérdezésben lévő összes streamforrást explicit módon el kell nevezni az.name()API használatával. Az alapértelmezett értékfalse.
Á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.