設定串流查詢的輸出,讓使用者提供寫入器處理。 處理邏輯可以指定為一個以列為輸入的函數,或是一個帶有 process(row) 和 可選 open(partition_id, epoch_id) 性 和 close(error) 方法的物件。
語法
foreach(f)
參數
| 參數 | 類型 | 說明 |
|---|---|---|
f |
可呼叫或物件 | 一個以 Row 為輸入的函式,或是一個帶有 process(row) 方法與可選 open Methods close 的物件。 |
退貨
DataStreamWriter
Notes
所提供的物件必須是可序列化的。 任何用於寫入資料(例如開啟連線)的初始化都應該在內部 open()進行,而非建構時。
Examples
import time
df = spark.readStream.format("rate").load()
使用函式處理每一列:
def print_row(row):
print(row)
q = df.writeStream.foreach(print_row).start()
time.sleep(3)
q.stop()
使用一個物件處理每一列,openprocess且有 、 ,以及close以下方法:
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()