Kompatibilita verzí prostředí

Important

Verze prostředí pro Lakeflow pipeline jsou ve veřejném náhledu.

Kanály s verzí environment spouští Python kód prostřednictvím Spark Connect. Tato stránka se zabývá tím, co je nekompatibilní, co se chová odlišně a jak Databricks skenuje pipeline na ovlivněné vzory.

Omezení

Verze prostředí ještě nejsou kompatibilní se všemi funkcemi kanálu. Spuštění kanálu se sadou verzí prostředí selže, pokud kód kanálu Python provede některou z následujících věcí:

  • Ztlumí stav relace Sparku uvnitř funkce zdobené dekorátorem kanálů. Mezi příklady patří spark.conf.set(...), spark.sql("USE CATALOG ...")a createOrReplaceTempView.
  • Používá rozhraní PySpark API, která nejsou ve Spark Connect k dispozici, včetně SparkContextRDDSQLContext, a všech rozhraní API Py4J. Podívejte se, co je podporováno ve Spark Connectu.

Pokud povolení verze prostředí v kanálu způsobí selhání, zakázání verze prostředí vrátí kanál do předchozího stavu.

Změny chování

Spark Connect má malý počet rozdílů v chování oproti klasickému modulu runtime PySpark. Úplný odkaz najdete v tématu Spark Connect vs. Classic Spark . Kompatibilita scan tyto vzorce předem detekuje a blokuje migraci, dokud nejsou vyřešeny, takže je můžete najít a opravit dříve, než ovlivní produkční data.

V kanálu se nejběžnější situace, kdy se chování může lišit, jsou:

Prokládání konstrukce datového rámce a mutaci relací

Když kanál vytvoří datový rámec, ztlumí se stav relace Sparku (například změní výchozí katalog nebo schéma, nastaví konfiguraci, nahradí dočasné zobrazení nebo znovu zaregistruje UDF), pak použije datový rámec:

  • Bez verze prostředí datový rámec používá stav relace před mutací .
  • U verze prostředí datový rámec používá stav relace po mutaci .

Příklad:

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

Bez verze mytable prostředí obsahuje [(1, "Original Row")]. S verzí mytable prostředí obsahuje [(2, "Replaced Row")].

UDF, které odkazují na proměnlivý stav Python

Když funkce definovaná uživatelem odkazuje na globální proměnnou Python, jejíž hodnota se po definování definovaného uživatelem změní:

  • Bez verze prostředí používá funkce definovaná uživatelem nejnovější hodnotu proměnné.
  • U verze prostředí používá funkce definovaná uživatelem hodnotu v okamžiku, kdy byla definovaná funkce definovaná uživatelem.

Příklad:

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")))

Bez verze my_mv prostředí obsahuje [("alex_b",)]. S verzí my_mv prostředí obsahuje [("alex_a",)].

Pokud kanál spoléhá na některý ze vzorů, před povolením verze prostředí ho auditujte.

Kontrola kompatibility

Skenování kompatibility najde vzory kódu ve vašem pipeline, které by pod verzí prostředí produkovaly jiné výsledky, takže je můžete opravit před automatickou migrací pipeline. Pokud je kontrola povolená v kanálu:

  • Každá aktualizace vygeneruje jednu BehaviorChangeInSparkConnectWARN událost v logu událostí pipeline za každý detekovaný vzor.
  • Pipeline není migrován do verze pro prostředí a nemůžete ji povolit sami, dokud nejsou vyřešena všechna varování kompatibility z předchozí úspěšné aktualizace.

Tato kontrola se nevztahuje na pipeline, které nemá předchozí aktualizaci nebo které už má nastavenou verzi prostředí.

Povolení kontroly v kanálu

Kontrolu kompatibility můžete povolit přidáním pipelines.environmentVersion.enableCompatibilityScan konfigurace kanálu. Konfiguraci můžete přidat prostřednictvím uživatelského rozhraní editoru kanálů nebo přidáním položky do kódu JSON konfigurace kanálu.

Prostřednictvím uživatelského rozhraní:

  1. V editoru kanálů klikněte na Nastavení.
  2. V nastavení kanálu vyhledejte část Konfigurace .
  3. Klikněte na ikonu Plus.Přidejte konfiguraci.
  4. Zadejte pipelines.environmentVersion.enableCompatibilityScan jako klíč a true jako hodnotu.
  5. Uložte nastavení kanálu.

V kódu JSON kanálu:

Do bloku přidejte následující položku configuration :

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

Zkontrolujte a vyřešit varování kompatibility

Chcete-li najít a odstranit vzory, které blokují verzi prostředí ve vašem pipeline:

  1. Spusť pipeline v režimu suchého běhu a poté dotazuj na události pipeline event logBehaviorChangeInSparkConnectWARN. Každá událost hlásí jeden zjištěný vzor. Úplný seznam kódů problémů, vzorů příkladů a navrhovaných oprav najdete v referenčních informacích k událostem kompatibility .
  2. Aktualizujte kód pipeline, aby odstranili zjištěné vzory po navržené opravě, a pipeline spusťte znovu.
  3. Opakujte, dokud úspěšná aktualizace nevyprodukuje žádné další kompatibilitní události. Pipeline pak lze automaticky migrovat a můžete si také sami povolit verzi prostředí.

