Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Интерфейс, используемый для записи потокового кадра данных во внешние системы хранения (например, файловых систем и хранилищ 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()