Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
O @materialized_view decorador pode ser usado para definir vistas materializadas em um gasoduto Lakeflow ou um oleoduto independente.
Para definir uma exibição materializada, aplique-se @materialized_view a uma consulta que executa uma leitura em lote em uma fonte de dados.
Você também pode aplicar @materialized_view fora de um pipeline, em um notebook em computação geral serverless, para definir uma visão materializada independente. Nesse caso, você não pode passar private, usar uma REFRESH instrução ou definir um cronograma de atualização. Volte a usar o decorador para refrescar a mesa, adicionando full_refresh=True para uma renovação completa.
Veja Definir tabelas com os decoradores de pipelines.
Sintaxe
from pyspark import pipelines as dp
@dp.materialized_view(
name="<name>",
comment="<comment>",
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"},
table_properties = {"<key>" : "<value>", "<key>" : "<value>"},
path = "<storage-location-path>",
partition_cols = ["<partition-column>", "<partition-column>"],
cluster_by_auto = False,
cluster_by = ["<clustering-column>", "<clustering-column>"],
schema = "schema-definition",
refresh_policy = None,
row_filter = "row-filter-clause",
private = False)
@dp.expect(...)
def <function-name>():
return (<query>)
Parâmetros
@dp.expect() é uma cláusula de expectativa opcional. Você pode incluir várias expectativas. Veja as expectativas.
| Parâmetro | Tipo | Description |
|---|---|---|
| função | function |
Obrigatório Uma função que retorna um DataFrame em lote do Apache Spark a partir de uma consulta definida pelo usuário. |
name |
str |
O nome da tabela. Se não for fornecido, o padrão será o nome da função. |
comment |
str |
Uma descrição da tabela. |
spark_conf |
dict |
Uma lista de configurações do Spark para a execução dessa consulta |
table_properties |
dict |
Um dict das propriedades da tabela para a tabela. |
path |
str |
Um local de armazenamento para dados de tabela. Se não estiver definido, use o local de armazenamento gerenciado para o esquema que contém a tabela. |
partition_cols |
list |
Uma lista de uma ou mais colunas a serem usadas para particionar a tabela. |
cluster_by_auto |
bool |
Habilite o agrupamento automático de líquidos na tabela. Isso pode ser combinado e cluster_by definir as colunas a serem usadas como chaves de clustering iniciais, seguidas pelo monitoramento e atualizações automáticas de seleção de chaves com base na carga de trabalho. Consulte clusterização automática de líquidos. |
cluster_by |
list |
Habilite o agrupamento líquido na tabela e defina as colunas a serem usadas como chaves de agrupamento. Consulte Usar clustering líquido para tabelas. |
schema |
str ou StructType |
Uma definição de esquema para a tabela. Os esquemas podem ser definidos como uma string DDL do SQL ou com um script Python StructType. |
private |
bool |
Crie uma tabela, mas não publique a tabela no metastore. Essa tabela está disponível para o pipeline, mas não está acessível fora do pipeline. As tabelas privadas persistem durante o tempo de vida do pipeline. O padrão é False. |
refresh_policy |
str |
Uma string que define a política de atualização para a visualização materializada. Um de: auto, incremental, , incremental_strictou full. Consulte a política de atualização e REFRESH a cláusula POLICY (pipelines).O padrão é auto. |
row_filter |
str |
(Versão prévia pública) Uma cláusula de filtro de linha para a tabela. Consulte Publicar as tabelas com os filtros de linha e as máscaras de coluna. |
Especificar um esquema é opcional e pode ser feito com o PySpark StructType ou o DDL do SQL. Ao especificar um esquema, opcionalmente, você pode incluir colunas geradas, máscaras de coluna e chaves primárias e estrangeiras. Consulte:
- Colunas geradas pelo Delta Lake
- Restrições no Azure Databricks
- Publicar tabelas com filtros de linha e máscaras de coluna.
Exemplos
from pyspark import pipelines as dp
# Specify a schema
sales_schema = StructType([
StructField("customer_id", StringType(), True),
StructField("customer_name", StringType(), True),
StructField("number_of_line_items", StringType(), True),
StructField("order_datetime", StringType(), True),
StructField("order_number", LongType(), True)]
)
@dp.materialized_view(
comment="Raw data on sales",
schema=sales_schema)
def sales():
return ("...")
# Specify a schema with SQL DDL, use a generated column, and set clustering columns
@dp.materialized_view(
comment="Raw data on sales",
schema="""
customer_id STRING,
customer_name STRING,
number_of_line_items STRING,
order_datetime STRING,
order_number LONG,
order_day_of_week STRING GENERATED ALWAYS AS (dayofweek(order_datetime))
""",
cluster_by = ["order_day_of_week", "customer_id"])
def sales():
return ("...")
# Use automatic liquid clustering to let Databricks choose the clustering columns
@dp.materialized_view(
comment="Raw data on sales",
cluster_by_auto=True)
def sales():
return ("...")
# Specify partition columns
@dp.materialized_view(
comment="Raw data on sales",
schema="""
customer_id STRING,
customer_name STRING,
number_of_line_items STRING,
order_datetime STRING,
order_number LONG,
order_day_of_week STRING GENERATED ALWAYS AS (dayofweek(order_datetime))
""",
partition_cols = ["order_day_of_week"])
def sales():
return ("...")
# Specify table constraints
@dp.materialized_view(
schema="""
customer_id STRING NOT NULL PRIMARY KEY,
customer_name STRING,
number_of_line_items STRING,
order_datetime STRING,
order_number LONG,
order_day_of_week STRING GENERATED ALWAYS AS (dayofweek(order_datetime)),
CONSTRAINT fk_customer_id FOREIGN KEY (customer_id) REFERENCES main.default.customers(customer_id)
""")
def sales():
return ("...")
# Specify a row filter and column mask
@dp.materialized_view(
schema="""
id int COMMENT 'This is the customer ID',
name string COMMENT 'This is the customer full name',
region string,
ssn string MASK catalog.schema.ssn_mask_fn USING COLUMNS (region)
""",
row_filter = "ROW FILTER catalog.schema.us_filter_fn ON (region, name)")
def sales():
return ("...")
# Specify a refresh policy
@dp.materialized_view(
refresh_policy = 'incremental_strict'
)
def sales():
return ("...")