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()