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