A foreachBatch használata tetszőleges adatgyűjtőkbe való íráshoz

Ez a lap bemutatja, hogyan használható foreachBatch a Strukturált streamelés szolgáltatással streamlekérdezés kimenetének írására olyan adatforrások számára, amelyek nem rendelkeznek meglévő streamelt fogadóval.

A kódminta streamingDF.writeStream.foreachBatch(...) lehetővé teszi kötegelt függvények alkalmazását a streamelési lekérdezés minden mikro kötegének kimeneti adataira. A foreachBatch használt függvények két paramétert használnak:

  • Olyan DataFrame, amely egy mikro köteg kimeneti adatait tartalmazza.
  • A(z) mikro tétel egyedi azonosítója.

A Delta Lake egyesítési műveletekhez a strukturált streamelésben a foreachBatch elemet kell használnia. Tekintse meg a foreachBatch használatát a streamelési lekérdezésekben az Upserttel kapcsolatban.

További DataFrame-műveletek alkalmazása

Számos DataFrame- és adathalmaz-művelet nem támogatott a streamelési DataFrame-ekben, mert a Spark ezekben az esetekben nem támogatja a növekményes tervek generálásának támogatását. Az foreachBatch() segítségével az egyes mikroköteg-kimenetekre alkalmazhatja ezen műveletek némelyikét. Például használhatja a foreachBatch() és az SQL MERGE INTO műveletet arra, hogy a streaming aggregációk kimenetét frissítési módban egy Delta Lake-táblába írja. További részletekért lásd: MERGE INTO.

Fontos

  • foreachBatch() csak legalább egyszer írható garanciát biztosít. Azonban a funkcióhoz megadott batchId-t használhatja a kimenet deduplikálására, így biztosítható az egyszeri végrehajtás garanciája. Mindkét esetben meg kell indokolnia a végpontok közötti szemantikát.
  • foreachBatch() nem működik a folyamatos feldolgozási móddal , mivel alapvetően egy streamelési lekérdezés mikroköteges végrehajtására támaszkodik. Ha folyamatos módban ír adatokat, használja foreach() helyette.
  • foreachBatch Állapotalapú operátor használata esetén fontos, hogy a feldolgozás befejezése előtt minden köteg teljesen fel legyen használva. Lásd: Az egyes kötegek adatkeretének teljes felhasználása

Üres adatkeretek kezelése

foreachBatch() lehet, hogy üres DataFrame-et kap, és a kódnak kezelnie kell ezt a forgatókönyvet. Ellenkező esetben a lekérdezés sikertelen lehet.

Ha például a Delta Lake a streamforrás, ezek a forgatókönyvek üres DataFrame-et adhatnak át a következőnek foreachBatch():

  • OPTIMIZE a feldolgozandó fájlok nélkül: Ha egy OPTIMIZE művelet a Delta Lake forrástábláján fut, de nincsenek feldolgozandó fájlok, a Strukturált streamelés eltolásnapló-bejegyzést ír a táblaverzió növeléséhez. Ez üres mikroköteget hoz létre a célhelyen, annak ellenére, hogy fájlok nincsenek beolvasva.
  • Fájlpruning a fizikai terv szintjén: Ha a predikátum leküldése vagy a fájlpruning kiküszöböli a fizikai terv szintjén lévő összes rekordot, az eredmény egy üres rögzítés a célállomáshoz.

A felhasználói kódnak üres adatkereteket kell kezelnie a megfelelő működéshez. Lásd az alábbi példákat:

Python

def process_batch(output_df, batch_id):
  # Process valid DataFrames only
  if not output_df.isEmpty():
    # business logic
    pass

streamingDF.writeStream.foreachBatch(process_batch).start()

Scala

  .foreachBatch(
          (outputDf: DataFrame, bid: Long) => {
             // Process valid DataFrames only
             if (!outputDf.isEmpty) {
                // business logic
             }
         }
  ).start()

Viselkedési változások a Databricks Runtime 14.0-ban foreachBatch

A Databricks Runtime 14.0-s és újabb verziókban a standard hozzáférési móddal konfigurált számításon a következő viselkedésváltozások érvényesek:

  • print() parancsok írnak kimenetet az illesztőprogram-naplókba.
  • A dbutils.widgets almodul nem érhető el a függvényen belül.
  • A függvényben hivatkozott fájloknak, moduloknak és objektumoknak szerializálhatónak kell lenniük, és elérhetőnek kell lenniük a Sparkban.

