foreachBatch (DataStreamWriter)

Задает выходные данные потокового запроса, обрабатываемого с помощью предоставленной функции. Поддерживается только в режиме выполнения микро-пакетной службы (то есть, если триггер не является непрерывным). В каждом микропакете предоставленная функция вызывается с выходными строками в виде кадра данных и идентификатора пакета. Идентификатор пакета можно использовать для дедупликации и транзакционно записи выходных данных во внешние системы.

Синтаксис

foreachBatch(func)

Параметры

Параметр Тип Описание
func Вызываемые Функция, которая принимает кадр данных и идентификатор пакета (int) в качестве входных данных.

Возвраты

DataStreamWriter

Примечания

В режиме Spark Connect предоставленная функция не имеет доступа к переменным, определенным вне него.

Примеры

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