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.
DataFrame-műveletek vagy SQL táblaértékfüggvények használatával lekérdezheti a strukturált streamelési állapot adatait és metaadatait. Ezekkel a függvényekkel megfigyelheti a strukturált streamelési állapotalapú lekérdezések állapotadatait, amelyek a figyeléshez és a hibakereséshez hasznosak lehetnek.
Az állapotadatok vagy metaadatok lekérdezéséhez olvasási hozzáféréssel kell rendelkeznie a streamelési lekérdezés ellenőrzőpont-elérési útjának eléréséhez. A cikkben ismertetett függvények írásvédett hozzáférést biztosítanak az állapotadatokhoz és metaadatokhoz. Az állapotinformációk lekérdezéséhez csak kötegelt olvasási szemantikát szabad használni.
Feljegyzés
A Lakeflow-folyamatok, streamelőtáblák és materializált nézetek állapotadatai nem kérdezhetők le. Az állapotadatok nem kérdezhetők le kiszolgáló nélküli számítással vagy standard hozzáférési móddal konfigurált számítással.
Követelmények
- Használja az alábbi számítási konfigurációk egyikét:
- Databricks Runtime 16.3 és újabb verziók standard hozzáférési móddal konfigurált számításon.
- A Databricks Runtime 14.3 LTS és újabb verziója dedikált vagy elkülönítés nélküli hozzáférési móddal konfigurált számítási egységen.
- A streamelési lekérdezés által használt ellenőrzőpont elérési útjához való olvasási hozzáférés.
Strukturált streamelési állapottároló olvasása
A támogatott Databricks-futtatókörnyezetekben végrehajtott strukturált streamelési lekérdezések állapottárolási információit olvashatja. Alkalmazza a következő szintaxist:
Python
df = (spark.read
.format("statestore")
.load("/checkpoint/path"))
Scala
val df = spark.read
.format("statestore")
.load("/checkpoint/path")
SQL
SELECT * FROM read_statestore('/checkpoint/path')
Állapotolvasó API-beállításai és sémája
A formátumbeállítások teljes listáját statestore az Állapottárban találja.
A kimeneti adatok a következő sémával rendelkeznek:
| Oszlop | Típus | Leírás |
|---|---|---|
key |
Struct (az állapotkulcsból származtatott további típus) | Az állapotalapú operátorrekord kulcsa az állapot-ellenőrzőpontban. |
value |
Struct (az állapotértékből származtatott további típus) | Az állapotalapú operátorrekord értéke az állapot-ellenőrzőpontban. |
partition_id |
Egész szám | Az állapotalapú operátorrekordot tartalmazó állapot-ellenőrzőpont partíciója. |
A Databricks Runtime 16.4 LTS és újabb verzióiban, amikor a(z) readChangeFeed beállítás értéke true, a kimeneti adatok sémája a következő:
| Oszlop | Típus | Leírás |
|---|---|---|
batch_id |
Hosszú | Az a kötegazonosító, amelyhez az állapotváltozás tartozik. |
change_type |
Sztring | A köteg által alkalmazott módosítás típusa: update beszúrásokhoz és frissítésekhez, delete törlésekhez. |
key |
Struct (az állapotkulcsból származtatott további típus) | Az állapotalapú operátorrekord kulcsa az állapot-ellenőrzőpontban. |
value |
Struct (az állapotértékből származtatott további típus) | Az állapotalapú operátorrekord értéke az állapot-ellenőrzőpontban.
null olyan rekordok esetében, ahol change_type van delete. |
partition_id |
Egész szám | Az állapotalapú operátorrekordot tartalmazó állapot-ellenőrzőpont partíciója. |
Lásd read_statestore táblaértékelt függvény.
Strukturált streamelési állapot változásainak olvasása
A Databricks Runtime 16.4 LTS-en és újabb verziókon érhető el. Ha azt szeretné megtekinteni, hogyan változik az állapot a mikrokötegek között ahelyett, hogy egyetlen mikroköteg teljes állapotát nézné meg, állítsa a readChangeFeed értékét true értékre, és adja meg a(z) changeStartBatchId értéket. Igény szerint adja meg a kívánt értéket changeEndBatchId. A lehetőségek teljes listáját az Állapottárban találja.
Például a köteg 2 állapotváltozásainak olvasása a legújabb véglegesített kötegen keresztül:
Python
df = (spark.read
.format("statestore")
.option("readChangeFeed", True)
.option("changeStartBatchId", 2)
.load("<checkpointLocation>")
)
Scala
val df = spark.read
.format("statestore")
.option("readChangeFeed", true)
.option("changeStartBatchId", 2)
.load("<checkpointLocation>")
SQL
SELECT * FROM read_statestore(
'<checkpointLocation>',
readChangeFeed => true,
changeStartBatchId => 2
);
A kimeneti séma további batch_id és change_type oszlopokat is tartalmaz. A teljes sémáért tekintse meg az Állapotolvasó API-beállításait és sémáját.
Strukturált streamelési állapot metaadatainak olvasása
Elérhető a Databricks Runtime 14.3 LTS vagy újabb verziójában. A strukturált streamelési lekérdezésekhez az állapot metaadatait olvashatja el:
Python
df = (spark.read
.format("state-metadata")
.load("<checkpointLocation>"))
Scala
val df = spark.read
.format("state-metadata")
.load("<checkpointLocation>")
SQL
SELECT * FROM read_state_metadata('/checkpoint/path')
A visszaadott adatok sémája a következő:
| Oszlop | Típus | Leírás |
|---|---|---|
operatorId |
Egész szám | Az állapotalapú streaming operátor egész azonosítója. |
operatorName |
Sztring | Az állapotalapú streamelő operátor neve. |
stateStoreName |
Sztring | Az operátor állapottárolójának neve. |
numPartitions |
Egész szám | Az állapottároló partícióinak száma. |
minBatchId |
Hosszú | Az állapot lekérdezéséhez elérhető minimális kötegazonosító. |
maxBatchId |
Hosszú | Az állapot lekérdezéséhez elérhető maximális kötegazonosító. |
Feljegyzés
A minBatchId és maxBatchId által megadott kötegazonosító értékek az ellenőrzőpont megírásának idejének állapotát tükrözik. A régi kötegek automatikusan törlődnek a mikroköteg végrehajtásával, így az itt megadott érték nem garantált, hogy továbbra is elérhető marad.
Lásd read_state_metadata táblaértékelt függvény.
Példa: Stream-stream illesztés egyik oldalának lekérdezése
A stream-stream illesztés bal oldalának lekérdezéséhez használja az alábbi szintaxist:
Python
left_df = (spark.read
.format("statestore")
.option("joinSide", "left")
.load("/checkpoint/path"))
Scala
val leftDf = spark.read
.format("statestore")
.option("joinSide", "left")
.load("/checkpoint/path")
SQL
SELECT * FROM read_statestore(
'/checkpoint/path',
joinSide => 'left'
);
Példa: Adatfolyam állapottárolójának lekérdezése több állapotalapú operátorral
Ez a példa az állapot metaadat-olvasójával gyűjti össze egy streamelési lekérdezés metaadatait több állapotalapú operátorral, majd a metaadat-eredményeket használja az állapotolvasó beállításaiként.
Az állapot metaadat-olvasója az ellenőrzőpont elérési útját használja egyetlen lehetőségként, ahogyan az alábbi szintaxisbeli példában is látható:
Python
df = (spark.read
.format("state-metadata")
.load("<checkpointLocation>"))
Scala
val df = spark.read
.format("state-metadata")
.load("<checkpointLocation>")
SQL
SELECT * FROM read_state_metadata('/checkpoint/path')
Az alábbi táblázat az állapottár metaadatainak példakimenetét mutatja be:
| operátorazonosító | operátor neve | állapotTárolóNév | numPartitions | minBatchId | maxBatchId |
|---|---|---|---|---|---|
| 0 | stateStoreSave | alapértelmezett | 200 | 0 | 13 |
| 1 | dedupeWithinWatermark | alapértelmezett | 200 | 0 | 13 |
Az operátor eredményének dedupeWithinWatermark lekéréséhez kérje le az állapotolvasót a operatorId beállítással, ahogyan az alábbi példában is látható:
Python
left_df = (spark.read
.format("statestore")
.option("operatorId", 1)
.load("/checkpoint/path"))
Scala
val leftDf = spark.read
.format("statestore")
.option("operatorId", 1)
.load("/checkpoint/path")
SQL
SELECT * FROM read_statestore(
'/checkpoint/path',
operatorId => 1
);