Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
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 ...")acreateOrReplaceTempView. - 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
BehaviorChangeInSparkConnectWARNudá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í:
- V editoru kanálů klikněte na Nastavení.
- V nastavení kanálu vyhledejte část Konfigurace .
- Klikněte na
Přidejte konfiguraci.
- Zadejte
pipelines.environmentVersion.enableCompatibilityScanjako klíč atruejako hodnotu. - 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:
- Spusť pipeline v režimu suchého běhu a poté dotazuj na události pipeline event log
BehaviorChangeInSparkConnectWARN. 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 . - Aktualizujte kód pipeline, aby odstranili zjištěné vzory po navržené opravě, a pipeline spusťte znovu.
- 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_typejebehavior_change_in_spark_connect. -
leveljeWARN. -
detailsbehavior_change_in_spark_connectobsahuje objekt, který má jednoissuepole. Hodnota problému je jedním z níže uvedených kódů. -
messageje č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
- Konfigurujte verze prostředí pro pipeline — přehled funkcí, automatickou migraci a jak si sami povolit verzi prostředí.
- Schéma protokolu událostí kanálu – schéma protokolu událostí celého kanálu
- Protokol událostí kanálu – dotazování protokolu událostí kanálu