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 @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. |