Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Вы можете создавать и обновлять автономные материализованные представления и потоковые таблицы из записной книжки с помощью Python. Это позволяет управлять автономными конвейерами наряду с другими рабочими процессами в блокнотах на Python.
Это можно сделать двумя способами.
- Определите таблицу с декораторами
pyspark.pipelines,@dp.materialized_viewи@dp.table. Используйте это, когда логику легче выразить в виде кода DataFrame. См. Определение таблиц с помощью декораторов конвейеров. - Отправьте те же SQL-инструкции, которые выполняет хранилище Databricks SQL, передав их в
spark.sql(). Это предоставляет полный SQL-интерфейс для автономного материализованного представления и потоковой таблицы, включая операторыREFRESHи расписания обновления. См. Отправить операторы SQL сspark.sql().
Для автономных конвейеров с исходным кодом Python требуется блокнот, подключённый к бессерверным вычислительным ресурсам общего назначения. Нельзя использовать Python для создания или обновления автономных конвейеров из хранилища Databricks SQL, так как хранилище выполняет SQL-инструкции, а не ноутбуки Python. Чтобы использовать хранилище SQL, см. Использование автономных материализованных представлений и Использование автономных потоковых таблиц.
Important
Создание и обновление автономных материализованных представлений и потоковой передачи таблиц из записной книжки на бессерверных общих вычислительных ресурсах доступно в бета-версии и доступно в отдельных регионах. См. записные книжки.
Requirements
Для создания и обновления автономных конвейеров с помощью Python требуется записная книжка, подключенная к бессерверным общим вычислениям в Databricks Runtime 18.1 или более поздней версии. Полный список требований, включая региональную доступность и разрешения, см. в записных книжках.
Определите таблицы с помощью декораторов конвейеров
Вы можете определить отдельное материализованное представление или потоковую таблицу с теми же декораторами, что и в конвейере Lakeflow. Каждая декорированная функция определяет одну таблицу. Когда вы запускаете ячейку, Azure Databricks создаёт таблицу и запускает бессерверный конвейер для её заполнения. Ячейка снова станет доступной после завершения обновления.
Предупреждение
Для декораторов конвейеров требуется версия бессерверной среды 5 или выше.
Определите материализированный вид
Используйте @dp.materialized_view на функции, которая возвращает пакетный DataFrame. Следующий пример создаёт материализованное представление daily_booking_revenue из таблицы bookings в образце набора данных 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"))
)
Чтобы определить таблицу на основе потокового чтения, используйте вместо этого @dp.table.
Определим таблицу потоков
Используйте @dp.table на функции, которая возвращает потоковый DataFrame. Следующий пример создаёт потоковую таблицу bookings_raw на основе потокового чтения той же таблицы 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")
Если функция возвращает пакетный DataFrame, @dp.table создаёт материализированный вид. Единственное исключение — replace_where, которое всегда дает потоковую таблицу. Следующий пример поддерживает ежедневную выручку от заселений начиная с 1 июля 2025 года в актуальном состоянии, без пересчёта более ранних дат:
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"))
)
Каждый запуск удаляет строки, соответствующие предикату, и пересчитывает только этот диапазон. См. раздел "Пакетная обработка с помощью потоков REPLACEWHERE".
Обновить таблицу
Чтобы обновить определённую таблицу с помощью декоратора, запустите код, который её определяет, например, заново запустив ячейку блокнота, запустив весь блокнот или запустив блокнот как задачу. Каждый запуск создаёт таблицу, если её нет, и обновляет её, если она есть.
Чтобы повторно обработать все данные, доступные в источнике, передайте full_refresh=True в любой из декораторов:
@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")
Нельзя использовать оператор REFRESH для таблицы, определённой с помощью декоратора, или планировать обновления с помощью SCHEDULE или TRIGGER ON UPDATE. Чтобы обновить график, определите таблицу в SQL или запланируйте блокнот как работу. Смотрите Задания Lakeflow.
Настройте таблицу
Декораторы принимают те же общие параметры набора данных, что и внутри конвейера, включая comment, table_properties, partition_cols, cluster_by, schema и 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"))
)
Для списка параметров см. materialized_view и таблицу.
private=True не поддерживается, поскольку приватная таблица может быть прочитана только другими наборами данных в том же конвейере.
Неподдерживаемые API
Отдельная таблица — это один набор данных с одним потоком, поэтому API, описывающие взаимосвязи между наборами данных, недоступны. Следующее вызывает ошибку за пределами конвейера:
-
@dp.temporary_viewиdp.create_streaming_table. -
@dp.append_flowи другие дополнительные потоки -
dp.create_auto_cdc_flowиdp.create_auto_cdc_from_snapshot_flow. -
@dp.replace_flowи параметрreplace_using, которые определяют потоки REPLACE USING. См. Частичная замена моментальных снимков с использованием потоков REPLACE USING. dp.create_sink- Ожидания, такие как
@dp.expectи@dp.expect_or_fail
Чтобы использовать их, вместо этого создайте конвейер Lakeflow. См. раздел Разработка кода конвейера с помощью Python.
Отправляйте SQL-выписки с помощью spark.sql()
В записной книжке Python передайте те же инструкции, что и в хранилище spark.sql()Databricks SQL. Автономный материализованный вид и синтаксис потоковой таблицы идентичен; только способ отправки инструкции отличается. Как и в хранилище, каждый оператор CREATE или REFRESH запускает бессерверный конвейер для выполнения этой операции.
Сеанс spark доступен по умолчанию в записных книжках Azure Databricks, поэтому импорт не требуется.
Создание материализованного представления
В следующем примере создается материализованное представление mv1 из базовой таблицы 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
""")
Полные CREATE MATERIALIZED VIEW сведения, такие как запланированные и триггерные обновления, см. в разделе "Создание материализованного представления".
Создать потоковую таблицу
В следующем примере создается потоковая таблица sales из таблицы raw_data:
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")
Полные CREATE STREAMING TABLE сведения, включая загрузку файлов с помощью автозагрузчика и планирования, см. в разделе "Использование автономных таблиц потоковой передачи".
Обновление материализованного представления или потоковой таблицы
Используйте инструкцию REFRESH для обновления автономной таблицы с последними данными из источника:
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")
При бессерверных общих вычислениях обновления синхронны. Асинхронные обновления (ключевое ASYNC слово) не поддерживаются. См. бессерверные общие вычисления.
Операторы параметризации
Чтобы передать значения из кода Python в инструкцию вместо жесткого кода, используйте именованные маркеры параметров в SQL и укажите их значения с помощью args аргументаspark.sql(). Используйте маркер, например :min_sales непосредственно для литеральных значений. Заключайте маркер в IDENTIFIER() только в том случае, если параметр является именем объекта, например именем таблицы, представления или схемы, поскольку идентификаторы нельзя подставлять как обычные строковые значения.
В следующем примере параметризуется имя материализованного представления и значение фильтра:
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,
})
Дополнительные сведения см. в разделах «Маркеры параметров» и IDENTIFIER«Предложение».
Выполнение других инструкций
Вы можете выполнить любой оператор для отдельного материализованного представления или потоковой таблицы из блокнота Python, передав его в spark.sql(), включая операторы планирования обновлений, изменения таблицы или удаления таблицы. Сведения об использовании материализованных представлений и потоковых таблиц, включая синтаксис SQL, см. в статье "Использование автономных материализованных представлений " и использование автономных таблиц потоковой передачи.
Ограничения
Отдельные материализованные представления и потоковые таблицы, созданные с использованием serverless general compute, имеют дополнительные ограничения, такие как отсутствие поддержки асинхронного обновления и отсутствие атрибуции затрат на уровне отдельных таблиц. Полный список см. в разделе "Бессерверные общие вычисления".
Поскольку эти конвейеры работают на бессерверных вычислительных ресурсах общего назначения, а не в SQL-хранилище, они не наследуют пользовательские теги от вышестоящего хранилища. Распространение тегов хранилища применяется только к system.billing.usage материализованным представлениям и потоковым таблицам, инструкции которых выполняются в SQL-хранилище. См. Как относить затраты к хранилищу SQL с помощью пользовательских тегов.
Таблицы, определенные с использованием декораторов конвейеров, имеют следующие дополнительные ограничения:
- Вы не можете обновить их с помощью
REFRESHили настроить обновление по расписанию с помощьюSCHEDULEилиTRIGGER ON UPDATE. См. Обновить таблицу. - Ожидания, дополнительные потоки, потоки сбора данных изменений (CDC), поглотители и временные просмотры не поддерживаются. См. Неподдерживаемые API.
- Функция
private=Trueне поддерживается.