Aller au contenu principal

foreachBatch (DataStreamWriter)

Définit la sortie de la query en streaming à traiter à l'aide de la fonction fournie. Pris en charge uniquement en mode d'exécution micro-batch (c'est-à-dire, lorsque le trigger n'est pas continu). Dans chaque micro-batch, la fonction fournie est appelée avec les lignes de sortie en tant que DataFrame et l'identifiant de batch. L'ID du batch peut être utilisé pour dédupliquer et écrire la sortie de manière transactionnelle vers des systèmes externes.

Syntaxe​

foreachBatch(func)

parameter​

parameter

Type

Description

func

appelable

Une fonction qui prend un DataFrame et un ID de batch (int) en entrée.

parameter

Type

Description

func

appelable

Une fonction qui prend un DataFrame et un ID de batch (int) en entrée.

Renvoie​

DataStreamWriter

Notes​

En mode Spark Connect, la fonction fournie n'a pas accès aux variables définies en dehors de celle-ci.

Exemples​

Python
import time
df = spark.readStream.format("rate").load()

def func(batch_df, batch_id):
batch_df.collect()

q = df.writeStream.foreachBatch(func).start()
time.sleep(3)
q.stop()