Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Декоратор @dp.append_flow создает потоки добавления данных или обратное наполнение для таблиц конвейера. Функция должна возвращать кадр данных потоковой передачи Apache Spark. См. сведения о загрузке и обработке данных с помощью потоков конвейера Lakeflow.
Потоки добавления могут быть направлены на таблицы потоков или приемники.
Синтаксис
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.append_flow(
target = "<target-table-name>",
name = "<flow-name>", # optional, defaults to function name
once = False, # optional
spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
comment = "<comment>", # optional
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>) #
Параметры
| Параметр | Тип | Description |
|---|---|---|
| function | function |
Обязательное. Функция, возвращающая кадр данных потоковой передачи Apache Spark из определяемого пользователем запроса. |
target |
str |
Обязательное. Имя таблицы или приемника, которое является целевым объектом потока добавления. |
name |
str |
Имя потока. Если этот параметр не указан, по умолчанию используется имя функции. |
once |
bool |
При необходимости определите поток как одноразовый, например, обратная закачка. Использование once=True меняет поток двумя способами:
|
comment |
str |
Описание потока. |
spark_conf |
dict |
Список конфигураций Spark для выполнения этого запроса |
import_checkpoint |
str |
Путь к существующей контрольной точке структурированного потока для импорта в поток, поэтому мигрировавший поток возобновляется с последнего фиксированного смещения вместо повторной обработки источника. Импорт контрольной точки находится в бета-версии. См. Миграция структурированной контрольной точки потока. |
Примеры
from pyspark import pipelines as dp
# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})
# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
return <streaming-query>
# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
return (
spark.read
.format("json")
.load("/path/to/backfill/")
)
# Create a Kafka sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "host:port",
"topic": "my_topic"
}
)
# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))
Миграция структурированной контрольной точки потока
Important
Импорт контрольной точки находится в бета-версии.
Используйте import_checkpoint миграцию существующей рабочей нагрузки структурированного потока в конвейер без повторной обработки исходника. Установите запрос checkpointLocation Structured Streaming, который может быть облачным хранилищем, том Unity Catalog или DBFS-путь. При первом обновлении конвейера поток клонирует эти контрольные точки в управляемое хранилище конвейера. Затем поток возобновляется с последнего фиксированного смещения с сохранением состояния (например, агрегации, ключи дедупликации и водяных знаков). Последующие обновления конвейера используют клонированную контрольную точку потока; исходная контрольная точка не изменяется.
Поток должен быть направлен на управляемую таблицу, созданную с помощью create_table , или раковину.
Остановите исходный запрос Structured Streaming перед запуском конвейера. Оригинальный запрос Structured Streaming можно повторно использовать после импорта, но нужно управлять его состоянием контрольной точки и следить, чтобы конвейер и запрос не записывались в одну и ту же таблицу одновременно, что может привести к дублирующимся данным.
Воссоздайте запрос Structured Streaming как поток конвейера, который записывает в новую таблицу и импортирует её контрольную точку:
from pyspark import pipelines as dp
# Create a new managed table for the pipeline
dp.create_table("target_table")
# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
target = "target_table",
import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
# The same source your original query read from.
return spark.readStream.table("source_table")
Контрольная точка импортируется только один раз — при первом обновлении конвейера; поздние обновления игнорируются import_checkpoint.
Полное обновление не импортирует контрольную точку; она начинается с новой, пустой контрольной точки и повторно обрабатывает исходный код. Для импорта другой контрольной точки используйте имя потока, которое ранее не использовалось для целевой таблицы; повторное использование существующего имени потока пропускает импорт.
Ограничения
- Импорт контрольной точки в уже существующую таблицу (например, исходную цель запроса Structured Streaming) не поддерживается. Нацеливайтесь на новую таблицу, которую создаёт конвейер, или на погрузчик.
-
import_checkpointподдерживается только на append_flow.