DataStreamReader

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

Синтаксис

# Access through SparkSession
spark.readStream

Методы

Метод Описание
format(source) Задает формат источника входных данных.
schema(schema) Задает схему потокового кадра данных.
option(key, value) Добавляет входной параметр для базового источника данных.
options(**options) Добавляет несколько вариантов ввода для базового источника данных.
load(path) Загружает кадр данных потоковой передачи из заданного пути и возвращает его.
json(path) Загружает поток JSON-файлов и возвращает кадр данных.
orc(path) Загружает поток файлов ORC и возвращает кадр данных.
parquet(path) Загружает поток файлов Parquet и возвращает кадр данных.
text(path) Загружает текстовый файловый поток и возвращает кадр данных.
csv(path) Загружает поток CSV-файла и возвращает кадр данных.
xml(path) Загружает поток XML-файлов и возвращает кадр данных.
table(tableName) Загружает таблицу delta потоковой передачи и возвращает кадр данных.
name(source_name) Присваивает имя источнику потоковой передачи для эволюции контрольных точек.
changes(tableName) Возвращает изменения на уровне строк (запись измененных данных) из указанной таблицы в виде потокового кадра данных.

Примеры

spark.readStream
# <...streaming.readwriter.DataStreamReader object ...>

Загрузка потока скорости, применение преобразования, запись в консоль и остановка через 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()