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.pipelinesdecorators,@dp.materialized_viewand@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, includingREFRESHstatements and refresh schedules. See Submit SQL statements withspark.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.
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.
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:
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:
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:
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:
@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:
@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_viewanddp.create_streaming_table@dp.append_flowand other additional flowsdp.create_auto_cdc_flowanddp.create_auto_cdc_from_snapshot_flow@dp.replace_flowand thereplace_usingparameter, which define REPLACE USING flows. See Partial snapshot replacement with REPLACE USING flows.dp.create_sink- Expectations, such as
@dp.expectand@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:
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:
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:
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:
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
REFRESHstatement or schedule refreshes withSCHEDULEorTRIGGER 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=Trueis not supported.