Совместимость версий среды

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 . Сканирование совместимости заранее обнаруживает эти закономерности и блокирует миграцию до тех пор, пока они не будут устранены, чтобы вы могли найти и исправить их до того, как они повлияют на производственные данные.

В конвейере наиболее распространенные ситуации, в которых поведение может отличаться:

Переключения структуры кадра данных и изменения сеанса

Когда конвейер создает кадр данных, затем изменяет состояние сеанса 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 конфигурации конвейера.

Через пользовательский интерфейс:

  1. В редакторе конвейера нажмите кнопку "Параметры".
  2. Найдите раздел "Конфигурация" в параметрах конвейера.
  3. Щелкните значок Добавьте конфигурацию.
  4. Введите pipelines.environmentVersion.enableCompatibilityScan в качестве ключа и true в качестве значения.
  5. Сохраните параметры конвейера.

В конвейере JSON:

Добавьте в блок следующую configuration запись:

"configuration": {
  "pipelines.environmentVersion.enableCompatibilityScan": "true"
}

Проверьте и устраните предупреждения о совместимости

Чтобы найти и очистить шаблоны, блокирующие версию среды в вашем пайплайне:

  1. Запустите конвейер в режиме пробного запуска , затем запросите журнал событий конвейера по событиям BehaviorChangeInSparkConnectWARN . Каждое событие сообщает один обнаруженный шаблон. Дополнительные сведения см. в справочнике по событиям совместимости для полного списка кодов проблем, примеров шаблонов и предлагаемых исправлений.
  2. Обновите код конвейера, чтобы удалить обнаруженные шаблоны после предложенного исправления, и запустите конвейер заново.
  3. Повторяйте, пока успешное обновление не приведёт к отсутствию событий совместимости. Конвейер можно мигрировать автоматически, а также самостоятельно включить версию среды.

Включение версии среды запускает те же проверки безопасности, независимо от того, мигрирует ли 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;

Сведения о том, как запросить журнал событий, см. в статье "Запрос журнала событий".

Дополнительные ресурсы