DataStreamWriter

Dış depolama sistemlerine (örneğin, dosya sistemleri ve anahtar-değer depoları) akış DataFrame yazmak için kullanılan arabirim. Buna erişmek için kullanın df.writeStream .

Sözdizimi

# Access through DataFrame
df.writeStream

Methods

Yöntem Açıklama
outputMode(outputMode) Akış DataFrame'in verilerinin havuza nasıl yazılması olduğunu belirtir. Seçenekler : append, completeve update.
format(source) Çıkış veri kaynağı biçimini belirtir.
option(key, value) Temel alınan veri kaynağı için bir çıkış seçeneği ekler.
options(**options) Temel alınan veri kaynağı için birden çok çıkış seçeneği ekler.
partitionBy(*cols) Çıktıyı dosya sistemindeki belirtilen sütunlara göre bölümler.
clusterBy(*cols) Çıkışı verilen sütunlara göre kümeler.
queryName(queryName) Akış sorgusunun adını belirtir.
trigger(**kwargs) Akış sorgusu yürütme tetikleyicisini ayarlar.
foreach(f) Verilen işlev veya nesne tarafından işlenecek akış sorgusunun çıkışını ayarlar.
foreachBatch(func) Verilen işlev tarafından işlenecek her mikrobatch'in çıkışını ayarlar.
start(path) Akış sorgusunun yürütülmesini başlatır ve bir StreamingQuery nesne döndürür.
table(tableName) Diğer ad için toTable(). Belirtilen tabloya veri yazar ve bir StreamingQuery nesne döndürür.
toTable(tableName) Akış sorgusunun yürütülmesini başlatır ve verilen tabloya sürekli olarak sonuç çıkartır.

Örnekler

Hız akışı yükleyin, dönüştürme uygulayın, konsola yazın ve 3 saniye sonra durdurun.

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