DataStreamWriter

Интерфейс, используемый для записи потокового кадра данных во внешние системы хранения (например, файловых систем и хранилищ key-value). Используйте df.writeStream для доступа к этому.

Синтаксис

# Access through DataFrame
df.writeStream

Методы

Метод Описание
outputMode(outputMode) Указывает, как данные потокового кадра данных записываются в приемник. Параметры: append, completeи update.
format(source) Указывает формат источника выходных данных.
option(key, value) Добавляет параметр вывода для базового источника данных.
options(**options) Добавляет несколько вариантов вывода для базового источника данных.
partitionBy(*cols) Секционирует выходные данные по заданным столбцам в файловой системе.
clusterBy(*cols) Кластеризация выходных данных по заданным столбцам.
queryName(queryName) Указывает имя потокового запроса.
trigger(**kwargs) Задает триггер для выполнения потокового запроса.
foreach(f) Задает выходные данные потокового запроса, обрабатываемого заданной функцией или объектом.
foreachBatch(func) Задает выходные данные каждого микробатча, обрабатываемого данной функцией.
start(path) Запускает выполнение потокового запроса и возвращает StreamingQuery объект.
table(tableName) Псевдоним для toTable(). Записывает данные в указанную StreamingQuery таблицу и возвращает объект.
toTable(tableName) Запускает выполнение потокового запроса, постоянно выводя результаты в указанную таблицу.

Примеры

Загрузка потока скорости, применение преобразования, запись в консоль и остановка через 3 секунды.

import time
df = spark.readStream.format("rate").load()
df = df.selectExpr("value % 3 as v")
q = df.writeStream.format("console").start()
time.sleep(3)
q.stop()