Mi az állapotalapú streamelés?

Ez a lap az állapotalapú strukturált streamelési lekérdezéseket ismerteti, beleértve az állapotalapú műveleteket, az optimalizálási javaslatokat, a több állapotalapú operátor láncolását és az állapot-újraegyensúlyozást.

Az állapotalapú strukturált streamelési lekérdezések növekményes frissítéseket igényelnek a köztes állapotinformációkhoz, míg az állapot nélküli strukturált streamelési lekérdezések csak a forrástól a fogadóig feldolgozott sorok adatait követik nyomon. Az állapot nélküli lekérdezésekhez elérhető optimalizálási funkciókért tekintse meg az állapot nélküli streamelési lekérdezések optimalizálása című témakört.

Állapotalapú műveletek

Az állapotalapú műveletek közé tartozik a stream-aggregáció, distincta dropDuplicatesstream-stream illesztések és az egyéni állapotalapú alkalmazások.

Az állapotalapú strukturált streamelési lekérdezésekhez szükséges köztes állapotinformációk váratlan késéshez és éles problémákhoz vezethetnek, ha helytelenül vannak konfigurálva.

A Databricks Runtime 13.3 LTS-ben vagy újabb verzióiban engedélyezheti a változásnapló-ellenőrzőpont-ellenőrzést a RocksDB-vel, hogy csökkentse az ellenőrzőpontok időtartamát és a végpontok közötti késést strukturált streamelési számítási feladatok esetében. A Databricks azt javasolja, hogy engedélyezze a változásnapló-ellenőrzőpontozást az összes strukturált streamelési állapotalapú lekérdezéshez. Lásd: Változásnapló-ellenőrzőpont engedélyezése.

Állapotalapú strukturált streamelési lekérdezések optimalizálása

A Databricks az állapotalapú strukturált streamelési lekérdezésekhez a következőket javasolja:

  • Számítási teljesítményre optimalizált példányokat használj munkásként.
  • Állítsa a shuffle partíciók számát a csomópont magjainak számához arányosan 1–2-szeresére.

Fontos

A rendszer az ellenőrzőpont létrehozásakor rögzíti az shuffle partíciók számát. A módosításnak spark.sql.shuffle.partitions nincs hatása egy olyan streamelési lekérdezésre, amely már rendelkezik ellenőrzőponttal – a lekérdezés továbbra is az eredeti partíciószámot használja. Új partíciószám alkalmazásához új ellenőrzőpont-hellyel kell elindítania a lekérdezést.

A Databricks Runtime 18.0-s vagy újabb verziójában az állapot nélküli streamelési lekérdezések új ellenőrzőpont nélkül támogatják a dinamikus shuffle partíciómódosításokat.

A Databricks Runtime 18 LTS-ben és újabb verziókban anélkül módosíthatja az állapotalapú lekérdezések partíciószámát, hogy elveszítené az ellenőrzőpont állapotát. Tekintse meg az igény szerinti állapot újraparticionálását az állapotalapú streamelési lekérdezésekhez.

  • Állítsa a spark.sql.streaming.noDataMicroBatches.enabled konfigurációt false-re a SparkSession környezetében. Ez megakadályozza, hogy a streamelt mikrokötegmotor feldolgozzon olyan mikro kötegeket, amelyek nem tartalmaznak adatokat. Ha ezt a konfigurációt false-ra állítja, az olyan állapotalapú műveleteket, amelyek vízjeleket vagy feldolgozási időtúllépéseket használnak, nem eredményezhet adatkimenetet, amíg új adatok nem érkeznek, ahelyett, hogy azonnal megtörténne.

A Databricks a RocksDB használatát javasolja változásnapló-ellenőrzőpontokkal az állapotalapú streamek állapotának kezeléséhez. Lásd: A RocksDB állapottároló konfigurálása az Azure Databricksben.

Megjegyzés

Az állapotkezelési séma nem módosítható a lekérdezések újraindítása között. Ha egy lekérdezést az alapértelmezett felügyelettel indítottak el, az állapottároló módosításához újra kell indítania egy új ellenőrzőpont-hellyel.

Több állapotalapú operátor használata a strukturált streamelésben

A Databricks Runtime 13.3 LTS-ben vagy újabb verzióiban az Azure Databricks fejlett támogatást nyújt a strukturált streamelési számítási feladatok állapotalapú operátorai számára. Több állapotalapú operátort is összefűzhet, ami azt jelenti, hogy egy művelet kimenetét, például egy ablakos összesítést egy másik állapotalapú művelethez, például egy illesztéshez táplálhatja.

