Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Запускает выполнение потокового запроса, постоянно выводя результаты в указанную таблицу по мере поступления новых данных. Возвращает объект StreamingQuery.
Синтаксис
toTable(tableName, format=None, outputMode=None, partitionBy=None, queryName=None, **options)
Параметры
| Параметр | Тип | Описание |
|---|---|---|
tableName |
str | Название таблицы. |
format |
str, необязательный | Формат, используемый для сохранения. |
outputMode |
str, необязательный | Запись данных в приемник: appendили completeupdate. |
partitionBy |
str или list, необязательный | Имена столбцов секционирования. Игнорируется для таблиц версии 2, которые уже существуют. |
queryName |
str, необязательный | Уникальное имя запроса. |
**options |
Все остальные параметры строки.
checkpointLocation Укажите для большинства потоков. |
Возвраты
StreamingQuery
Примечания
Для таблиц partitionBy версии 1 столбцы всегда учитываются. Для таблиц версии 2 учитывается только в том случае, partitionBy если таблица еще не существует.
Примеры
Сохраните поток данных в таблицу:
import tempfile
import time
_ = spark.sql("DROP TABLE IF EXISTS my_table2")
with tempfile.TemporaryDirectory(prefix="toTable") as d:
q = spark.readStream.format("rate").option(
"rowsPerSecond", 10).load().writeStream.toTable(
"my_table2",
queryName='that_query',
outputMode="append",
format='parquet',
checkpointLocation=d)
time.sleep(3)
q.stop()
spark.read.table("my_table2").show()
_ = spark.sql("DROP TABLE my_table2")