Povolení verze prostředí provádí stejné bezpečnostní kontroly, ať už Databricks automaticky migruje pipeline, nebo si to environment_version nastavíte sami. Pipeline s nevyřešenými varováními kompatibility nepřechází do verze prostředí, dokud nejsou varování vyřešena. Pokud migrace nemůže být bezpečně dokončena nebo z jakéhokoliv důvodu selže, zastaví se před zápisem jakýchkoli dat a pipeline pokračuje v běhu na předchozím runtime.

Když aktualizace skončí z některého z těchto důvodů, záznam událostí pipeline a chybová zpráva o aktualizaci popisují příčinu a kroky k jejímu vyřešení. Postupujte podle těchto kroků a spusťte pipeline znovu, abyste dokončili migraci. Pokud si myslíte, že varování kompatibility je falešně pozitivní, vyřešte označený vzorec nebo kontaktujte podporu Azure Databricks.

Referenční informace k událostem kompatibility

Když skenování kompatibility běží na pipeline, vygeneruje jednu BehaviorChangeInSparkConnectWARN událost v záznamu událostí pipeline za každý detekovaný vzor. Když předchozí úspěšná aktualizace detekovala jakékoli vzory, pipeline není migrováno do verze prostředí, dokud nejsou vzory vyřešeny.

Každá událost hlásí jeden kód problému, který identifikuje zjištěné informace. Pokud chcete vyhledat kód, najděte ho v tabulce Kódy problému – každý řádek odkazuje na oddíl kategorie, který obsahuje ukázkový vzor a navrhované opravy.

Obrazec události

BehaviorChangeInSparkConnect události se řídí standardním schématem protokolu událostí kanálu:

  • event_type je behavior_change_in_spark_connect.
  • level je WARN.
  • details behavior_change_in_spark_connect obsahuje objekt, který má jedno issue pole. Hodnota problému je jedním z níže uvedených kódů.
  • message je čitelný popis zjištěného vzoru.

Kódy problémů

Kategorie Kód problému Description
Databáze a katalogové mutaci USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR Výchozí katalog se po vytvoření datového rámce změnil. Existující datový rámec může přeložit tabulky pomocí nového výchozího katalogu.
Databáze a katalogové mutaci USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR USE CATALOG byl volán mimo funkci zdobenou dekorátorem kanálů. Výchozí katalog se může neočekávaně změnit pro následné operace.
Databáze a katalogové mutaci USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR Výchozí databáze se po vytvoření datového rámce změnila. Existující datový rámec může přeložit tabulky pomocí nové výchozí databáze.
Databáze a katalogové mutaci USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR USE DATABASE byl volán mimo funkci zdobenou dekorátorem kanálů. Výchozí databáze se může neočekávaně změnit pro následné operace.
Dychtivé provádění ve funkcích toku CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku volá příkaz kontrolního bodu.
Dychtivé provádění ve funkcích toku CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku dychtivě vytvoří zobrazení datového rámce (createOrReplaceTempView nebo podobné).
Dychtivé provádění ve funkcích toku CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku vytvoří profil prostředku.
Dychtivé provádění ve funkcích toku GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku volá spark.resources nebo související rozhraní API prostředků.
Dychtivé provádění ve funkcích toku MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku provádí s cílovou tabulkou dychtivou činnost MERGE INTO .
Dychtivé provádění ve funkcích toku ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku provádí dychtivou operaci Spark ML.
Dychtivé provádění ve funkcích toku REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku zaregistruje Python zdroj dat.
Dychtivé provádění ve funkcích toku STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku pracuje s aktivním popisovačem dotazu streamování.
Dychtivé provádění ve funkcích toku STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku zaregistruje nebo odebere naslouchací proces dotazu streamování.
Dychtivé provádění ve funkcích toku STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Volání funkce spark.streams toku pro správu streamovaných dotazů.
Dychtivé provádění ve funkcích toku WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku provádí dychtivou DataFrameWriterV2 operaci.
Dychtivé provádění ve funkcích toku WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku provádí dychtivou DataFrame.write operaci.
Dychtivé provádění ve funkcích toku WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED Funkce toku spustí dotaz streamování (writeStream.start()).
Změny konfigurace Sparku CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED spark.conf.set() nebo spark.conf.unset() byl volán uvnitř funkce zdobené dekorátorem kanálů. Tato možnost není podporována ve verzi prostředí.
Změny konfigurace Sparku SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR spark.conf.set() byl volán mimo funkci dekorátoru kanálů po vytvoření datového rámce. Změna konfigurace může mít vliv na existující datový rámec v době provádění.
Změny konfigurace Sparku UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR spark.conf.unset() byl volán mimo funkci dekorátoru kanálů po vytvoření datového rámce. Změna konfigurace může mít vliv na existující datový rámec v době provádění.
Dočasné nahrazení zobrazení REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR Globální dočasné zobrazení bylo nahrazeno po vytvoření datového rámce odkazujícího na něj. Nahrazení se může projevit v existujícím datovém rámci.
Dočasné nahrazení zobrazení REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR Dočasné zobrazení bylo nahrazeno po vytvoření datového rámce odkazujícího na něj. Nahrazení se může projevit v existujícím datovém rámci.
UDF a UDTF mutací OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR Funkce definovaná uživatelem byla po vytvoření datového rámce znovu zaregistrována se stejným názvem. Existující datový rámec může použít novou definici definované uživatelem.
UDF a UDTF mutací OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR UDTF se po vytvoření datového rámce znovu zaregistroval se stejným názvem. Existující datový rámec může použít novou definici UDTF.
UDF a UDTF mutací UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR Funkce definovaná uživatelem odkazuje na globální proměnlivou proměnnou Python. U verze prostředí používá funkce definovaná uživatelem hodnotu proměnné v okamžiku, kdy byla definovaná funkce definovaná uživatelem, nikoli při vyvolání.
UDF a UDTF mutací UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR UDTF odkazuje na globální proměnlivou proměnnou Python. U verze prostředí používá UDTF hodnotu proměnné v okamžiku, kdy byl definovaný UDTF, nikoli při vyvolání.

