replace_flow

Important

Эта функция доступна в бета-версии.

@dp.replace_flow Декоратор создаёт поток ЗАМЕНЫ ИСПОЛЬЗОВАНИЯ для таблицы потоков в вашем пайплайне. При каждом обновлении поток заменяет все строки в целевой таблице, соответствующие ключевым replace_using столбцам, и оставляет все остальные строки нетронутыми. Функция должна возвращать кадр данных потоковой передачи Apache Spark. См. раздел Частичная замена снимков с ЗАМЕНОЙ ИСПОЛЬЗОВАНИЕМ потоков.

Используйте @dp.replace_flow тогда, когда исходный код — серия частичных снимков, ключевых по столбцам. Чтобы определить целевой таблицу и поток в одном операторе, передайте replace_using и sequence_by в @dp.table.

Синтаксис

from pyspark import pipelines as dp

dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.

@dp.replace_flow(
  target = "<target-table-name>",
  replace_using = ["<key-column>", "<key-column>"],
  sequence_by = "<sequence-column>",
  name = "<flow-name>", # optional, defaults to function name
  comment = "<comment>", # optional
  spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}) # optional
def <function-name>():
  return (<streaming-query>)

Parameters

Parameter Type Description
function function Required. Функция, возвращающая кадр данных потоковой передачи Apache Spark из определяемого пользователем запроса.
target str Required. Название стола потока, которая является целевой целью потока.
replace_using list Required. Ключевые столбцы, определяющие, какие целевые строки заменить. Укажите хотя бы одну колонку. Ключевые столбцы нельзя повторять, и тип каждого ключевого столбца должен быть сортируемым.
sequence_by str или Column Required. Колонка, которая распоряжает обновления. Для каждого ключа выигрывает самая высокая последовательность, а низшая последовательность никогда не перезаписывает старшую, уже находящуюся в целевом ключе.
name str Имя потока. Если этот параметр не указан, по умолчанию используется имя функции.
comment str Описание потока.
spark_conf dict Список конфигураций Spark для выполнения этого запроса.

Examples

from pyspark import pipelines as dp

# Keep the latest row for each order from a stream of partial snapshots
dp.create_streaming_table("orders_current")

@dp.replace_flow(
  target = "orders_current",
  replace_using = ["order_id"],
  sequence_by = "updated_at"
)
def orders_flow():
  return spark.readStream.table("order_updates")

Используйте более одного ключевого столбца, когда запись идентифицируется комбинацией столбцов:

from pyspark import pipelines as dp

dp.create_streaming_table("accounts_current")

@dp.replace_flow(
  target = "accounts_current",
  replace_using = ["region", "account_id"],
  sequence_by = "updated_at"
)
def accounts_flow():
  return spark.readStream.table("account_updates")

Ограничения

  • Таблица потока поддерживает один REPLACE USING поток и не может объединяться REPLACE USING с другим типом потока, таким как дополнительный поток, автоматический поток CDC или поток.REPLACE WHERE
  • Запрос должен быть потоковым запросом. @dp.replace_flow отклоняет источник, не транслирующий.
  • REPLACE USING потоки требуют Databricks Runtime 18.2 и выше.