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 @table decorador pode ser usado para definir tabelas de streaming em um pipeline Lakeflow ou pipeline independente.
Para definir uma tabela de streaming, aplique @table a uma consulta que execute uma leitura de streaming em uma fonte de dados ou use a função create_streaming_table().
Observação
No módulo mais antigo dlt, o operador @table era usado para criar tabelas de streaming e visões materializadas. O @table operador no pyspark.pipelines módulo ainda funciona dessa forma, mas o Databricks recomenda usar o @materialized_view operador para criar exibições materializadas.
Você também pode aplicar @table fora de um pipeline, em um notebook em computação geral serverless, para definir uma tabela de streaming independente, com algumas limitações.
Veja Definir tabelas com os decoradores de pipelines.
Sintaxe
from pyspark import pipelines as dp
@dp.table(
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",
row_filter = "row-filter-clause",
private = False,
replace_using = ["<key-column>", "<key-column>"],
sequence_by = "<sequence-column>")
@dp.expect(...)
def <function-name>():
return (<query>)
Parâmetros
@dp.expect() é uma cláusula opcional de expectativa de pipelines do Lakeflow. 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 de streaming do Apache Spark 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. Para obter as propriedades de coluna com suporte na cadeia de caracteres DDL, consulte a seção Parâmetros de CREATE STREAMING TABLE. |
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.Tabelas privadas foram criadas anteriormente com o temporary parâmetro. |
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. |
replace_using |
list |
(Beta) As colunas-chave que identificam quais linhas de destino devem ser substituídas, que definem a consulta da tabela como um fluxo de SUBSTITUIR USANDO. Especifique pelo menos uma coluna. Requer uma fonte de streaming e sequence_by.
Veja Substituição parcial de snapshot com SUBSTITUIR USING flows. |
sequence_by |
str ou Column |
(Beta) A coluna que ordena as atualizações para um replace_using fluxo, de modo que a sequência mais alta para uma chave vence. É obrigatório quando replace_using for definido. |
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.table(
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.table(
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 ("...")
# Specify a schema with an identity column
@dp.table(
comment="Raw data on sales",
schema="""
order_id BIGINT GENERATED ALWAYS AS IDENTITY,
customer_name STRING,
order_datetime STRING
""")
def sales():
return ("...")
# Use automatic liquid clustering to let Databricks choose the clustering columns
@dp.table(
comment="Raw data on sales",
cluster_by_auto=True)
def sales():
return ("...")
# Specify partition columns
@dp.table(
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.table(
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.table(
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 ("...")
# Define the table's query as a REPLACE USING flow, keeping the latest row for each payment_id
@dp.table(
name="payments_current",
replace_using=["payment_id"],
sequence_by="payment_date")
def payments_current():
return spark.readStream.table("samples.wanderbricks.payments")