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