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