Aller au contenu principal

toTable (DataStreamWriter)

start l'exécution de la query de streaming, produisant continuellement des résultats dans la table donnée à mesure que de nouvelles données arrivent. Renvoie un objet StreamingQuery.

Syntaxe​

toTable(tableName, format=None, outputMode=None, partitionBy=None, queryName=None, **options)

parameter​

parameter

Type

Description

tableName

str

Nom de la table.

format

str, facultatif

Le format utilisé pour enregistrer.

outputMode

str, facultatif

Comment les données sont écrites dans le récepteur : append, complete, ou update.

partitionBy

chaîne ou liste, facultatif

Noms des colonnes de partitionnement. Ignoré pour les tables v2 qui existent déjà.

queryName

str, facultatif

Nom unique pour la query.

**options

-

Toutes les autres options de chaîne. Fournissez un checkpointLocation pour la plupart des Stream.

parameter

Type

Description

tableName

str

Nom de la table.

format

str, facultatif

Le format utilisé pour enregistrer.

outputMode

str, facultatif

Comment les données sont écrites dans le récepteur : append, complete, ou update.

partitionBy

chaîne ou liste, facultatif

Noms des colonnes de partitionnement. Ignoré pour les tables v2 qui existent déjà.

queryName

str, facultatif

Nom unique pour la query.

**options

-

Toutes les autres options de chaîne. Fournissez un checkpointLocation pour la plupart des Stream.

Renvoie​

StreamingQuery

Notes​

Pour les tables v1, les colonnes partitionBy sont toujours respectées. Pour les tables v2, partitionBy n’est respecté que si la table n’existe pas encore.

Exemples​

Enregistrer un Stream de données dans une table :

Python
import tempfile
import time
_ = spark.sql("DROP TABLE IF EXISTS my_table2")
with tempfile.TemporaryDirectory(prefix="toTable") as d:
q = spark.readStream.format("rate").option(
"rowsPerSecond", 10).load().writeStream.toTable(
"my_table2",
queryName='that_query',
outputMode="append",
format='parquet',
checkpointLocation=d)
time.sleep(3)
q.stop()
spark.read.table("my_table2").show()
_ = spark.sql("DROP TABLE my_table2")