Usar Python com pipelines autônomos
Você pode criar e fazer refresh de views materializadas e tabelas de transmissão autônomas a partir de um Notebook usando Python. Isso permite gerenciar pipelines autônomos junto com seus outros fluxos de trabalho de notebook baseados em Python.
Há duas maneiras de fazer isso:
- Defina a tabela com os decoradores
pyspark.pipelines,@dp.materialized_viewe@dp.table. Use isto quando a lógica for mais fácil de expressar como código DataFrame. Consulte Definir tabelas com os decoradores de pipelines. - Envie as mesmas instruções SQL executadas por um warehouse do Databricks SQL, passando-as para
spark.sql(). Isso oferece toda a interface SQL autônoma de visualização materializada e tabela de transmissão, incluindo instruçõesREFRESHe programações de refresh. Consulte Enviar instruções SQL comspark.sql().
Fonte Python para pipelines autônomos exige um Notebook anexado ao **compute serverless de uso geral**. Não é possível usar Python para criar ou refresh pipelines autônomos de um warehouse do Databricks SQL, porque um warehouse executa comandos SQL, não Notebooks Python. Para usar um SQL warehouse em vez disso, consulte Usar visões materializadas autônomas e Usar tabelas de transmissão autônomas.
Requisitos
Para criar e refresh pipelines autônomos com Python, é necessário um Notebook anexado a um compute serverless geral no Databricks Runtime 18,1 ou acima. Para a lista completa de requisitos, incluindo disponibilidade regional e permissões, consulte Notebooks.
Definir tabelas com os decoradores de pipeline
Você pode definir uma view materializada ou tabela de transmissão autônoma com os mesmos decoradores usados em um LakeFlow Pipelines. Each decorated function defines one table. When you run the cell, Databricks creates the table and runs a serverless pipeline to populate it. The cell returns when the update finishes.
Os decoradores de pipelines exigem a versão de ambiente serverless 5 ou acima.
Define a materialized view
Use @dp.materialized_view em uma função que retorna um DataFrame em lotes. O exemplo a seguir cria a visualização materializada daily_booking_revenue a partir da tabela bookings no dataset de exemplo Wanderbricks:
from pyspark import pipelines as dp
from pyspark.sql import functions as F
@dp.materialized_view(name="main.default.daily_booking_revenue")
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)
Para definir uma tabela a partir de uma leitura de transmissão, use @dp.table em vez disso.
Definir uma tabela de transmissão
Use @dp.table em uma função que retorne um DataFrame de transmissão. O exemplo a seguir cria a tabela de transmissão bookings_raw a partir de uma leitura de transmissão da mesma tabela bookings:
from pyspark import pipelines as dp
@dp.table(name="main.default.bookings_raw")
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")
Se a função retornar um DataFrame em lotes, o @dp.table criará uma materialized view em vez disso. A única exceção é replace_where, que sempre resulta em uma tabela de transmissão. O exemplo a seguir mantém a receita diária de check-ins em ou após 1º de julho de 2025 atualizada, sem recomputar datas anteriores:
from pyspark import pipelines as dp
from pyspark.sql import functions as F
@dp.table(
name="main.default.booking_revenue_rw",
replace_where=F.col("check_in") >= F.to_date(F.lit("2025-07-01")),
)
def booking_revenue_rw():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)
Cada execução exclui as linhas que correspondem ao predicado e recalcula apenas esse intervalo. Consulte processamento em lotes com REPLACE WHERE flows.
Refresh de uma tabela
Para refresh uma tabela definida com um decorador, execute o código que a define novamente, por exemplo, executando novamente a célula do notebook, executando o notebook inteiro ou executando o notebook como um Job. Cada execução cria a tabela se ela não existir e a refresh se já existir.
Para reprocessar todos os dados disponíveis na fonte, passe full_refresh=True em qualquer um dos decoradores:
@dp.table(name="main.default.bookings_raw", full_refresh=True)
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")
You can't use a REFRESH statement on a table defined with a decorator, or programar refreshes with SCHEDULE or TRIGGER ON UPDATE. To refresh on a programar, define the table in SQL, or programar the Notebook as a Job. See Lakeflow Jobs.
Configurar a tabela
Os decoradores aceitam os mesmos parâmetros de dataset comuns que aceitam dentro de um pipeline, incluindo comment, table_properties, partition_cols, cluster_by, schema e spark_conf:
@dp.materialized_view(
name="main.default.daily_booking_revenue",
comment="Daily booking revenue.",
table_properties={"quality": "gold"},
cluster_by=["check_in"],
)
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)
For the parameter list, see materialized_view and table.
private=True não é compatível, pois uma tabela privada só pode ser lida por outros datasets no mesmo pipeline.
APIs sem compatibilidade
Uma tabela autônoma é um único dataset com um único fluxo, portanto, as APIs que descrevem as relações entre datasets não estão disponíveis. O seguinte gera um erro fora de um pipeline:
@dp.temporary_viewedp.create_streaming_table@dp.append_flowe outros fluxos adicionaisdp.create_auto_cdc_flowedp.create_auto_cdc_from_snapshot_flow@dp.replace_flowe o parâmetroreplace_using, que definem fluxos REPLACE USING. Consulte Substituição parcial de snapshot com fluxos REPLACE USING.dp.create_sink- Expectativas, como
@dp.expecte@dp.expect_or_fail
Para usar esses recursos, crie um Lakeflow pipeline em vez disso. Consulte Desenvolver código de pipeline com Python.
Enviar comandos SQL com spark.sql()
Em um Notebook Python, passe as mesmas instruções que executaria de um SQL warehouse do Databricks para spark.sql(). A sintaxe da view materializada autônoma e da tabela de transmissão é idêntica; apenas a forma como a instrução é submetida difere. Assim como em um warehouse, cada instrução CREATE ou REFRESH executa um pipeline serverless para processar a operação.
A sessão spark está disponível por default em Notebooks Databricks, portanto, nenhuma importação é necessária.
Criar uma visualização materializada
O exemplo a seguir cria a visualização materializada mv1 a partir da tabela-base base_table1:
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW mv1
AS SELECT
date,
sum(sales) AS sum_of_sales
FROM base_table1
GROUP BY date
""")
Para detalhes completos de CREATE MATERIALIZED VIEW, como atualizações agendadas e acionadas, consulte Criar uma visualização materializada.
Criar tabela de transmissão
O exemplo a seguir cria a tabela de transmissão sales da tabela raw_data:
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")
Para obter todos os detalhes de CREATE STREAMING TABLE, incluindo o carregamento de arquivos com o Auto Loader e a programação, consulte Usar tabelas de transmissão autônomas.
refresh uma view materializada ou tabela de transmissão
Utilize uma instrução REFRESH para atualizar uma tabela autônoma com os dados mais recentes de sua origem:
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")
No compute geral serverless, as atualizações são síncronas. Atualizações assíncronas (a palavra-chave ASYNC) não têm suporte. Consulte Compute geral serverless.
Parametrizar declarações
Para passar valores do seu código Python para uma instrução em vez de codificá-los rigidamente, use marcadores de parâmetros nomeados no SQL e forneça seus valores por meio do argumento args de spark.sql(). Utilize um marcador como :min_sales diretamente para valores literais. Envolva o marcador em IDENTIFIER() apenas quando o parâmetro for um nome de objeto, como uma tabela, view ou esquema, porque os identificadores não podem ser substituídos como valores de string simples.
O exemplo a seguir parametriza tanto o nome da visualização materializada quanto um valor de filtro.
mv_name = "main.sales.regional_sales"
min_sales = 1000
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW IDENTIFIER(:mv)
AS SELECT
region,
sum(sales) AS sum_of_sales
FROM base_table1
WHERE sales > :min_sales
GROUP BY region
""", args={
"mv": mv_name,
"min_sales": min_sales,
})
Para obter mais informações, consulte Marcadores de parâmetro e cláusula IDENTIFIER.
Executar outros comandos
É possível executar qualquer instrução de view materializada independente ou de tabela de transmissão de um Notebook Python, ao passá-la para spark.sql(), incluindo instruções para programar refresh, alterar uma tabela ou descartar uma tabela. Para entender como usar visualizações materializadas e tabelas de transmissão, incluindo sintaxe SQL, consulte Usar visualizações materializadas autônomas e Usar tabelas de transmissão autônomas.
Limitações
Visualizações materializadas autônomas e tabelas de transmissão criadas em compute serverless geral possuem limitações adicionais, como a falta de suporte para refreshes assíncronos e a ausência de atribuição de custos por tabela. Para a lista completa, consulte Compute Geral Serverless.
Como esses pipelines estão em execução em Serverless compute geral em vez de um SQL Warehouse, eles não herdam tags personalizadas de um warehouse delimitador. A propagação de tags do warehouse para system.billing.usage aplica-se apenas a materialized views e transmissão tables cujas instruções têm execução a partir de um SQL Warehouse. Consulte Atribuir custos ao SQL warehouse com tags personalizadas.
As tabelas definidas com os decoradores de pipelines têm as seguintes limitações adicionais:
- Não é possível refresh-los com uma instrução
REFRESHnem refresh comSCHEDULEouTRIGGER ON UPDATE. Consulte Refresh uma tabela. - Expectativas, fluxos adicionais, fluxos de captura de dados de alterações (CDC), coletores (sinks) e views temporárias não são aceitos. Consulte APIs não suportadas.
private=Truenão é suportado.