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.
A ForEachBatch fogadó egy adatfolyamot mikrokötegek sorozataként dolgoz fel. Minden tétel feldolgozható a Python segítségével, az Apache Spark strukturált streameléséhez hasonló egyéni logikával. A Lakeflow-folyamatok ForEachBatch-fogadójával átalakíthatja, egyesítheti vagy írhat streamelési adatokat egy vagy több olyan célra, amely natív módon nem támogatja a streamelési írásokat.
A ForEachBatch fogadó a következő funkciókat biztosítja:
- Egyéni logika minden mikroköteghez: A ForEachBatch egy rugalmas stream kimenet. Tetszőleges műveleteket alkalmazhat (például egy külső táblába való egyesítést, több célhelyre való írást vagy upserts végrehajtását) Python kóddal.
- Teljes frissítés támogatása: A folyamatok folyamatonként kezelik az ellenőrzőpontokat, így az ellenőrzőpontok automatikusan alaphelyzetbeállnak a folyamat teljes frissítésekor. A ForEachBatch kimenettel Ön felelős a lefelé irányuló adatok alaphelyzetbe állításának kezeléséért, amikor ez bekövetkezik.
- Unity Catalog-támogatás: A ForEachBatch fogadó támogatja a Unity Catalog összes funkcióját, például a Unity Catalog-kötetek vagy -táblázatok olvasását vagy írását.
- Korlátozott karbantartás: A csővezeték nem követi nyomon, hogy milyen adatokat írnak a ForEachBatch kimenetből, ezért nem tudja eltávolítani azokat az adatokat. Ön a felelős az alsóbb rétegbeli adatkezelésért.
- Eseménynapló-bejegyzések: A folyamat eseménynaplója rögzíti az egyes ForEachBatch-fogadók létrehozását és használatát. Ha a Python függvény nem szerializálható, egy figyelmeztető bejegyzés jelenik meg az eseménynaplóban további javaslatokkal.
Megjegyzés:
- A ForEachBatch fogadó olyan streamelési lekérdezésekhez készült, mint a
append_flow. Nem kötegelt folyamatokhoz vagyAutoCDCszemantikához készült. - Az ezen az oldalon ismertetett ForEachBatch-fogadóelem folyamatvonalakhoz készült. Az Apache Spark strukturált adatfolyam-kezelést szintén támogatja
foreachBatch. A strukturált streamelésrőlforeachBatchtovábbi információt a ForeachBatch használata tetszőleges adatgyűjtőkbe való íráshoz című témakörben talál.
Mikor érdemes ForEachBatch mosogatót használni?
Használjon ForEachBatch-fogadót, ha a folyamat olyan funkciókat igényel, amelyek nem érhetők el beépített fogadóformátumban, például delta: vagy kafka. A tipikus használati esetek a következők:
- Egyesítés vagy frissítés Delta Lake-táblába: Egyéni egyesítési logika futtatása minden egyes mikro köteghez (például frissített rekordok kezelése).
- Írás több vagy nem támogatott célhelyre: Az egyes kötegek kimenetét több olyan táblába vagy külső tárolórendszerbe írja, amelyek nem támogatják a streamelést (például bizonyos JDBC-fogadókba).
- Egyéni logika vagy átalakítások alkalmazása: Adatok kezelése közvetlenül Python (például speciális kódtárak vagy speciális átalakítások használatával).
A beépített sinkekkel, illetve az egyéni sinkek Python használatával történő létrehozásával kapcsolatos információkért lásd: Sinkek a Lakeflow folyamatokban.
Az @dp.foreach_batch_sink() Python API-referenciát lásd: foreach_batch_sink.
Teljes frissítés
Mivel a ForEachBatch streamelési lekérdezést használ, a folyamat nyomon követi az egyes folyamatok ellenőrzőpont-könyvtárát. Teljes frissítéskor:
- Az ellenőrzési pont könyvtára visszaállítva.
- A fogadó függvény (
foreach_batch_sinkUDF) egy teljesen újbatch_idciklust lát 0-tól kezdve. - A célrendszerben lévő adatokat a folyamat nem távolítja el automatikusan (mivel a folyamat nem tudja, hol vannak megírva az adatok). Ha tiszta lap forgatókönyvre van szüksége, manuálisan kell elvetnie vagy csonkítania azokat a külső táblákat vagy helyeket, amelyeket a ForEachBatch cél tölt fel.
A Unity Catalog funkcióinak használata
A Spark strukturált streamelés foreach_batch_sink összes meglévő Unity Catalog-képessége továbbra is elérhető marad.
Ez magában foglalja a felügyelt vagy külső Unity Catalog-táblákba való írást is. Unity Catalog által felügyelt vagy külső táblákba pontosan ugyanúgy írhatja a mikroadatcsomagokat, mint bármely Apache Spark Structured Streaming feladatban.
Eseménynapló-bejegyzések
Amikor létrehoz egy ForEachBatch fogadót, egy SinkDefinition esemény "format": "foreachBatch" hozzáadódik a folyamat eseménynaplójához.
Ez lehetővé teszi a ForEachBatch kimeneti egységek használatának nyomon követését és figyelmeztetések megtekintését a kimeneti egységekkel kapcsolatban.
A Databricks Connect használata
Ha a megadott függvény nem szerializálható (a Databricks Connect fontos követelménye), az eseménynapló tartalmaz egy WARN bejegyzést, amely javasolja a kód egyszerűsítését vagy újrabontását, ha a Databricks Connect támogatása szükséges.
Ha például a ForEachBatch UDF-ben lévő paraméterek lekérésére használja a dbutils-t, az argumentumot megkaphatja, mielőtt felhasználná az UDF-ben.
# Instead of accessing parameters within the UDF...
def foreach_batch(df, batchId):
value = dbutils.widgets.get ("X") + str (i)
# ...get the parameters first, and use them within the UDF:
argX = dbutils.widgets.get ("X")
def foreach_batch(df, batchId):
value = argX + str (i)
Ajánlott eljárások
- Tartsa tömören a ForEachBatch függvényt: Kerülje a szálkezelést, a nagy erőforrástár-függőségeket vagy a nagy memóriabeli adatmanipulációkat. Az összetett vagy állapotalapú logika szerializálási hibákhoz vagy teljesítménybeli szűk keresztmetszetekhez vezethet.
- Ellenőrizze az ellenőrzőpontmappát: Streaming lekérdezések esetén a folyamatlánc az ellenőrzőpontokat folyamonként kezeli, nem fogadónként. Ha több folyam is van a csővezetékben, mindegyik folyamnak saját ellenőrzőpont-könyvtára van.
- Külső függőségek ellenőrzése: Ha külső rendszerekre vagy kódtárakra támaszkodik, ellenőrizze, hogy telepítve vannak-e az összes fürtcsomóponton vagy a tárolóban.
-
Vegye figyelembe a Databricks Connectet: Ha a környezet a jövőben áttérhet a Databricks Connectre, ellenőrizze, hogy a kód szerializálható-e, és nem támaszkodik-e
dbutilsaforeach_batch_sinkUDF-en belül.
Korlátozások
- Nincs lehetőség a ForEachBatch kezelésére: Mivel az egyéni Python kód bárhol írhat adatokat, a folyamat nem tudja kezelni vagy nyomon követni az adatokat. Saját adatkezelési vagy adatmegőrzési szabályzatokat kell kezelnie a célhelyekhez, amelyekbe ír.
- Mikroköteg metrikái: A folyamatok streamelési metrikákat gyűjtenek, de egyes forgatókönyvek hiányos vagy szokatlan metrikákat okozhatnak a ForEachBatch használatakor. Ennek oka a ForEachBatch mögöttes rugalmassága, amely megnehezíti az adatfolyamok és sorok nyomon követését a rendszer számára.
-
Több célhelyre való írás támogatása több olvasás nélkül: Egyes ügyfelek a ForEachBatch használatával egyszer olvashatnak egy forrásból, majd több célhelyre írhatnak. Ennek eléréséhez bele kell foglalnia a ForEachBatch függvényben a
df.persistvagy adf.cacheelemet. Ezekkel a beállításokkal Azure Databricks az adatokat csak egyszer próbálja beolvasni. E beállítások nélkül a lekérdezés több olvasást eredményez. Ez nem szerepel a következő kód példákban. -
A Databricks Connect használata: Ha a folyamat a Databricks Connecten fut,
foreachBatcha felhasználó által definiált függvényeknek (UDF) szerializálhatónak kell lenniük, és nem használhatókdbutils. A csővezeték figyelmeztetést ad, ha nem szerializálható UDF-et észlel, de nem áll le a csővezeték. -
Nem szerializálható logika: A helyi objektumokra, osztályokra vagy nem használható erőforrásokra hivatkozó kód megtörhet a Databricks Connect-környezetekben. Használjon tiszta Python modulokat, és győződjön meg arról, hogy a hivatkozások (például
dbutils) nem használhatók, ha a Databricks Connect követelmény.
Példák
Példa alapszintű szintaxisra
from pyspark import pipelines as dp
# Create a ForEachBatch sink
@dp.foreach_batch_sink(name = "my_foreachbatch_sink")
def feb_sink(df, batch_id):
# Custom logic here. You can perform merges,
# write to multiple destinations, etc.
return
# Create source data for example:
@dp.table()
def example_source_data():
return spark.range(5)
# Add sink to an append flow:
@dp.append_flow(
target="my_foreachbatch_sink",
)
def my_flow():
return spark.readStream.format("delta").table("example_source_data")
Mintaadatok használata egyszerű feldolgozási folyamathoz
Ez a példa a NYC Taxi mintát használja. Feltételezi, hogy a munkaterület rendszergazdája engedélyezte a Databricks nyilvános adathalmazok katalógusát. A célrendszer esetében módosítsa a my_catalog.my_schema katalógust és sémát, amelyhez hozzáféréssel rendelkezik.
from pyspark import pipelines as dp
from pyspark.sql.functions import current_timestamp
# Create foreachBatch sink
@dp.foreach_batch_sink(name = "my_foreach_sink")
def my_foreach_sink(df, batch_id):
# Custom logic here. You can perform merges,
# write to multiple destinations, etc.
# For this example, we are adding a timestamp column.
enriched = df.withColumn("processed_timestamp", current_timestamp())
# Write to a Delta location
enriched.write \
.format("delta") \
.mode("append") \
.saveAsTable("my_catalog.my_schema.trips_sink_delta")
# Return is optional here, but generally not used for the sink
return
# Create an append flow that reads sample data,
# and sends it to the ForEachBatch sink
@dp.append_flow(
target="my_foreach_sink",
)
def taxi_source():
df = spark.readStream.table("samples.nyctaxi.trips")
return df
Írás több célhelyre
Ez a példa több célhelyre is ír. Bemutatja, hogyan használható a txnVersion és txnAppId a Delta Lake-táblák írásainak idempotenssé tételéhez. További részletekért lásd: Idempotens táblaírások használataforeachBatch.
Tegyük fel, hogy két táblába írunk, table_a és table_btegyük fel, hogy egy kötegen belül az írás table_a sikeres lesz, míg az írás table_b sikertelen lesz. Ha újrafuttatják a köteget, a (txnVersion, txnAppId) pár lehetővé teszi a Deltának, hogy figyelmen kívül hagyja a duplikált írást table_a, és csak a köteget írja be table_b.
from pyspark import pipelines as dp
app_id = "my-app-name" # different applications that write to the same table should have unique txnAppId
# Create the ForEachBatch sink
@dp.foreach_batch_sink(name="user_events_feb")
def user_events_handler(df, batch_id):
# Optionally do transformations, logging, or merging logic
# ...
# Write to a Delta table
df.write \
.format("delta") \
.mode("append") \
.option("txnVersion", batch_id) \
.option("txnAppId", app_id) \
.saveAsTable("my_catalog.my_schema.example_table_1")
# Also write to a JSON file location
df.write \
.format("json") \
.mode("append") \
.option("txnVersion", batch_id) \
.option("txnAppId", app_id) \
.save("/tmp/json_target")
return
# Create source data for example
@dp.table()
def example_source():
return spark.range(5)
# Create the append flow, and target the ForEachBatch sink
@dp.append_flow(target="user_events_feb", name="user_events_flow")
def read_user_events():
return spark.readStream.format("delta").table("example_source")
Az spark.sql() használata
A ForEachBatch-fogadóban is használható spark.sql() , ahogy az alábbi példában is látható.
from pyspark import pipelines as dp
from pyspark.sql import Row
@dp.foreach_batch_sink(name = "example_sink")
def feb_sink(df, batch_id):
df.createOrReplaceTempView("df_view")
df.sparkSession.sql("MERGE INTO target_table AS tgt " +
"USING df_view AS src ON tgt.id = src.id " +
"WHEN MATCHED THEN UPDATE SET tgt.id = src.id * 10 " +
"WHEN NOT MATCHED THEN INSERT (id) VALUES (id)"
)
return
# Create target delta table
spark.range(5).write.format("delta").mode("overwrite").saveAsTable("target_table")
# Create source table
@dp.table()
def src_table():
return spark.range(5)
@dp.append_flow(
target="example_sink",
)
def example_flow():
return spark.readStream.format("delta").table("source_table")
Összevonás külső Delta Lake-táblával
from pyspark import pipelines as dp
from pyspark.sql.functions import col
from delta.tables import DeltaTable
@dp.foreach_batch_sink(name = "external_merge_feb")
def foreachBatchFunc(df, batchId):
out = DeltaTable.forName(df.sparkSession, $table)
out.alias("target") \
.merge(df.alias("source"), "source.value = target.value") \
.whenMatchedUpdateAll() \
.whenNotMatchedInsertAll() \
.whenNotMatchedBySourceDelete() \
.execute()
@dp.update_flow(
target="external_merge_feb",
name="merge_flow"
)
def read_data():
return (
spark.readStream.format("delta")
.load("/tmp/source_delta_table")
.filter(col("value").isNotNull())
)
Gyakran ismételt kérdések (FAQ)
Használhatom dbutils a ForEachBatch mosogatómat?
Ha nem Databricks Connect-környezetben tervezi futtatni a folyamatot, dbutils akkor működni fog. A Databricks Connect dbutils használata esetén azonban nem érhető el a függvényen foreachBatch belül. A folyamat figyelmeztetéseket adhat ki, ha dbutils használatot észlel, hogy segítsen elkerülni a töréseket.
Használhatok több adatfolyamot egyetlen ForEachBatch-célállomással?
Igen. Több olyan folyamatot is meghatározhat @dp.append_flow amelyek mindegyike ugyanazt a fogadónevet célozza, de mindegyik saját ellenőrzőpontokat tart fenn.
Kezeli a folyamat a cél adatainak megőrzését vagy törlését?
Nem. Mivel a ForEachBatch-fogadó bármilyen tetszőleges helyre vagy rendszerre írhat, a folyamat nem tudja automatikusan kezelni vagy törölni az adatokat a célhelyen. Ezeket a műveleteket az egyéni kód vagy külső folyamatok részeként kell kezelnie.
Hogyan háríthatom el a ForEachBatch függvény szerializálási hibáit vagy kudarcait?
Nézze meg a klasztervezérlő naplókat vagy a folyamatvonal eseménynaplóit. A Spark Connecttel kapcsolatos szerializációs problémák esetén ellenőrizze, hogy a függvény csak szerializálható Python objektumoktól függ, és nem hivatkozik-e letiltott objektumokra (például nyitott fájlfogópontokra vagy dbutils).