Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Задает выходные данные потокового запроса, обрабатываемого с помощью предоставленного модуля записи. Логику обработки можно указать как функцию, которая принимает строку в качестве входных данных или как объект с process(row) необязательными open(partition_id, epoch_id) и close(error) методами.
Синтаксис
foreach(f)
Параметры
| Параметр | Тип | Описание |
|---|---|---|
f |
вызываемый или объект | Функция, которая принимает строку в качестве входных данных или объект с методом process(row) и необязательными open методами close . |
Возвраты
DataStreamWriter
Примечания
Предоставленный объект должен быть сериализуемым. Любая инициализация для записи данных (например, открытие подключения) должна выполняться внутри open(), а не во время строительства.
Примеры
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()
Обработайте каждую строку с помощью объекта и openprocessclose методов:
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()