Vízjelek alkalmazása az adatfeldolgozási küszöbértékek szabályozásához

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_time oszlop egy 10 perces vízjel és egy 5 perces guruló ablak definiálására szolgál.
  • A rendszer minden megfigyelt id eseté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.
  • slideDuration kisebbnek vagy egyenlőnek kell lennie a windowDuration.

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() és session_window()timeColumn argumentumának TimestampType vagy TimestampNTZType tí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 complete kimeneti módot ablakos összesítésekkel használva az összes ablakállapot határozatlan ideig tartható. Használja a(z) append kimeneti 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.