Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Akış sorgusunun yürütülmesini başlatır ve yeni veriler geldikçe sonuçları verilen tabloya sürekli olarak gönderir. StreamingQuery nesnesi döndürür.
Sözdizimi
toTable(tableName, format=None, outputMode=None, partitionBy=None, queryName=None, **options)
Parametreler
| Parametre | Türü | Açıklama |
|---|---|---|
tableName |
str | Tablonun adı. |
format |
str, isteğe bağlı | Kaydetmek için kullanılan biçim. |
outputMode |
str, isteğe bağlı | Verilerin havuza nasıl yazıldı: append, completeveya update. |
partitionBy |
str veya list, isteğe bağlı | Bölümleme sütunlarının adları. Zaten var olan v2 tabloları için yoksayıldı. |
queryName |
str, isteğe bağlı | Sorgunun benzersiz adı. |
**options |
Diğer tüm dize seçenekleri. Çoğu akış için bir checkpointLocation sağlayın. |
İadeler
StreamingQuery
Notlar
v1 tablolarında partitionBy sütunlara her zaman saygı gösterilir. v2 tabloları için yalnızca partitionBy tablo henüz mevcut değilse dikkate alır.
Örnekler
Veri akışını tabloya kaydetme:
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")