Strukturált streamelési állapotinformációk olvasása

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
);