Skip to main content

Use Python with standalone pipelines

You can create and refresh standalone materialized views and streaming tables from a notebook using Python. This lets you manage standalone pipelines alongside your other Python-based notebook workflows.

There are two ways to do this:

  • Define the table with the pyspark.pipelines decorators, @dp.materialized_view and @dp.table. Use this when the logic is easier to express as DataFrame code. See Define tables with the pipelines decorators.
  • Submit the same SQL statements a Databricks SQL warehouse runs by passing them to spark.sql(). This gives you the full standalone materialized view and streaming table SQL surface, including REFRESH statements and refresh schedules. See Submit SQL statements with spark.sql().

Python source for standalone pipelines requires a notebook attached to serverless general compute. You can't use Python to create or refresh standalone pipelines from a Databricks SQL warehouse, because a warehouse runs SQL statements, not Python notebooks. To use a SQL warehouse instead, see Use standalone materialized views and Use standalone streaming tables.

Beta

Creating and refreshing standalone materialized views and streaming tables from a notebook on serverless general compute is in Beta and available in select regions. See Notebooks.

Requirements​

To create and refresh standalone pipelines with Python, you need a notebook attached to serverless general compute on Databricks Runtime 18.1 or above. For the complete list of requirements, including regional availability and permissions, see Notebooks.

Define tables with the pipelines decorators​

You can define a standalone materialized view or streaming table with the same decorators you use in a Lakeflow pipeline. 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.

warning

The pipelines decorators require serverless environment version 5 or above.

Define a materialized view​

Use @dp.materialized_view on a function that returns a batch DataFrame. The following example creates the materialized view daily_booking_revenue from the bookings table in the Wanderbricks sample dataset:

Python
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"))
)

To define a table from a streaming read, use @dp.table instead.

Define a streaming table​

Use @dp.table on a function that returns a streaming DataFrame. The following example creates the streaming table bookings_raw from a streaming read of the same bookings table:

Python
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")

If the function returns a batch DataFrame, @dp.table creates a materialized view instead. The one exception is replace_where, which always results in a streaming table. The following example keeps daily revenue for check-ins on or after July 1, 2025 up to date, without recomputing earlier dates:

Python
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"))
)

Each run deletes the rows matching the predicate and recomputes only that range. See Batch processing with REPLACE WHERE flows.

Refresh a table​

To refresh a table you defined with a decorator, run the code that defines it again, for example by re-running the notebook cell, running the whole notebook, or running the notebook as a job. Each run creates the table if it doesn't exist and refreshes it if it does.

To reprocess all data available in the source, pass full_refresh=True on either decorator:

Python
@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 schedule refreshes with SCHEDULE or TRIGGER ON UPDATE. To refresh on a schedule, define the table in SQL, or schedule the notebook as a job. See Lakeflow Jobs.

Configure the table​

The decorators accept the same common dataset parameters that they accept inside a pipeline, including comment, table_properties, partition_cols, cluster_by, schema, and spark_conf:

Python
@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 is not supported, because a private table can only be read by other datasets in the same pipeline.

Unsupported APIs​

A standalone table is a single dataset with a single flow, so the APIs that describe relationships between datasets aren't available. The following raise an error outside a pipeline:

  • @dp.temporary_view and dp.create_streaming_table
  • @dp.append_flow and other additional flows
  • dp.create_auto_cdc_flow and dp.create_auto_cdc_from_snapshot_flow
  • @dp.replace_flow and the replace_using parameter, which define REPLACE USING flows. See Partial snapshot replacement with REPLACE USING flows.
  • dp.create_sink
  • Expectations, such as @dp.expect and @dp.expect_or_fail

To use these, author a Lakeflow pipeline instead. See Develop pipeline code with Python.

Submit SQL statements with spark.sql()​

In a Python notebook, pass the same statements you would run from a Databricks SQL warehouse to spark.sql(). The standalone materialized view and streaming table syntax is identical; only the way you submit the statement differs. As with a warehouse, each CREATE or REFRESH statement runs a serverless pipeline to process the operation.

The spark session is available by default in Databricks notebooks, so no import is required.

Create a materialized view​

The following example creates the materialized view mv1 from the base table base_table1:

Python
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW mv1
AS SELECT
date,
sum(sales) AS sum_of_sales
FROM base_table1
GROUP BY date
""")

For full CREATE MATERIALIZED VIEW details, such as scheduled and triggered refreshes, see Create a materialized view.

Create a streaming table​

The following example creates the streaming table sales from the raw_data table:

Python
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")

For full CREATE STREAMING TABLE details, including loading files with Auto Loader and scheduling, see Use standalone streaming tables.

Refresh a materialized view or streaming table​

Use a REFRESH statement to update a standalone table with the latest data from its source:

Python
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")

On serverless general compute, refreshes are synchronous. Asynchronous refreshes (the ASYNC keyword) are not supported. See Serverless general compute.

Parameterize statements​

To pass values from your Python code into a statement instead of hardcoding them, use named parameter markers in the SQL and supply their values through the args argument of spark.sql(). Use a marker such as :min_sales directly for literal values. Wrap the marker in IDENTIFIER() only when the parameter is an object name, such as a table, view, or schema, because identifiers can't be substituted as plain string values.

The following example parameterizes both the materialized view name and a filter value:

Python
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,
})

For more information, see Parameter markers and IDENTIFIER clause.

Run other statements​

You can run any standalone materialized view or streaming table statement from a Python notebook by passing it to spark.sql(), including statements to schedule refreshes, alter a table, or drop a table. To understand how to use materialized views and streaming tables, including SQL syntax, see Use standalone materialized views and Use standalone streaming tables.

Limitations​

Standalone materialized views and streaming tables created on serverless general compute have additional limitations, such as no support for asynchronous refreshes and no per-table cost attribution. For the full list, see Serverless general compute.

Because these pipelines run on serverless general compute rather than a SQL warehouse, they do not inherit custom tags from an enclosing warehouse. Warehouse tag propagation to system.billing.usage applies only to materialized views and streaming tables whose statements run from a SQL warehouse. See Attribute costs to the SQL warehouse with custom tags.

Tables defined with the pipelines decorators have the following additional limitations:

  • You can't refresh them with a REFRESH statement or schedule refreshes with SCHEDULE or TRIGGER ON UPDATE. See Refresh a table.
  • Expectations, additional flows, change data capture (CDC) flows, sinks, and temporary views are not supported. See Unsupported APIs.
  • private=True is not supported.

Additional resources​