Databáze a katalogové mutaci

Tyto problémy se vygenerují, když kód kanálu ztlumí výchozí databázi nebo katalog. Ve verzi prostředí můžou datové rámce vytvořené před tím, než se mutací podařilo přeložit tabulky pomocí nové databáze nebo katalogu.

Příklad vzoru, který aktivuje událost:

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()

Bez verze prostředí se df přeloží events z marketing katalogu. S verzí prostředí se df přeloží events z sales katalogu.

Navrhovaná oprava: Plně opravňující názvy tabulek, takže překlad nezávisí na výchozím katalogu nebo databázi, a vyhněte se změně výchozího katalogu nebo databáze mezi vytvořením a použitím datového rámce.

from pyspark import pipelines as dp

df = spark.read.table("marketing.default.events")

@dp.materialized_view
def events_summary():
  return df.groupBy("region").count()

Změny konfigurace Sparku

Tyto problémy se vygenerují, když kód kanálu ztlumí konfiguraci Sparku způsobem, který může změnit chování datového rámce v rámci verze prostředí.

Příklad vzoru, který aktivuje událost:

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")

Bez verze prostředí používá přetypování hodnotu conf při vytváření datového rámce. Ve verzi prostředí se přetypování používá spark.sql.ansi.enabled=true a může selhat při neplatném vstupu.

Navrhovaná oprava: Před vytvořením datového rámce nastavte všechny požadované konfigurace Sparku v horní části souboru kanálu. Pro konfiguraci jednotlivých dotazů použijte nastavení kanálu configuration ve specifikaci kanálu.

Dočasné nahrazení zobrazení

Tyto problémy se vygenerují, když kód kanálu nahradí dočasné zobrazení po vytvoření datového rámce odkazujícího na něj. U verze prostředí může existující datový rámec odrážet nový obsah zobrazení.

Příklad vzoru, který aktivuje událost:

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

Bez verze mytable prostředí obsahuje [(1, "Original Row")]. S verzí mytable prostředí obsahuje [(2, "Replaced Row")].

Navrhovaná oprava: Vytvořte každé dočasné zobrazení jednou a nenahrazovat ho. Pokud potřebujete více zobrazení se souvisejícími daty, zadejte každý jedinečný název.

UDF a UDTF mutací

Tyto problémy se vygenerují, když kód kanálu ztlumí UDF nebo UDTF způsoby, které mění chování v rámci verze prostředí.

Příklad vzoru, který aktivuje událost:

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")))

Bez verze my_mv prostředí obsahuje [("alex_b",)]. S verzí my_mv prostředí obsahuje [("alex_a",)].

Navrhované opravy: Předejte hodnoty do funkce definované uživatelem jako argumenty místo jejich zachytávání z globálních Python nebo nastavte globální před definováním funkce definovanou uživatelem a potom ho neztlumujte.

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")))

Dychtivé provádění ve funkcích toku

Tyto problémy se vygenerují, když kód kanálu provádí dychtivý příkaz Sparku uvnitř funkce zdobené dekorátorem kanálů (@tableatd @materialized_view.). Očekává se, že funkce toku definují a vrací datový rámec; příkazy, které zapisují data, spravují streamované dotazy, registrují prostředky nebo spouští operace ML, nejsou povoleny uvnitř funkce toku se sadou verzí prostředí.

Navrhovaná oprava: Přesuňte dychtivou operaci mimo funkci toku a vraťte datový rámec z funkce toku. Vedlejší účinky, jako je zápis do tabulky nebo spuštění streamovacího dotazu, patří mimo definici kanálu; modul kanálu zpracovává materializaci datového rámce vráceného funkcí toku.

Vyhledání událostí kompatibility v protokolu událostí

Následující dotaz vrátí všechny události kompatibility kanálu seřazené jako poslední:

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;

Počet událostí podle kódu problému v nedávných aktualizacích:

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;

Informace o dotazování protokolu událostí najdete v tématu Dotazování protokolu událostí.

Dodatečné zdroje