Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Эта функция доступна в бета-версии.
Используйте функцию create_table() в конвейере для создания управляемой таблицы, написанной одним или несколькими объявлениями append_flow . Связывание create_table() вызова с одним или несколькими @append_flow(target=...) декораторами, которые записываются в таблицу. Несколько потоков могут использовать одну и ту же управляемую таблицу.
Сведения об эквиваленте SQL см. в разделеCREATE TABLE ... FLOW.
Syntax
from pyspark import pipelines as dp
dp.create_table(
name = "<table-name>",
comment = "<comment>",
spark_conf={"<key>" : "<value>", "<key>" : "<value>"},
table_properties={"<key>" : "<value>", "<key>" : "<value>"},
partition_cols=["<partition-column>", "<partition-column>"],
path="<storage-location-path>",
schema="schema-definition",
expect_all = {"<key>" : "<value>", "<key>" : "<value>"},
expect_all_or_drop = {"<key>" : "<value>", "<key>" : "<value>"},
expect_all_or_fail = {"<key>" : "<value>", "<key>" : "<value>"},
cluster_by = ["<clustering-column>", "<clustering-column>"],
cluster_by_auto = False,
row_filter = "row-filter-clause",
private = False
)
Parameters
| Parameter | Тип | Описание |
|---|---|---|
name |
str |
Required. Имя таблицы. |
comment |
str |
Описание таблицы. |
spark_conf |
dict |
Список конфигураций Spark для выполнения этого запроса. |
table_properties |
dict |
Набор dictсвойств таблицы для таблицы. |
partition_cols |
list |
Список одного или нескольких столбцов, используемых для секционирования таблицы. |
path |
str |
Расположение хранилища для данных таблицы. Если не задано, используйте управляемое расположение хранилища для схемы, содержащей таблицу. |
schema |
str или StructType |
Определение схемы для таблицы. Схемы можно определить в виде строки DDL SQL или с использованием Python StructType. |
expect_all, , expect_all_or_dropexpect_all_or_fail |
dict |
Ограничения качества данных для таблицы. Обеспечивает то же поведение и использует тот же синтаксис, что и функции декоратора ожиданий, но реализованы в качестве параметра. См. ожидания. |
cluster_by |
list |
Включите кластеризацию жидкости в таблице и определите столбцы, используемые в качестве ключей кластеризации. См. раздел "Использование кластеризации жидкости" для таблиц. |
cluster_by_auto |
bool |
Включите автоматическое кластеризация жидкости в таблице. Можно объединить с cluster_by определением начальных ключей кластеризации. См. автоматическая кластеризация жидкости. |
row_filter |
str |
(общественная предварительная версия) Условие фильтра строк для таблицы. См. публикуйте таблицы с фильтрами строк и масками столбцов. |
private |
bool |
При Trueсоздании частной таблицы, которая не публикуется в каталоге и доступна только в конвейере. По умолчанию — False. |
Ограничения
- Управляемые таблицы не поддерживают потоки изменений для отслеживания данных (CDC).
create_auto_cdc_flow()илиcreate_auto_cdc_from_snapshot_flow()назначение управляемой таблицы завершается ошибкой. Используйте create_streaming_table() для целевых объектов CDC. - Управляемые таблицы поддерживают только
append_flow. Заменить потоки (replace_flow/FLOW ... REPLACE WHERE) не поддерживаются. - Управляемые таблицы поддерживаются только в конвейерах с каталогом Unity.
- Нельзя повторно использовать имя существующей потоковой таблицы для управляемой таблицы.
Example
from pyspark import pipelines as dp
dp.create_table("combined")
@dp.append_flow(target="combined")
def from_a():
return spark.readStream.table("source_a")
@dp.append_flow(target="combined")
def from_b():
return spark.readStream.table("source_b")