foreach (DataStreamWriter)

Sağlanan yazıcı kullanılarak işlenecek akış sorgusunun çıkışını ayarlar. İşleme mantığı, giriş olarak bir satır alan bir işlev olarak veya ve isteğe bağlı process(row) ve open(partition_id, epoch_id) yöntemlerine sahip close(error) bir nesne olarak belirtilebilir.

Sözdizimi

foreach(f)

Parametreler

Parametre Türü Açıklama
f çağrılabilir veya nesne Giriş olarak Satır alan bir işlev veya yöntemi ve isteğe bağlı process(row) ve open yöntemleri olan bir close nesne.

İadeler

DataStreamWriter

Notlar

Sağlanan nesne serileştirilebilir olmalıdır. Veri yazmak için herhangi bir başlatma (örneğin, bir bağlantı açma) içinde yapılmalıdır open(), oluşturma zamanında değil.

Örnekler

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

İşlev kullanarak her satırı işleyin:

def print_row(row):
    print(row)

q = df.writeStream.foreach(print_row).start()
time.sleep(3)
q.stop()

, ve open yöntemlerine sahip processclosebir nesnesi kullanarak her satırı işleyin:

class RowPrinter:
    def open(self, partition_id, epoch_id):
        print("Opened %d, %d" % (partition_id, epoch_id))
        return True

    def process(self, row):
        print(row)

    def close(self, error):
        print("Closed with error: %s" % str(error))

q = df.writeStream.foreach(RowPrinter()).start()
time.sleep(3)
q.stop()