Aller au contenu principal

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

Python
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 batch_id de 0 représente le start d'un Stream ou le start d'une refresh complète. Le code foreach_batch_sink devrait gérer correctement un refresh complet pour les sources de données en aval. Pour plus d'informations, consultez Full refresh.

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 batch_id de 0 représente le start d'un Stream ou le start d'une refresh complète. Le code foreach_batch_sink devrait gérer correctement un refresh complet pour les sources de données en aval. Pour plus d'informations, consultez Full refresh.

Sur cette page