Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Importante
Esta característica se encuentra en su versión beta.
Use la create_table() función de una canalización para crear una tabla administrada escrita por una o varias declaraciones append_flow . Empareja la create_table() llamada con uno o varios @append_flow(target=...) decoradores que escriben en la tabla. Varios flujos pueden tener como destino la misma tabla administrada.
Para obtener el equivalente de SQL, consulte 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
| Parámetro | Tipo | Description |
|---|---|---|
name |
str |
Required. El nombre de la tabla. |
comment |
str |
Una descripción de la tabla. |
spark_conf |
dict |
Lista de configuraciones de Spark para la ejecución de esta consulta. |
table_properties |
dict |
Una dict de las propiedades de la tabla para la tabla. |
partition_cols |
list |
Lista de una o varias columnas que se van a usar para crear particiones en la tabla. |
path |
str |
Una ubicación de almacenamiento para los datos de tabla. Si no se establece, use la ubicación de almacenamiento administrada para el esquema que contiene la tabla. |
schema |
str o StructType |
Definición de esquema para la tabla. Los esquemas se pueden definir como una cadena de DDL de SQL o con un StructType de Python. |
expect_all, , expect_all_or_drop, expect_all_or_fail |
dict |
Restricciones de calidad de datos para la tabla. Proporciona el mismo comportamiento y usa la misma sintaxis que las funciones de decorador de expectativas, pero se implementa como un parámetro. Consulte Expectativas. |
cluster_by |
list |
Habilite la agrupación en clústeres líquidos en la tabla y defina las columnas que se usarán como claves de agrupación en clústeres. Consulte Uso de clústeres líquidos para tablas. |
cluster_by_auto |
bool |
Habilite la agrupación automática de líquidos en la tabla. Se puede combinar con cluster_by para definir las claves de agrupación en clústeres iniciales. Consulte Agrupación automática de líquidos. |
row_filter |
str |
(Versión preliminar pública) Una cláusula de filtro de fila para la tabla. Vea Publicación de tablas con filtros de fila y máscaras de columna. |
private |
bool |
Cuando True, crea una tabla privada que no se publica en el catálogo y solo es accesible dentro de la canalización. Tiene como valor predeterminado False. |
Limitaciones
- Las tablas administradas no admiten flujos de cambio de captura de datos modificados (CDC).
create_auto_cdc_flow()ocreate_auto_cdc_from_snapshot_flow()se produce un error en el destino de una tabla administrada. Use create_streaming_table() para destinos CDC. - Las tablas administradas solo
append_flowadmiten . No se admiten flujos de reemplazo (replace_flow/FLOW ... REPLACE WHERE). - Las tablas administradas solo se admiten en canalizaciones con El catálogo de Unity.
- No se puede reutilizar el nombre de una tabla de streaming existente para una tabla administrada.
Ejemplo
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")