foreachBatch (DataStreamWriter)

Sağlanan işlev kullanılarak işlenecek akış sorgusunun çıkışını ayarlar. Yalnızca mikro toplu yürütme modunda (tetikleyici sürekli olmadığında) desteklenir. Her mikro toplu işlemde, sağlanan işlev çıktı satırları DataFrame ve toplu iş tanımlayıcısı olarak çağrılır. Toplu iş kimliği, çıktıyı yinelenenleri kaldırıp dış sistemlere işlem yoluyla yazmak için kullanılabilir.

Sözdizimi

foreachBatch(func)

Parametreler

Parametre Türü Açıklama
func Callable Giriş olarak DataFrame ve toplu iş kimliği (int) alan bir işlev.

İadeler

DataStreamWriter

Notlar

Spark Connect modunda, sağlanan işlevin dışında tanımlanan değişkenlere erişimi yoktur.

Örnekler

import time
df = spark.readStream.format("rate").load()

def func(batch_df, batch_id):
    batch_df.collect()

q = df.writeStream.foreachBatch(func).start()
time.sleep(3)
q.stop()