Meglévő kötegelt adatforrások újrafelhasználása

A foreachBatch() használatával meglévő kötegelt adatok íróit használhatja olyan adatgyűjtőkhöz, amelyek esetleg nem rendelkeznek strukturált streamelési támogatással. Íme néhány példa:

Számos más kötegelt adatforrás használható a foreachBatch()-ből. Lásd: Csatlakozás adatforrásokhoz és külső szolgáltatásokhoz.

Több helyre írni

Ha egy streamelési lekérdezés kimenetét több helyre kell írnia, a Databricks több strukturált streamíró használatát javasolja a legjobb párhuzamosság és átviteli sebesség érdekében.

Ha foreachBatch több fogadóba ír, szerializálja a streamelési írások végrehajtását, ami növelheti az egyes mikrokötegek késését.

Ha több Delta Lake-táblába való íráshoz használja a(z) foreachBatch elemet, lásd: A foreachBatch használata idempotens táblaírásokhoz.

Az egyes kötegek adatkeretének teljes felhasználása

Ha állapotalapú operátorokat használ (például használ dropDuplicatesWithinWatermark), minden köteg iterációnak a teljes DataFrame-et kell használnia, vagy újra kell indítania a lekérdezést. Ha nem használja fel a teljes DataFrame-et, a streamelési lekérdezés a következő köteggel meghiúsul.

Ez több esetben is előfordulhat. Az alábbi példák bemutatják, hogyan javítható ki a DataFrame-et nem megfelelően használó lekérdezések.

A köteg egy részhalmazának szándékos használata

Ha csak a köteg egy részhalmaza érdekli, a következő kóddal rendelkezhet.

from pyspark.sql.functions import expr

stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
# creates a stateful operator:
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()

def partial_func(batch_df, batch_id):
  batch_df.show(2)

q = streamWithWatermark.writeStream \
  .foreachBatch(partial_func) \
  .option("checkpointLocation", checkpoint_dir) \
  .trigger(processingTime='2 seconds') \
  .start()

Ebben az esetben csak a batch_df.show(2) köteg első két elemét kezeli, ami várható, de ha több elem is van, azokat fel kell használni. Az alábbi kód a teljes DataFrame-et használja.

from pyspark.sql.functions import expr

stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()

# function to do nothing with a row
def do_nothing(row):
  pass

def partial_func(batch_df, batch_id):
  batch_df.show(2)
  batch_df.foreach(do_nothing) # silently consume the rest of the batch

q = streamWithWatermark.writeStream \
  .foreachBatch(partial_func) \
  .option("checkpointLocation", checkpoint_dir) \
  .trigger(processingTime='2 seconds') \
  .start()

Itt a do_nothing függvény csendben figyelmen kívül hagyja a DataFrame többi részét.

A hiba kezelése batch feldolgozás során

A hibakezeléshez foreachBatcha Databricks azt javasolja, hogy engedélyezze a streamelési lekérdezés gyors meghiúsulását, és ehelyett az újrapróbálkozási logika kezeléséhez a vezénylési rétegre, például a Lakeflow-feladatokra vagy az Apache Airflow-ra támaszkodjon. Ez sokkal biztonságosabb, mint összetett újrapróbálkozások hurkoinak létrehozása a kódban, ahol adatvesztés fordulhat elő.

Az alábbiakban az írási célon alapuló irányelveket talál:

Target Examples Útmutatás
DataFrame-műveletek Delta Lake-táblák Az txnAppId és txnVersion írási opciókat kell használnia, és txnVersion-t a batchId-hoz kell kapcsolnia, hogy garantálja az idempotenciát és biztosítsa az adatok helyességét az újrapróbálkozások során. Ne kapja meg és próbálkozzon újra a kivételekkel helyileg. Ehelyett a Databricks azt javasolja, hogy engedélyezze a hibák propagálását, hogy a Spark metrikák pontosak maradjanak, az adatok ne ismétlődjenek, és az orchestrátor tiszta módon ismét kísérelje meg a teljes tételt.
Egyéni kód és külső célhelyek .collect(), OLTP-adatbázisok, üzenetsorok, API-k A saját idempotencia megvalósítása. Fel kell tételezni, hogy bármely műveletet újra lehet próbálni és próbálják is különböző kötegekben. Ha a batchId változatlan marad, a művelet eredményének is változatlannak kell maradnia. Előfordulhat, hogy csak átmeneti hibák esetén próbálkozik újra, például rövid kapcsolati időtúllépéseknél, de legyen rendkívül körültekintő, hogy elkerülje a részleges vagy ismétlődő írásokat, ha az újrapróbálkozás végül meghiúsul. A legbiztonságosabb módszer a hibák propagálása, és a vezérlő számára a teljes adatcsomag újrapróbálkozásának lehetővé tétele.

