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 a vízjelekkel kapcsolatos fogalmakat ismerteti, és javaslatokat biztosít a vízjelek gyakori állapotalapú streamelési műveletekben való használatára.
A streamelési lekérdezések idővel állapotadatokat halmoznak fel. A vízjelek automatikusan eltávolítják a régi állapotadatokat a memóriahibák és a nagyobb feldolgozási késés elkerülése érdekében.
Mi az a vízjel?
A feldolgozás során a strukturált streamelés megőrzi az állapotot a mikrokötegek között. A streaming lekérdezések állapotot használnak az eredmények növekményes frissítéséhez, ahelyett, hogy minden egyes mikroköteg után mindent újraszámolnának. A vízjelek szabályozzák azt a küszöbértéket, amikor egy lekérdezés leállítja az állapotentitások feldolgozását.
Az állami szervezetek gyakori példái a következők:
- Összesítések egy időablakban.
- Egyedi kulcsok két stream közötti illesztésben.
Ha vízjelet szeretne deklarálni egy streamelési DataFrame-en, adjon meg egy időbélyegmezőt és egy késési küszöbértéket. Az új adatok érkezésekor az állapotkezelő nyomon követi a megadott mező legutóbbi időbélyegét, és csak a késési küszöbértéken belüli rekordokat dolgozza fel.
A lekérdezések mindig feldolgozzák a küszöbértéken belülre érkező rekordokat. A lekérdezések továbbra is feldolgozhatják a küszöbértéken kívülre érkező rekordokat, de ez nem garantált.
Az alábbi példa egy 10 perces vízjel küszöbértéket alkalmaz egy időablakkal számolt összesítésre.
Python
from pyspark.sql.functions import window
(df
.withWatermark("event_time", "10 minutes")
.groupBy(
window("event_time", "5 minutes"),
"id")
.count()
)
Scala
import org.apache.spark.sql.functions.window
df
.withWatermark("event_time", "10 minutes")
.groupBy(
window($"event_time", "5 minutes"),
$"id")
.count()
Ebben a példában:
- A
event_timeoszlop egy 10 perces vízjel és egy 5 perces guruló ablak definiálására szolgál. - A rendszer minden megfigyelt
idesetén darabszámot gyűjt minden nem átfedő 5 perces időablakban. - Az állapotinformáció az egyes darabszámokhoz mindaddig megőrződik, amíg az ablak vége 10 perccel korábbra nem esik, mint a legutóbb megfigyelt
event_time.
Fontos
Egy groupBy() és window() műveletben az oszlopokra a nevük alapján hivatkozzon, "<colName>" vagy col("<colName>") használatával, hogy az eseményidő-jelölő megmaradjon. A Scalában is használhatja $colName.
Hogyan befolyásolják a vízjelek a feldolgozási időt és az átviteli sebességet?
A kimeneti módok azt szabályozzák, hogy egy vízjeleket használó lekérdezés mikor ír adatokat a célhelyre. A vízjelek nélkülözhetetlenek az állapotalapú streamelés átviteli sebességének szabályozásához, mivel csökkentik az állapotinformációk teljes mennyiségét a memóriában. Nem minden kimeneti mód támogatott minden állapotalapú művelethez. Lásd: Vízjelek és kimeneti mód ablakos aggregációkhoz.
A vízjel időtartamának kiválasztása kompromisszumokkal jár:
- A rövidebb vízjelek csökkentik a lekérdezések késését, mivel a lekérdezések kevesebb állapotinformációt tárolnak, és minden egyes vízjel-időtartam befejeződése után megírják az eredményeket. A rövid vízjelek azonban kevéssé tolerálják a késve érkező adatokat.
- A hosszabb vízjelek nagy toleranciával rendelkeznek a késői adatokkal szemben. A hosszú vízjelek azonban növelik a lekérdezések késését, mivel a lekérdezéseknek több állapotinformációt kell tárolniuk, és hosszabb vízjel-időtartam után várniuk kell az eredmények írására.
Vízjelek és kimeneti mód ablakos aggregációkhoz
Az alábbi táblázat az időbélyegen és vízjelen lévő összesítéssel rendelkező lekérdezések feldolgozási viselkedését mutatja be:
| Kimeneti mód | Működés |
|---|---|
| Hozzáfűzés | A lekérdezés a vízjel küszöbértékének lejárta után sorokat ír a céltáblába. A késési küszöbérték alapján minden írás késik. A régi összesítési állapot a küszöbérték leteltét követően el lesz ejtve. |
| Frissít | A lekérdezés az eredmények kiszámításakor sorokat ír a céltáblába, és a lekérdezés új adatok érkezésekor frissítheti és felülírhatja a sorokat. A régi összesítési állapot a küszöbérték leteltét követően el lesz ejtve. |
| Kész | Az összesítési állapot nincs elhagyva. A lekérdezés minden eseményindítóhoz újraírja a céltáblát. |
Vízjelek és kimeneti módok stream-stream illesztésekhez
A több stream közötti illesztések csak a hozzáfűzési módot támogatják. A lekérdezések minden köteghez egyező rekordokat írnak.
Belső illesztések esetén a Databricks azt javasolja, hogy minden streamelési adatforráson állítson be vízjel küszöbértéket, hogy a lekérdezés elvethesse a régi rekordok állapotadatait. Vízjelek nélkül a Structured Streaming minden egyes aktiváláskor megpróbálja összekapcsolni a join mindkét oldaláról származó összes kulcsot, ami befolyásolhatja a teljesítményt.
Külső illesztések esetén a vízjelezés kötelező. Ha egy rekord nem egyezik, a lekérdezés null értéket ír a kulcshoz. Mivel az illesztések csak a hozzáfűzési módot támogatják, a nem egyező rekordok addig nem íródnak ki, amíg le nem jár a késési küszöbérték.
A késve érkező adatok küszöbértékének szabályozása több vízjelre vonatkozó szabállyal
Több strukturált streambemenet esetén több vízjelet is beállíthat a késői adatok tűréshatárainak szabályozásához. A vízjelekkel szabályozhatja az állapotinformációkat és a késést.
A streamelési lekérdezések több bemeneti adatfolyamot is tartalmazhatnak, amelyek egyesítve vagy egyesítve vannak. Állapotalapú műveletek esetén a bemeneti adatfolyamok mindegyike eltérő küszöbértéket igényelhet a késői adattűréshez. Adja meg ezeket a küszöbértékeket az egyes bemeneti adatfolyamok használatával withWatermark("eventTime", delay) . Az alábbiakban látható egy példa lekérdezés stream-stream összekapcsolásokkal.
Python
input_stream1 = ... # delays up to 1 hour
input_stream2 = ... # delays up to 2 hours
(input_stream1.withWatermark("eventTime1", "1 hour")
.join(
input_stream2.withWatermark("eventTime2", "2 hours"),
joinCondition)
)
Scala
val inputStream1 = ... // delays up to 1 hour
val inputStream2 = ... // delays up to 2 hours
inputStream1.withWatermark("eventTime1", "1 hour")
.join(
inputStream2.withWatermark("eventTime2", "2 hours"),
joinCondition)
A lekérdezés állapotalapú műveletekkel történő futtatásakor a strukturált stream egyenként nyomon követi az egyes bemeneti streamek maximális eseményidejének időtartamát, kiszámítja a vízjeleket a megfelelő késleltetés alapján, és egyetlen globális vízjelet határoz meg. Alapértelmezés szerint a strukturált streamelés a minimumot használja globális vízjelként. Ha egy stream lemarad a többitől, a minimális globális vízjel megakadályozza, hogy a lekérdezés véletlenül későn jelöljön meg adatokat. Ez például akkor fordulhat elő, ha az egyik stream leállítja az adatok fogadását a felsőbb rétegbeli hibák miatt. A globális vízjel biztonságosan mozog a leglassabb stream ütemében, és szükség esetén késlelteti a lekérdezés kimenetét.
A késleltetés csökkentése érdekében állítsa a(z) spark.sql.streaming.multipleWatermarkPolicy értékét max értékre (az alapértelmezett érték min), hogy a leggyorsabb adatfolyam vízjelét használja globális vízjelként. Ez a konfiguráció azonban elveti az adatokat a leglassabb streamekből. A Databricks azt javasolja, hogy körültekintően alkalmazza ezt a konfigurációt.
Vízjelek alkalmazása különböző műveletekre
A distinct művelet az állapotban minden egyedi rekordot nyilvántartja. Vízjel nélkül az állapot határozatlan ideig nő, és memóriaproblémákat okozhat. Adjon meg vízjelet egy időbélyegmezőn a kötött állapothoz, és távolítsa el a régi rekordokat a küszöbérték túllépése után.
Az alábbi példa vízjelet alkalmaz egy distinct műveletre:
Python
streamingDf = spark.readStream. ... # columns: eventTime, id, value, ...
# Apply watermark before distinct operation
(streamingDf
.withWatermark("eventTime", "1 hour")
.distinct()
)
Scala
val streamingDf = spark.readStream. ... // columns: eventTime, id, value, ...
// Apply watermark before distinct operation
streamingDf
.withWatermark("eventTime", "1 hour")
.distinct()
Ebben a példában a streamelési lekérdezés eltávolítja azokat az ismétlődő rekordokat, amelyek a legutóbbi megfigyelést eventTimekövető 1 órán belül érkeznek. A lekérdezés a küszöbérték lejárta után elveti a deduplikáció állapotadatait.
Fontos
Ha az összes oszlop helyett egy adott oszlopot szeretne deduplikálni, használja dropDuplicates() vagy dropDuplicatesWithinWatermark() helyette distinct. Lásd: Duplikátumok eltávolítása a vízjel határain belül.
Duplikált elemek elvetése a vízjelen belül
A Databricks Runtime 13.3 LTS vagy újabb verziójában egyedi azonosítóval deduplikálhatja a rekordokat egy vízjel küszöbértékén belül.
A strukturált streamelés pontosan egyszeri feldolgozást biztosít, de nem deduplikálja a rekordokat az adatforrásokból. Bármely dropDuplicatesWithinWatermark mező ismétlődéseit eltávolíthatja, még akkor is, ha a mezők eltérnek az ismétlődő rekordoktól, például az esemény időpontjától vagy az érkezési időtől.
A(z) dropDuplicatesWithinWatermark használatával a lekérdezések mindig eltávolítják a vízjelküszöbön belül érkező duplikált rekordokat. A lekérdezések a küszöbértéken kívülre érkező rekordokat is deduplikálhatják, de ez nem garantált. Annak érdekében, hogy a lekérdezések minden ismétlődést elvetjenek, állítsa a vízjel küszöbértékét nagyobbra, mint az ismétlődő események közötti maximális időbélyeg-különbség.
A metódus használatához meg kell adnia egy vízjelet dropDuplicatesWithinWatermark :
Python
streamingDf = spark.readStream. ...
# deduplicate using guid column with watermark based on eventTime column
(streamingDf
.withWatermark("eventTime", "10 hours")
.dropDuplicatesWithinWatermark(["guid"])
)
Scala
val streamingDf = spark.readStream. ... // columns: guid, eventTime, ...
// deduplicate using guid column with watermark based on eventTime column
streamingDf
.withWatermark("eventTime", "10 hours")
.dropDuplicatesWithinWatermark(Seq("guid"))
Példák felhasználási esetekre
Az alábbi példák speciális ablakozási használati eseteket mutatnak be:
Az óránkénti értékesítési összesítések kiszámításához használjon egymást nem átfedő ablakokat
A gördülő ablakok rögzített méretűek, és nem átfedő intervallumokkal rendelkeznek. Minden bemeneti sor pontosan egy ablakhoz tartozik. Használjon bukó ablakokat diszkrét időszakos aggregációk, például óránkénti értékesítési végösszegek kiszámításához:
Python
from pyspark.sql.functions import window, sum
hourly_sales = (orders
.withWatermark("timestamp", "1 hour")
.groupBy(window("timestamp", "1 hour"))
.agg(sum("amount").alias("total_sales"))
)
Scala
import org.apache.spark.sql.functions.{window, sum}
val hourlySales = orders
.withWatermark("timestamp", "1 hour")
.groupBy(window($"timestamp", "1 hour"))
.agg(sum($"amount").alias("total_sales"))
Ebben a példában:
-
window("timestamp", "1 hour")a rendeléseket nem átfedésben lévő 1 órás intervallumokba csoportosítja, például 5–6 és 6–7 óra között. -
withWatermark("timestamp", "1 hour")az egyes ablakok aggregátumát az állapotban tárolja, amíg az ablak végének időbélyege legalább 1 órával korábbi, mint a maximális rendelési időbélyeg.
Tolóablakok használata a gördülő aggregátumok kiszámításához
A csúszóablakok rögzített méretűek, átfedő intervallumokkal. Egyetlen sor több ablakhoz is tartozhat. Használjon csúszóablakokat gördülő összesítések kiszámításához, például egy gördülő 6 órás időszak alatti értékesítésekhez:
Python
from pyspark.sql.functions import window, sum
rolling_sales = (orders
.withWatermark("timestamp", "1 hour")
.groupBy(window("timestamp", "6 hours", slideDuration="1 hour"))
.agg(sum("amount").alias("total_sales"))
)
Scala
import org.apache.spark.sql.functions.{window, sum}
val rollingSales = orders
.withWatermark("timestamp", "1 hour")
.groupBy(window($"timestamp", "6 hours", "1 hour"))
.agg(sum($"amount").alias("total_sales"))
Ebben a példában:
-
window("timestamp", "6 hours", slideDuration="1 hour")a rendeléseket 6 órás időközökbe csoportosítja, amelyek 1 órával haladnak előre, például 5–11 és 6–12 óra között. -
withWatermark("timestamp", "1 hour")az egyes ablakok aggregátumát az állapotában tárolja mindaddig, amíg az ablak záró időbélyege 1 órával korábbi, mint a maximális rendelési időbélyeg. -
slideDurationkisebbnek vagy egyenlőnek kell lennie awindowDuration.
Munkamenetablakok használata a felhasználói tevékenység ellenőrzéséhez
A munkamenetablakok mérete nem rögzített. Megnyílik egy ablak, amikor egy sor megérkezik, és egy új sort nem tartalmazó résidő után bezáródik. Munkamenetablakok használatával összesítheti a tevékenységkitöréseket hosszú tétlenségi időszakok között, például egy felhasználó oldalnézeteit egy 30 perces időszakon belül:
Python
from pyspark.sql.functions import session_window, sum
sessionized_page_views = (activity
.withWatermark("timestamp", "1 hour")
.groupBy("user_id", session_window("timestamp", gapDuration="30 minutes"))
.agg(sum("page_views").alias("total_page_views"))
)
Scala
import org.apache.spark.sql.functions.{session_window, sum}
val sessionizedPageViews = activity
.withWatermark("timestamp", "1 hour")
.groupBy($"user_id", session_window($"timestamp", "30 minutes"))
.agg(sum($"page_views").alias("total_page_views"))
Ebben a példában:
-
session_window("timestamp", gapDuration="30 minutes")megnyílik egy ablak, amikor megérkezik az első oldalnézet. Minden további oldalmegtekintés, amely 30 percen belül érkezik, meghosszabbítja ezt az időablakot. Ha 30 percen belül nem érkezik oldalnézet, az ablak bezárul, és a következő oldalnézet új ablakot indít el. -
withWatermark("timestamp", "1 hour")az egyes munkamenetek összesítését állapotban tartja, amíg az ablak záró időbélyege 1 órával régebbi, mint a maximális oldalnézet-időbélyeg. - A
window()éssession_window()timeColumnargumentumánakTimestampTypevagyTimestampNTZTypetípusúnak kell lennie. - Az
current_timestamp()ablakokat az eseményidő helyett a feldolgozási idő alapján definiálhatja. - Az ablak időtartamát mikroszekundumtól akár napokig is megadhatja. A havi időtartamok és a hosszabb időtartamok nem támogatottak.
- A
completekimeneti módot ablakos összesítésekkel használva az összes ablakállapot határozatlan ideig tartható. Használja a(z)appendkimeneti módot megfelelő vízjellel az állapotnövekedés korlátozására és a nagy adatkészletek esetén fellépő memóriaproblémák megelőzésére. A kimeneti mód viselkedéséről további információt az ablakos aggregációk vízjelei és kimeneti módjai című témakörben talál.