append_flow

Декоратор @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 меняет поток двумя способами:
  • Возвращаемое значение. streaming-query. в данном случае это должен быть пакетный DataFrame, а не потоковый DataFrame.
  • Поток выполняется один раз по умолчанию. Если конвейер обновляется с полным обновлением, ONCE поток снова запускается для повторного создания данных.
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.