Íme néhány példa a kivételtípusokra és a javaslatokra, hogy hogyan kezelje őket a foreachBatch:

Kivétel típusa Examples Javasolt művelet
Átmeneti fogadóhibák SQLTransientConnectionException, HTTP 429, időtúllépések Elfogás: újrapróbálás, vagy küldés holtlevelek üzenetsorába
Duplikált vagy kulcsfontosságú kényszer megsértése, ha a fogadó idempotens SQLIntegrityConstraintViolationException Fogás: naplózás és letiltás
Egyéni újrapróbálható hibák Burkolt socket-kivételek, újrapróbálható adatbázishibák Fogás: metrikák növekménye és szabályozott folytatás engedélyezése
Logikai vagy sémahibák NullPointerException, AttributeError sémaeltérés Propagálás: engedélyezi, hogy a Spark hibát okozzon a lekérdezés során
Nem újrapróbálkozható fogadóhibák vagy nem észlelt logikai hibák ValueError, PermissionError Propagálás: engedélyezi, hogy a Spark hibát okozzon a lekérdezés során
Kritikus hibák OutOfMemoryError, sérült állapot, adatintegritási szabálysértések Propagálás: engedélyezi, hogy a Spark hibát okozzon a lekérdezés során

Példakódok: kivételkezelés

Az alábbi példák szándékosan okoznak hibát foreach a hiba kezelésére szolgáló különböző megközelítések megjelenítésében:

from pyspark.sql.functions import expr

stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
# creates a stateful operator:
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()

def foreach_func(row):
  # handle the row, but in this case, for the sample, will just raise an error:
  raise Exception('error')

def partial_func(batch_df, batch_id):
  try:
    batch_df.foreach(foreach_func)
  except Exception as e:
    print(e) # or whatever error handling you want to have


q = streamWithWatermark.writeStream \
  .foreachBatch(partial_func) \
  .option("checkpointLocation", checkpoint_dir) \
  .trigger(processingTime='2 seconds') \
  .start()

A fenti kód kezeli és észrevétlenül elnyomja a hibát, és lehet, hogy nem használja fel a köteg többi részét. A helyzet kezelésére két lehetőség áll rendelkezésre.

Először újraindíthatja a hibát, amely visszakerül az orchestration rétegbe a köteg újrapróbálkozásához. Ez megoldhatja a hibát, ha az egy átmeneti probléma, vagy jelezheti az operációs csapatnak, hogy próbálják meg manuálisan kijavítani. Ehhez módosítsa a partial_func kódot úgy, hogy így nézzen ki:

def partial_func(batch_df, batch_id):
  try:
    batch_df.foreach(foreach_func)
  except Exception as e:
    print(e) # or whatever error handling you want to have
    raise e # re-raise the issue

Másodszor, ha ki szeretné kapni a kivételt, és figyelmen kívül szeretné hagyni a köteg többi részét, módosíthatja a kódot úgy, hogy a do_nothing függvény használatával csendben figyelmen kívül hagyja a köteg többi részét.

from pyspark.sql.functions import expr

stream = spark.readStream.format("rate").option("rowsPerSecond", "100").load()
# creates a stateful operator:
streamWithWatermark = stream.withWatermark("timestamp", "15 minutes").dropDuplicatesWithinWatermark()

def foreach_func(row):
  # handle the row, but in this case, for the sample, will just raise an error:
  raise Exception('error')

# function to do nothing with a row
def do_nothing(row):
    pass

def partial_func(batch_df, batch_id):
  try:
    batch_df.foreach(foreach_func)
  except Exception as e:
    print(e) # or whatever error handling you want to have
    batch_df.foreach(do_nothing) # silently consume the remainder of the batch

