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