A Databricks Runtime 16.2-s vagy újabb verzióiban több állapotalapú operátorral rendelkező számítási feladatokban is használható transformWithState . Lásd: Készíts egy egyedi állapotos alkalmazást .transformWithState

Az alábbi példák számos használható mintát mutatnak be.

Fontos

Több állapotalapú operátor használatakor a következő korlátozások érvényesek:

  • Az örökölt egyéni állapotalapú operátorok (FlatMapGroupWithState és applyInPandasWithState) nem támogatottak.
  • Csak a hozzáfűző kimeneti mód támogatott.

Láncolt időablak-összesítés

Python

words = ...  # streaming DataFrame of schema { timestamp: Timestamp, word: String }

# Group the data by window and word and compute the count of each group
windowedCounts = words.groupBy(
    window(words.timestamp, "10 minutes", "5 minutes"),
    words.word
).count()

# Group the windowed data by another window and word and compute the count of each group
anotherWindowedCounts = windowedCounts.groupBy(
    window(window_time(windowedCounts.window), "1 hour"),
    windowedCounts.word
).count()

Scala

import spark.implicits._

val words = ... // streaming DataFrame of schema { timestamp: Timestamp, word: String }

// Group the data by window and word and compute the count of each group
val windowedCounts = words.groupBy(
  window($"timestamp", "10 minutes", "5 minutes"),
  $"word"
).count()

// Group the windowed data by another window and word and compute the count of each group
val anotherWindowedCounts = windowedCounts.groupBy(
  window($"window", "1 hour"),
  $"word"
).count()

Időablak-összesítés két különböző streamben, majd stream-stream ablak illesztése

Python

clicksWindow = clicksWithWatermark.groupBy(
  clicksWithWatermark.clickAdId,
  window(clicksWithWatermark.clickTime, "1 hour")
).count()

impressionsWindow = impressionsWithWatermark.groupBy(
  impressionsWithWatermark.impressionAdId,
  window(impressionsWithWatermark.impressionTime, "1 hour")
).count()

clicksWindow.join(impressionsWindow, "window", "inner")

Scala

val clicksWindow = clicksWithWatermark
  .groupBy(window("clickTime", "1 hour"))
  .count()

val impressionsWindow = impressionsWithWatermark
  .groupBy(window("impressionTime", "1 hour"))
  .count()

clicksWindow.join(impressionsWindow, "window", "inner")

Folyamok közötti időintervallum-illesztés, amelyet időablakos összesítés követ

Python

joined = impressionsWithWatermark.join(
  clicksWithWatermark,
  expr("""
    clickAdId = impressionAdId AND
    clickTime >= impressionTime AND
    clickTime <= impressionTime + interval 1 hour
    """),
  "leftOuter"                 # can be "inner", "leftOuter", "rightOuter", "fullOuter", "leftSemi"
)

joined.groupBy(
  joined.clickAdId,
  window(joined.clickTime, "1 hour")
).count()

Scala

val joined = impressionsWithWatermark.join(
  clicksWithWatermark,
  expr("""
    clickAdId = impressionAdId AND
    clickTime >= impressionTime AND
    clickTime <= impressionTime + interval 1 hour
  """),
  joinType = "leftOuter"      // can be "inner", "leftOuter", "rightOuter", "fullOuter", "leftSemi"
)

joined
  .groupBy($"clickAdId", window($"clickTime", "1 hour"))
  .count()

Állapotkiegyenlítés strukturált streameléshez

Az állapot-kiegyensúlyozás alapértelmezés szerint engedélyezve van a Lakeflow-pipeline-ok összes streaming számítási feladata esetében. A Databricks Runtime 11.3 LTS-ben vagy újabb verziójában az alábbi konfigurációs beállítást állíthatja be a Spark-fürtkonfigurációban az állapot-újraegyensúlyozás engedélyezéséhez:

spark.sql.streaming.statefulOperator.stateRebalancing.enabled true

Az állapot-újraegyensúlyozás előnyös az állapotalapú strukturált streamelési folyamatok számára, amelyek fürt átméretezési eseményeken mennek keresztül. Az állapot nélküli streamelési műveletek a fürt méretének megváltoztatásától függetlenül nem előnyösek.

Megjegyzés

A számítási automatikus skálázás korlátozottan képes a strukturált streaming számítási feladatok fürtméretének csökkentésére. 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.

A klaszter átméretezési eseményei az állapot újraegyensúlyozását váltják ki. A mikro kötegek nagyobb késéssel járhatnak az események újraegyensúlyozása során, mivel az állapot betöltődik a felhőbeli tárolóból az új végrehajtókba.