q = streamWithWatermark.writeStream \
  .foreachBatch(partial_func) \
  .option("checkpointLocation", checkpoint_dir) \
  .trigger(processingTime='2 seconds') \
  .start()

Sikertelen rekordok írása kézbesítetlen levelek üzenetsorába

Alapértelmezés szerint a lekérdezések azonnal meghiúsulnak, amikor rossz rekordok érkeznek. Ezeket a megszakításokat elkerülheti, ha egy másodlagos Delta Lake-táblát állít be halottbetűs üzenetsorként (DLQ).

A DLQ használatával a rendszer a sikertelen rekordokat a másodlagos táblába irányítja, és zavartalanul folytatja az érvényes adatok feldolgozását. A DLQ-tábla lehetővé teszi a rossz rekordok későbbi vizsgálatát és újrafeldolgozását.

Ezt a módszert akkor használja, ha:

  • A stream változatos adatokat tartalmaz, amelyek sérthetik a sémakorlátozásokat vagy az üzleti szabályokat.
  • A naplózási vagy megfelelőségi szabályok megkövetelik az összes rekord megőrzését.

Example

Az alábbi példa a foreachBatch használatával az egyes mikrokötegeket érvényes és érvénytelen rekordokra osztja, majd mindkét részhalmazt külön Delta Lake-táblába írja.

Python

from pyspark.sql.functions import current_timestamp, lit

main_table = "catalog.schema.orders"
dlq_table = "catalog.schema.orders_dlq"
app_id = "orders-streaming-job"

def process_orders(batch_df, batch_id):
    if batch_df.isEmpty():
        return

    valid_condition = "order_amount > 0 AND customer_id IS NOT NULL"

    # Write valid records to the main table with idempotency options
    batch_df.filter(valid_condition).write \
        .format("delta") \
        .mode("append") \
        .option("txnVersion", batch_id) \
        .option("txnAppId", app_id) \
        .saveAsTable(main_table)

    # Route invalid records to the dead-letter queue
    invalid_df = batch_df.filter(f"NOT ({valid_condition})")
    if not invalid_df.isEmpty():
        invalid_df \
            .withColumn("dlq_batch_id", lit(batch_id)) \
            .withColumn("dlq_ingest_time", current_timestamp()) \
            .write \
            .format("delta") \
            .mode("append") \
            .option("txnVersion", batch_id) \
            .option("txnAppId", app_id) \
            .saveAsTable(dlq_table)

spark.readStream \
    .format("delta") \
    .table("catalog.schema.raw_orders") \
    .writeStream \
    .foreachBatch(process_orders) \
    .option("checkpointLocation", "/path/to/checkpoint") \
    .start()

Scala

import org.apache.spark.sql.DataFrame
import org.apache.spark.sql.functions.{current_timestamp, lit}

val mainTable = "catalog.schema.orders"
val dlqTable = "catalog.schema.orders_dlq"
val appId = "orders-streaming-job"

def processOrders(batchDf: DataFrame, batchId: Long): Unit = {
  if (batchDf.isEmpty) return

  val validCondition = "order_amount > 0 AND customer_id IS NOT NULL"

  // Write valid records to the main table with idempotency options
  batchDf.filter(validCondition).write
    .format("delta")
    .mode("append")
    .option("txnVersion", batchId)
    .option("txnAppId", appId)
    .saveAsTable(mainTable)

  // Route invalid records to the dead-letter queue
  val invalidDf = batchDf.filter(s"NOT ($validCondition)")
  if (!invalidDf.isEmpty) {
    invalidDf
      .withColumn("dlq_batch_id", lit(batchId))
      .withColumn("dlq_ingest_time", current_timestamp())
      .write
      .format("delta")
      .mode("append")
      .option("txnVersion", batchId)
      .option("txnAppId", appId)
      .saveAsTable(dlqTable)
  }
}

spark.readStream
  .format("delta")
  .table("catalog.schema.raw_orders")
  .writeStream
  .foreachBatch(processOrders _)
  .option("checkpointLocation", "/path/to/checkpoint")
  .start()

Mindkét írási művelet a txnVersion és txnAppId használatával biztosítja az idempotenciát. Ha a Spark ugyanazt a batchId köteget újra feldolgozza, a Delta Lake kihagyja a duplikált írási műveletet. Lásd Idempotens táblaírások használataforeachBatch.