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.
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.enabledkonfigurációtfalse-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ótfalse-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ésapplyInPandasWithState) 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.