Notiz
Zougrëff op dës Säit erfuerdert Autorisatioun. Dir kënnt probéieren, Iech unzemellen oder Verzeechnesser ze änneren.
Zougrëff op dës Säit erfuerdert Autorisatioun. Dir kënnt probéieren, Verzeechnesser ze änneren.
Important
Environment versions for Lakeflow pipelines are in Public Preview.
Pipelines with an environment version set run Python code through Spark Connect. This page covers what is incompatible, what behaves differently, and how Databricks scans a pipeline for affected patterns.
Limitations
Environment versions are not yet compatible with all pipeline functionality. A pipeline run with an environment version set fails if the pipeline's Python code does any of the following:
- Mutates Spark session state inside a function decorated with a pipelines decorator. Examples include
spark.conf.set(...),spark.sql("USE CATALOG ..."), andcreateOrReplaceTempView. - Uses PySpark APIs that are unavailable in Spark Connect, including
SparkContext,RDD,SQLContext, and any Py4J APIs. See What is supported in Spark Connect.
If enabling an environment version on a pipeline causes it to fail, disabling the environment version returns the pipeline to its previous state.
Behavior changes
Spark Connect has a small number of behavior differences from the classic PySpark runtime. See Spark Connect vs. classic Spark for the full reference. The Compatibility scan detects these patterns ahead of time and blocks migration until they are addressed, so you can find and fix them before they affect production data.
In a pipeline, the most common situations where behavior may differ are:
Interleaved DataFrame construction and session mutation
When a pipeline constructs a DataFrame, then mutates Spark session state (for example, changes the default catalog or schema, sets a config, replaces a temp view, or re-registers a UDF), then uses the DataFrame:
- Without an environment version, the DataFrame uses the pre-mutation session state.
- With an environment version, the DataFrame uses the post-mutation session state.
For example:
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
Without an environment version, mytable contains [(1, "Original Row")]. With an environment version, mytable contains [(2, "Replaced Row")].
UDFs that reference mutable Python state
When a UDF references a Python global variable whose value changes after the UDF is defined:
- Without an environment version, the UDF uses the latest value of the variable.
- With an environment version, the UDF uses the value at the time the UDF was defined.
For example:
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")))
Without an environment version, my_mv contains [("alex_b",)]. With an environment version, my_mv contains [("alex_a",)].
If a pipeline relies on either pattern, audit it before enabling an environment version.
Compatibility scan
The compatibility scan finds code patterns in your pipeline that would produce different results under an environment version, so you can fix them before a pipeline is migrated automatically. When the scan is enabled on a pipeline:
- Each update emits one
BehaviorChangeInSparkConnectWARNevent in the pipeline event log per detected pattern. - The pipeline is not migrated to an environment version, and you cannot enable one yourself, until all compatibility warnings from the previous successful update are addressed.
This check doesn't apply to a pipeline that has no previous update, or that already has an environment version set.
Enable the scan on a pipeline
You can enable the compatibility scan by adding the pipelines.environmentVersion.enableCompatibilityScan pipeline configuration. You can add configuration through the pipeline editor UI or by adding an entry to the pipeline configuration JSON.
Through the UI:
- From the pipeline editor, click Settings.
- Find the Configuration section in pipeline settings.
- Click
Add configuration.
- Enter
pipelines.environmentVersion.enableCompatibilityScanas the key andtrueas the value. - Save the pipeline settings.
In the pipeline JSON:
Add the following entry to the configuration block:
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Review and resolve compatibility warnings
To find and clear the patterns that block an environment version on your pipeline:
- Run the pipeline in dry run mode, then query the pipeline event log for
BehaviorChangeInSparkConnectWARNevents. Each event reports one detected pattern. See Compatibility events reference for the full list of issue codes, example patterns, and suggested fixes. - Update the pipeline code to remove the detected patterns following the suggested fix, and run the pipeline again.
- Repeat until a successful update emits no more compatibility events. The pipeline can then be migrated automatically, and you can also enable an environment version yourself.
Enabling an environment version runs the same safety checks whether Databricks migrates the pipeline automatically or you set environment_version yourself. A pipeline with unresolved compatibility warnings does not move to an environment version until the warnings are resolved. If the migration can't be completed safely, or fails for any reason, it stops before writing any data and the pipeline continues to run on its previous runtime.
When an update stops for one of these reasons, the pipeline event log and the update error message describe the cause and the steps to resolve it. Follow those steps and run the pipeline again to complete the migration. If you believe a compatibility warning is a false positive, resolve the flagged pattern or contact Azure Databricks support.
Compatibility events reference
When the compatibility scan runs on a pipeline, it emits one BehaviorChangeInSparkConnect WARN event in the pipeline event log per detected pattern. When the previous successful update detected any patterns, the pipeline is not migrated to an environment version until the patterns are addressed.
Each event reports a single issue code that identifies what was detected. To look up a code, find it in the Issue codes table — each row links to the category section that contains an example pattern and the suggested fix.
Event shape
BehaviorChangeInSparkConnect events follow the standard pipeline event log schema:
event_typeisbehavior_change_in_spark_connect.levelisWARN.detailscontains thebehavior_change_in_spark_connectobject, which has a singleissuefield. The issue value is one of the codes listed below.messageis a human-readable description of the detected pattern.
Issue codes
| Category | Issue code | Description |
|---|---|---|
| Database and catalog mutations | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
The default catalog was changed after a DataFrame was created. The existing DataFrame may resolve tables using the new default catalog. |
| Database and catalog mutations | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE CATALOG was called outside a function decorated by a pipelines decorator. The default catalog may change unexpectedly for subsequent operations. |
| Database and catalog mutations | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
The default database was changed after a DataFrame was created. The existing DataFrame may resolve tables using the new default database. |
| Database and catalog mutations | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE DATABASE was called outside a function decorated by a pipelines decorator. The default database may change unexpectedly for subsequent operations. |
| Eager execution within flow functions | CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function calls a checkpoint command. |
| Eager execution within flow functions | CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function eagerly creates a DataFrame view (createOrReplaceTempView or similar). |
| Eager execution within flow functions | CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function creates a resource profile. |
| Eager execution within flow functions | GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function calls spark.resources or a related resource API. |
| Eager execution within flow functions | MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function performs an eager MERGE INTO on a target table. |
| Eager execution within flow functions | ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function performs an eager Spark ML operation. |
| Eager execution within flow functions | REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function registers a Python data source. |
| Eager execution within flow functions | STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function operates on an active streaming query handle. |
| Eager execution within flow functions | STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function registers or removes a streaming query listener. |
| Eager execution within flow functions | STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function calls spark.streams to manage streaming queries. |
| Eager execution within flow functions | WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function performs an eager DataFrameWriterV2 operation. |
| Eager execution within flow functions | WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function performs an eager DataFrame.write operation. |
| Eager execution within flow functions | WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
The flow function starts a streaming query (writeStream.start()). |
| Spark configuration mutations | CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED |
spark.conf.set() or spark.conf.unset() was called inside a function decorated by a pipelines decorator. This is not supported with an environment version. |
| Spark configuration mutations | SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.set() was called outside a function decorated by a pipelines decorator after a DataFrame was created. The config change may affect the existing DataFrame at execution time. |
| Spark configuration mutations | UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.unset() was called outside a function decorated by a pipelines decorator after a DataFrame was created. The config change may affect the existing DataFrame at execution time. |
| Temporary view replacements | REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
A global temporary view was replaced after a DataFrame referencing it was created. The replacement may be reflected in the existing DataFrame. |
| Temporary view replacements | REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
A temporary view was replaced after a DataFrame referencing it was created. The replacement may be reflected in the existing DataFrame. |
| UDF and UDTF mutations | OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
A UDF was re-registered with the same name after a DataFrame referencing it was created. The existing DataFrame may use the new UDF definition. |
| UDF and UDTF mutations | OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
A UDTF was re-registered with the same name after a DataFrame referencing it was created. The existing DataFrame may use the new UDTF definition. |
| UDF and UDTF mutations | UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
A UDF references a global mutable Python variable. With an environment version, the UDF uses the value of the variable at the time the UDF was defined, not at invocation time. |
| UDF and UDTF mutations | UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
A UDTF references a global mutable Python variable. With an environment version, the UDTF uses the value of the variable at the time the UDTF was defined, not at invocation time. |
Database and catalog mutations
These issues are emitted when pipeline code mutates the default database or catalog. With an environment version, DataFrames constructed before the mutation may resolve tables using the new database or catalog.
Example pattern that triggers an event:
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()
Without an environment version, df resolves events from the marketing catalog. With an environment version, df resolves events from the sales catalog.
Suggested fix: Fully qualify table names so resolution does not depend on the default catalog or database, and avoid changing the default catalog or database between DataFrame creation and use.
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 configuration mutations
These issues are emitted when pipeline code mutates Spark configuration in ways that can change DataFrame behavior under an environment version.
Example pattern that triggers an event:
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")
Without an environment version, the cast uses the conf value at DataFrame creation time. With an environment version, the cast uses spark.sql.ansi.enabled=true and may fail on invalid input.
Suggested fix: Set all required Spark configurations at the top of the pipeline file, before any DataFrame is created. For per-query configuration, use the pipeline's configuration setting in the pipeline spec.
Temporary view replacements
These issues are emitted when pipeline code replaces a temporary view after a DataFrame referencing it was created. With an environment version, the existing DataFrame may reflect the new view contents.
Example pattern that triggers an event:
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
Without an environment version, mytable contains [(1, "Original Row")]. With an environment version, mytable contains [(2, "Replaced Row")].
Suggested fix: Create each temporary view a single time and do not replace it. If you need multiple views with related data, give each a distinct name.
UDF and UDTF mutations
These issues are emitted when pipeline code mutates a UDF or UDTF in ways that change behavior under an environment version.
Example pattern that triggers an event:
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")))
Without an environment version, my_mv contains [("alex_b",)]. With an environment version, my_mv contains [("alex_a",)].
Suggested fix: Pass values into the UDF as arguments instead of capturing them from Python globals, or set the global before defining the UDF and do not mutate it afterward.
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")))
Eager execution within flow functions
These issues are emitted when pipeline code performs an eager Spark command inside a function decorated by a pipelines decorator (@table, @materialized_view, etc.). Flow functions are expected to define and return a DataFrame; eager commands that write data, manage streaming queries, register resources, or run ML operations are not allowed inside a flow function with an environment version set.
Suggested fix: Move the eager operation outside the flow function and return a DataFrame from the flow function instead. Side-effects such as writing to a table or starting a streaming query belong outside the pipeline definition; the pipeline engine handles materialization of the DataFrame returned by the flow function.
Find compatibility events in the event log
The following query returns all compatibility events for a pipeline, ordered most recent first:
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;
To count events by issue code across recent updates:
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;
For how to query the event log, see Query the event log.
Additional resources
- Configure environment versions for pipelines — feature overview, automatic migration, and how to enable an environment version yourself.
- Pipeline event log schema — full pipeline event log schema.
- Pipeline event log — how to query the pipeline event log.