A ForEachBatch használata az adatfolyamok tetszőleges adatgyűjtőinek írásához

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 vagy AutoCDC szemantiká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ől foreachBatchtová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_sink UDF) egy teljesen új batch_id ciklust 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

  1. 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.
  2. 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.
  3. 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.
  4. 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 dbutils a foreach_batch_sink UDF-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.persist vagy a df.cache elemet. 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, foreachBatch a felhasználó által definiált függvényeknek (UDF) szerializálhatónak kell lenniük, és nem használhatók dbutils. 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).