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