foreach_batch_sink

A @dp.foreach_batch_sink() dekoratőr definiál egy ForEachBatch-fogadót, amely a streameket mikro kötegek sorozataként dolgozza fel, amelyeket az egyéni logikával Python kezel. A fogadóra a target részeként hivatkozhat az átalakított adatok megírásához. Elméleti útmutatásért, megfontolandó szempontokért és példákért tekintse meg a ForEachBatch használatát a folyamatok tetszőleges adatgyűjtőinek írásához.

Syntax

from pyspark import pipelines as dp

@dp.foreach_batch_sink(name="<name>")
def batch_handler(df, batch_id):
    """
    Required:
      - `df`: a Spark DataFrame representing the rows of this micro-batch.
      - `batch_id`: unique integer ID for each micro-batch in the query.
    """
    # Your custom write or transformation logic here
    # Example:
    # df.write.format("some-target-system").save("...")
    #
    # To access the sparkSession inside the batch handler, use df.sparkSession.

Paraméterek

Paraméter Description
név Opcionális. Egyedi név a csatornán belüli nyelő azonosításához. Alapértelmezés szerint a UDF neve, ha nincs megadva.
batch_handler Ez a felhasználó által definiált függvény (UDF), amelyet minden egyes mikroköteghez meghívunk.
Df A Spark DataFrame, amely az aktuális mikro-batch adatait tartalmazza.
batch_id A mikroköteg azonosító egész száma. A Spark ezt az azonosítót minden eseményindító-időközhöz növeli.
A batch_id vagy 0 a stream vagy a teljes frissítés kezdetét jelzi. A foreach_batch_sink kódnak megfelelően kell kezelnie az alsóbb rétegbeli adatforrások teljes frissítését. További információ: Teljes frissítés.