Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Версии окружающей среды для трубопроводов Lakeflow находятся в публичном предварительном просмотре.
Конвейеры с версией environment запускайте код Python через Spark Connect. На этой странице рассматривается, что несовместимо, что ведёт себя иначе, а также то, как Databricks сканирует конвейер на наличие затронутых паттернов.
Ограничения
Версии среды еще не совместимы со всеми функциями конвейера. Запуск конвейера с набором версий среды завершается ошибкой, если Python код конвейера выполняет любое из следующих действий:
- Изменяет состояние сеанса Spark внутри функции, украшенной декоратором конвейеров. Например,
spark.conf.set(...),spark.sql("USE CATALOG ...")иcreateOrReplaceTempView. - Использует API PySpark, недоступные в Spark Connect, включая
SparkContext,RDDSQLContextи любые API Py4J. См. раздел "Что поддерживается в Spark Connect".
Если включение версии среды в конвейере приводит к сбою, отключение версии среды возвращает конвейер в предыдущее состояние.
Изменения в поведении
Spark Connect имеет небольшое количество различий в поведении от классической среды выполнения PySpark. Полный справочник см. в разделе Spark Connect и классический Spark . Сканирование совместимости заранее обнаруживает эти закономерности и блокирует миграцию до тех пор, пока они не будут устранены, чтобы вы могли найти и исправить их до того, как они повлияют на производственные данные.
В конвейере наиболее распространенные ситуации, в которых поведение может отличаться:
- Переключения структуры кадра данных и изменения сеанса
- UDFs, ссылающиеся на изменяемое состояние Python
Переключения структуры кадра данных и изменения сеанса
Когда конвейер создает кадр данных, затем изменяет состояние сеанса Spark (например, изменяет каталог или схему по умолчанию, задает конфигурацию, заменяет временное представление или повторно регистрирует UDF), а затем использует кадр данных:
- Без версии среды кадр данных использует состояние сеанса предварительной мутации .
- С версией среды кадр данных использует состояние сеанса после мутации .
Рассмотрим пример.
from pyspark import pipelines as dp
spark.createDataFrame([(1, "Original Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
df = spark.sql("SELECT * FROM my_view")
spark.createDataFrame([(2, "Replaced Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
@dp.materialized_view
def mytable():
return df
Без версии mytable среды содержится [(1, "Original Row")]. С версией mytable среды содержится [(2, "Replaced Row")].
UDFs, ссылающиеся на изменяемое состояние Python
Когда UDF ссылается на глобальную переменную Python, значение которой изменяется после определения UDF:
- Без версии среды UDF использует последнее значение переменной.
- С версией среды UDF использует значение во время определения UDF.
Рассмотрим пример.
from pyspark import pipelines as dp
from pyspark.sql.functions import col, udf
suffix = "a"
@udf
def my_udf(s):
return s + suffix
suffix = "b"
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(my_udf(col("name")))
Без версии my_mv среды содержится [("alex_b",)]. С версией my_mv среды содержится [("alex_a",)].
Если конвейер использует любой шаблон, выполните аудит перед включением версии среды.
Проверка совместимости
Сканирование совместимости обнаруживает шаблоны кода в вашем конвейере, которые дают разные результаты в версии среды, так что вы можете исправить их до автоматической миграции конвейера. Если сканирование включено в конвейере:
- Каждое обновление генерирует одно
BehaviorChangeInSparkConnectWARNсобытие в журнале событий конвейера для каждого обнаруженного шаблона. - Конвейер не мигрируется в версию среды, и вы не можете включить её самостоятельно, пока не будут устранены все предупреждения о совместимости при предыдущем успешном обновлении.
Эта проверка не применяется к конвейеру, у которого не было предыдущего обновления или у которого уже установлена версия среды.
Включение проверки в конвейере
Вы можете включить проверку совместимости, добавив конфигурацию конвейера pipelines.environmentVersion.enableCompatibilityScan . Вы можете добавить конфигурацию через пользовательский интерфейс редактора конвейера или добавить запись в JSON конфигурации конвейера.
Через пользовательский интерфейс:
- В редакторе конвейера нажмите кнопку "Параметры".
- Найдите раздел "Конфигурация" в параметрах конвейера.
- Щелкните
Добавьте конфигурацию.
- Введите
pipelines.environmentVersion.enableCompatibilityScanв качестве ключа иtrueв качестве значения. - Сохраните параметры конвейера.
В конвейере JSON:
Добавьте в блок следующую configuration запись:
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Проверьте и устраните предупреждения о совместимости
Чтобы найти и очистить шаблоны, блокирующие версию среды в вашем пайплайне:
- Запустите конвейер в режиме пробного запуска , затем запросите журнал событий конвейера по событиям
BehaviorChangeInSparkConnectWARN. Каждое событие сообщает один обнаруженный шаблон. Дополнительные сведения см. в справочнике по событиям совместимости для полного списка кодов проблем, примеров шаблонов и предлагаемых исправлений. - Обновите код конвейера, чтобы удалить обнаруженные шаблоны после предложенного исправления, и запустите конвейер заново.
- Повторяйте, пока успешное обновление не приведёт к отсутствию событий совместимости. Конвейер можно мигрировать автоматически, а также самостоятельно включить версию среды.
Включение версии среды запускает те же проверки безопасности, независимо от того, мигрирует ли Databricks конвейер автоматически или вы сами устанавливаете environment_version . Конвейер с неразрешенными предупреждениями о совместимости не переходит в версию среды, пока предупреждения не будут устранены. Если миграция не удаётся безопасно завершить или по какой-либо причине не удаётся, она останавливается до записи данных, и конвейер продолжает работать в предыдущем режиме.
Когда обновление останавливается по одной из этих причин, журнал событий конвейера и сообщение об ошибке обновления описывают причину и шаги по её устранению. Следуйте этим шагам и запустите конвейер заново, чтобы завершить миграцию. Если вы считаете, что предупреждение о совместимости является ложным срабатыванием, устраните отмеченный шаблон или обратитесь в службу поддержки Azure Databricks.
Справочник по событиям совместимости
Когда сканирование совместимости выполняется на конвейере, оно генерирует одно BehaviorChangeInSparkConnectWARN событие в журнале событий конвейера для каждого обнаруженного шаблона. Когда предыдущее успешное обновление обнаружило какие-либо шаблоны, конвейер не мигрируется в версию среды, пока эти шаблоны не будут устранены.
Каждое событие сообщает один код проблемы, определяющий обнаруженное событие. Чтобы найти код, найдите его в таблице кодов проблем — каждая строка ссылается на раздел категории, содержащий пример шаблона и предлагаемое исправление.
Фигура события
BehaviorChangeInSparkConnect События следуют стандартной схеме журнала событий конвейера:
-
event_typeравноbehavior_change_in_spark_connect. -
levelравноWARN. -
detailsсодержитbehavior_change_in_spark_connectобъект, имеющий одноissueполе. Значение проблемы является одним из кодов, перечисленных ниже. -
message— это удобочитаемое описание обнаруженного шаблона.
Коды проблем
| Категория | Код проблемы | Description |
|---|---|---|
| Изменения базы данных и каталога | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Каталог по умолчанию был изменен после создания кадра данных. Существующий кадр данных может разрешать таблицы с помощью нового каталога по умолчанию. |
| Изменения базы данных и каталога | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE CATALOG вызывается вне функции, украшенной декоратором конвейеров. Каталог по умолчанию может неожиданно измениться для последующих операций. |
| Изменения базы данных и каталога | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
База данных по умолчанию была изменена после создания кадра данных. Существующий кадр данных может разрешать таблицы с помощью новой базы данных по умолчанию. |
| Изменения базы данных и каталога | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE DATABASE вызывается вне функции, украшенной декоратором конвейеров. База данных по умолчанию может неожиданно измениться для последующих операций. |
| Активное выполнение в функциях потока | CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока вызывает команду контрольной точки. |
| Активное выполнение в функциях потока | CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока с нетерпением создает представление кадра данных (createOrReplaceTempView или аналогичное). |
| Активное выполнение в функциях потока | CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока создает профиль ресурса. |
| Активное выполнение в функциях потока | GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Вызовы spark.resources функции потока или связанный API ресурсов. |
| Активное выполнение в функциях потока | MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока выполняет стремлении MERGE INTO к целевой таблице. |
| Активное выполнение в функциях потока | ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока выполняет страстную операцию Машинного обучения Spark. |
| Активное выполнение в функциях потока | REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока регистрирует источник данных Python. |
| Активное выполнение в функциях потока | STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока работает с активным дескриптором потокового запроса. |
| Активное выполнение в функциях потока | STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока регистрирует или удаляет прослушиватель потокового запроса. |
| Активное выполнение в функциях потока | STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока вызывает управление spark.streams запросами потоковой передачи. |
| Активное выполнение в функциях потока | WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока выполняет страстную DataFrameWriterV2 операцию. |
| Активное выполнение в функциях потока | WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока выполняет страстную DataFrame.write операцию. |
| Активное выполнение в функциях потока | WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Функция потока запускает потоковый запрос (writeStream.start()). |
| Изменения конфигурации Spark | CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED |
spark.conf.set() или spark.conf.unset() вызывается внутри функции, украшенной декоратором конвейеров. Это не поддерживается версией среды. |
| Изменения конфигурации Spark | SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.set() вызывается вне функции, украшенной декоратором конвейеров после создания кадра данных. Изменение конфигурации может повлиять на существующий кадр данных во время выполнения. |
| Изменения конфигурации Spark | UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.unset() вызывается вне функции, украшенной декоратором конвейеров после создания кадра данных. Изменение конфигурации может повлиять на существующий кадр данных во время выполнения. |
| Временные замены представлений | REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Глобальное временное представление было заменено после создания ссылки на кадр данных. Замена может быть отражена в существующем кадре данных. |
| Временные замены представлений | REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Временное представление было заменено после создания ссылки на кадр данных. Замена может быть отражена в существующем кадре данных. |
| Изменения UDF и UDTF | OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Пользовательская служба была повторно зарегистрирована с тем же именем после создания ссылки на кадр данных. Существующий кадр данных может использовать новое определение UDF. |
| Изменения UDF и UDTF | OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
UDTF был повторно зарегистрирован с тем же именем после создания ссылки на кадр данных. Существующий кадр данных может использовать новое определение UDTF. |
| Изменения UDF и UDTF | UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDF ссылается на глобальную переменную Python изменяемой Python. С версией среды UDF использует значение переменной во время определения UDF, а не во время вызова. |
| Изменения UDF и UDTF | UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDTF ссылается на глобальную переменную Python. С версией среды UDTF использует значение переменной во время определения UDTF, а не во время вызова. |
Изменения базы данных и каталога
Эти проблемы возникают при изменении кода конвейера базы данных или каталога по умолчанию. С версией среды кадры данных, созданные до изменения, могут разрешать таблицы с помощью новой базы данных или каталога.
Пример шаблона, который активирует событие:
from pyspark import pipelines as dp
spark.sql("USE CATALOG marketing")
df = spark.read.table("events")
spark.sql("USE CATALOG sales") # changes the default catalog after df was created
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Без версии df среды разрешается events из marketing каталога. В версии df среды разрешается events из sales каталога.
Предлагаемое исправление: Полное определение имен таблиц, поэтому разрешение не зависит от каталога или базы данных по умолчанию, и избегайте изменения каталога или базы данных по умолчанию между созданием и использованием кадра данных.
from pyspark import pipelines as dp
df = spark.read.table("marketing.default.events")
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Изменения конфигурации Spark
Эти проблемы возникают, когда код конвейера мутирует конфигурацию Spark способами, которые могут изменять поведение кадра данных в версии среды.
Пример шаблона, который активирует событие:
from pyspark import pipelines as dp
df = spark.read.table("events")
spark.conf.set("spark.sql.ansi.enabled", "true") # changes session conf after df was created
@dp.materialized_view
def events_strict():
return df.selectExpr("CAST(price AS INT) AS price")
Без версии среды приведение использует значение conf во время создания кадра данных. При использовании версии среды приведение использует spark.sql.ansi.enabled=true недопустимые входные данные и может завершиться ошибкой.
Предлагаемое исправление: Задайте все необходимые конфигурации Spark в верхней части файла конвейера перед созданием любого кадра данных. Для конфигурации каждого запроса используйте параметр конвейера configuration в спецификации конвейера.
Временные замены представлений
Эти проблемы возникают, когда код конвейера заменяет временное представление после создания ссылки на кадр данных. В версии среды существующий кадр данных может отражать новое содержимое представления.
Пример шаблона, который активирует событие:
from pyspark import pipelines as dp
spark.createDataFrame([(1, "Original Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
df = spark.sql("SELECT * FROM my_view")
spark.createDataFrame([(2, "Replaced Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
@dp.materialized_view
def mytable():
return df
Без версии mytable среды содержится [(1, "Original Row")]. С версией mytable среды содержится [(2, "Replaced Row")].
Предлагаемое исправление: Создайте каждое временное представление за один раз и не замените его. Если требуется несколько представлений со связанными данными, присвойте каждому отдельному имени.
Изменения UDF и UDTF
Эти проблемы возникают, когда код конвейера мутирует UDF или UDTF способами изменения поведения в версии среды.
Пример шаблона, который активирует событие:
from pyspark import pipelines as dp
from pyspark.sql.functions import col, udf
suffix = "a"
@udf
def my_udf(s):
return s + suffix
suffix = "b"
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(my_udf(col("name")))
Без версии my_mv среды содержится [("alex_b",)]. С версией my_mv среды содержится [("alex_a",)].
Suggested fix: передать значения в UDF в качестве аргументов вместо записи из Python глобальных или задать глобальный перед определением UDF и не мутировать его после этого.
from pyspark import pipelines as dp
from pyspark.sql.functions import col, lit, udf
@udf
def append_suffix(s, suffix):
return s + suffix
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(append_suffix(col("name"), lit("b")))
Активное выполнение в функциях потока
Эти проблемы возникают, когда код конвейера выполняет страстную команду Spark внутри функции, украшенной декоратором конвейеров (@table, @materialized_viewи т. д.). Ожидается, что функции потока определяют и возвращают кадр данных; команды, которые записывают данные, управляют потоковыми запросами, регистрируют ресурсы или выполняют операции машинного обучения, не допускаются в функции потока с набором версий среды.
Предлагаемое исправление: Переместите страстную операцию вне функции потока и верните кадр данных из функции потока. Побочные эффекты, такие как запись в таблицу или запуск потокового запроса, относятся за пределами определения конвейера; Обработчик конвейера обрабатывает материализацию кадра данных, возвращаемого функцией потока.
Поиск событий совместимости в журнале событий
Следующий запрос возвращает все события совместимости для конвейера, упорядоченные в первую очередь:
SELECT
timestamp,
message,
details:behavior_change_in_spark_connect:issue AS issue
FROM event_log(<pipeline-id>)
WHERE event_type = 'behavior_change_in_spark_connect'
AND level = 'WARN'
ORDER BY timestamp DESC;
Чтобы подсчитать события по коду проблемы в последних обновлениях:
SELECT
details:behavior_change_in_spark_connect:issue AS issue,
COUNT(*) AS occurrences
FROM event_log(<pipeline-id>)
WHERE event_type = 'behavior_change_in_spark_connect'
AND level = 'WARN'
GROUP BY 1
ORDER BY occurrences DESC;
Сведения о том, как запросить журнал событий, см. в статье "Запрос журнала событий".
Дополнительные ресурсы
- Настройте версии окружений для конвейеров — обзор функций, автоматическую миграцию и как самостоятельно включить версию среды.
- Схема журнала событий конвейера — полная схема журнала событий конвейера .
- Журнал событий конвейера — как запросить журнал событий конвейера.