foreach_batch_sink
Le décorateur @dp.foreach_batch_sink() définit un récepteur ForEachBatch, qui traite un stream comme une série de micro-batchs que vous gérez en Python avec une logique personnalisée. Vous référencez le récepteur en tant que target dans un flux d'ajout pour écrire les données transformées. Pour des conseils conceptuels, des considérations et des exemples, consultez Utiliser ForEachBatch pour écrire dans des récepteurs de données arbitraires dans les pipelines.
Syntaxe
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.
parameter
parameter | Description |
|---|---|
Nom | Facultatif. Un nom unique pour identifier le récepteur dans le pipeline. default, il s'agit du nom de l'UDF, s'il n'est pas inclus. |
batch_handler | Il s'agit de la fonction définie par l'utilisateur (UDF) appelée pour chaque micro-batch. |
df | Spark DataFrame contenant des données pour le micro-batch actuel. |
batch_id | L'ID entier du micro-batch. Spark incrémente cet ID pour chaque intervalle de trigger. Un |