Aller au contenu principal

DataStreamWriter

Interface permettant d'écrire un DataFrame en streaming dans des systèmes de stockage externes (par exemple, les systèmes de fichiers et les magasins de clé-valeur). Utilisez df.writeStream pour y accéder.

Syntaxe

Python
# Access through DataFrame
df.writeStream

Méthodes

Méthode

Description

outputMode(outputMode)

Spécifie comment les données d'un DataFrame de streaming sont écrites vers le récepteur. Les options sont append, complete et update.

format(source)

Spécifie le format de la source de données de sortie.

option(key, value)

Ajoute une option de sortie pour la source de données sous-jacente.

options(**options)

Ajoute plusieurs options de sortie pour la source de données sous-jacente.

partitionBy(*cols)

Partitionne la sortie par les colonnes données sur le système de fichiers.

clusterBy(*cols)

Clusters la sortie par les colonnes données.

queryName(queryName)

Spécifie le nom de la query de streaming.

trigger(**kwargs)

Définit le Trigger pour l'exécution de la query en streaming.

foreach(f)

Définit la sortie de la query en streaming à traiter par la fonction ou l'objet donné.

foreachBatch(func)

Définit la sortie de chaque microbatch à traiter par la fonction donnée.

start(path)

Démarre l’exécution de la query de streaming et renvoie un objet StreamingQuery.

table(tableName)

Alias pour toTable(). Écrit les données dans la table spécifiée et renvoie un objet StreamingQuery.

toTable(tableName)

start l'exécution de la query de streaming, produisant continuellement des résultats vers la table spécifiée.

Méthode

Description

outputMode(outputMode)

Spécifie comment les données d'un DataFrame de streaming sont écrites vers le récepteur. Les options sont append, complete et update.

format(source)

Spécifie le format de la source de données de sortie.

option(key, value)

Ajoute une option de sortie pour la source de données sous-jacente.

options(**options)

Ajoute plusieurs options de sortie pour la source de données sous-jacente.

partitionBy(*cols)

Partitionne la sortie par les colonnes données sur le système de fichiers.

clusterBy(*cols)

Clusters la sortie par les colonnes données.

queryName(queryName)

Spécifie le nom de la query de streaming.

trigger(**kwargs)

Définit le Trigger pour l'exécution de la query en streaming.

foreach(f)

Définit la sortie de la query en streaming à traiter par la fonction ou l'objet donné.

foreachBatch(func)

Définit la sortie de chaque microbatch à traiter par la fonction donnée.

start(path)

Démarre l’exécution de la query de streaming et renvoie un objet StreamingQuery.

table(tableName)

Alias pour toTable(). Écrit les données dans la table spécifiée et renvoie un objet StreamingQuery.

toTable(tableName)

start l'exécution de la query de streaming, produisant continuellement des résultats vers la table spécifiée.

Exemples

Chargez un stream de taux, appliquez une transformation, écrivez sur la console et arrêtez après 3 secondes.

Python
import time
df = spark.readStream.format("rate").load()
df = df.selectExpr("value % 3 as v")
q = df.writeStream.format("console").start()
time.sleep(3)
q.stop()