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 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 megadottbatchId-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áljaforeach()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():
-
OPTIMIZEa feldolgozandó fájlok nélkül: Ha egyOPTIMIZEmű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.widgetsalmodul 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.