Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Important
A Lakeflow pipelines környezeti verziói Public Preview (Public Preview) állapotban találhatók.
A környezeti verzióval rendelkező folyamatok Python kódot futtatnak a Spark Connect. Ez az oldal bemutatja, mi az összeegyeztetés, mi viselkedik másként, és hogyan szkennel a Databricks egy pipeline-t az érintett minták után.
Limitations
A környezeti verziók még nem kompatibilisek az összes folyamatfunkcióval. A környezeti verziókészlettel futtatott folyamatok meghiúsulnak, ha a folyamat Python kódja az alábbiak bármelyikét elvégzi:
- A Spark-munkamenet állapotát egy folyamatdekorátorral dekorált függvényen belül mutálja. Ilyenek például a következők:
spark.conf.set(...),spark.sql("USE CATALOG ...")éscreateOrReplaceTempView. - Olyan PySpark API-kat használ, amelyek nem érhetők el a Spark Connectben, beleértve
SparkContextaRDDpy4JSQLContextAPI-kat is. Tekintse meg a Spark Connect által támogatott elemet.
Ha egy környezeti verzió engedélyezése egy folyamaton sikertelenséget okoz, a környezeti verzió letiltása visszaadja a folyamatot az előző állapotának.
Viselkedésbeli változások
A Spark Connectnek kis számú viselkedésbeli eltérése van a klasszikus PySpark-futtatókörnyezettől. A teljes referencia a Spark Connect és a klasszikus Spark között található. A Kompatibilitási vizsgálat előre felismeri ezeket a mintákat, és blokkolja a migrációt, amíg nem kezelik őket, így megtalálhatod és javíthatod őket, mielőtt a termelési adatokat befolyásolnák.
A folyamatokban a leggyakoribb helyzetek, amikor a viselkedés eltérő lehet:
Interleaved DataFrame építés és munkamenet-mutáció
Amikor egy folyamat létrehoz egy DataFrame-et, akkor a Spark-munkamenet állapotát mutálja (például módosítja az alapértelmezett katalógust vagy sémát, beállít egy konfigurációt, lecserél egy ideiglenes nézetet, vagy újra regisztrál egy UDF-et), majd a DataFrame-et használja:
- Környezeti verzió nélkül a DataFrame a mutáció előtti munkamenet állapotát használja.
- A környezeti verzióval a DataFrame a mutáció utáni munkamenet állapotát használja.
Például:
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
A környezet verziója nélkül a mytable következőt tartalmazza [(1, "Original Row")]: . A környezeti verzió tartalmazza a mytable következőt [(2, "Replaced Row")]: .
UDF-ek, amelyek a mutable Python állapotra hivatkoznak
Ha egy UDF egy Python globális változóra hivatkozik, amelynek értéke az UDF definiálása után változik:
- Környezeti verzió nélkül az UDF a változó legújabb értékét használja.
- Egy környezeti verzió esetén az UDF az UDF meghatározásakor használt értéket használja.
Például:
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")))
A környezet verziója nélkül a my_mv következőt tartalmazza [("alex_b",)]: . A környezeti verzió tartalmazza a my_mv következőt [("alex_a",)]: .
Ha egy folyamat bármelyik mintára támaszkodik, a környezeti verzió engedélyezése előtt naplózhatja azt.
Kompatibilitási vizsgálat
A kompatibilitási szkennelés olyan kódmintákat talál a csővezetékben, amelyek környezeti verzióban eltérő eredményeket adnak, így ezeket meg tudod javítani, mielőtt a csővezeték automatikusan migrálna. Ha a vizsgálat engedélyezve van egy folyamatban:
- Minden frissítés egy
BehaviorChangeInSparkConnectWARNeseményt bocsát ki a csővezeték eseménynaplójában minden észlelt mintánként. - A csővezetéket nem migrálják környezeti verzióba, és magad nem engedélyezheted egyet, amíg az előző sikeres frissítésből származó kompatibilitási figyelmeztetések nem teljesülnek.
Ez a pənge nem vonatkozik olyan pipeline-re, amelynek nincs korábbi frissítése, vagy amelynek már van környezeti verziója beállítva.
A folyamat vizsgálatának engedélyezése
A kompatibilitási vizsgálatot a folyamatkonfiguráció hozzáadásával pipelines.environmentVersion.enableCompatibilityScan engedélyezheti. Konfigurációt a folyamatszerkesztő felhasználói felületén vagy egy bejegyzés hozzáadásával adhat hozzá a folyamatkonfiguráció JSON-jához.
A felhasználói felületen keresztül:
- A folyamatszerkesztőben kattintson a Beállítások gombra.
- Keresse meg a Konfiguráció szakaszt a folyamatbeállítások között.
- Kattintson
Konfiguráció hozzáadása.
- Adja meg
pipelines.environmentVersion.enableCompatibilityScankulcsként éstrueértékként. - Mentse a folyamat beállításait.
A folyamat JSON-fájljában:
Adja hozzá a következő bejegyzést a configuration blokkhoz:
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Kompatibilitási figyelmeztetések felülvizsgálata és megoldása
A folyamatokon a környezeti verziót blokkoló minták megtalálásához és törléséhez:
- Futtasd a pipeline-t dry run módban, majd kérdezd le a pipeline event log-ot eseményekről
BehaviorChangeInSparkConnectWARN. Minden esemény egy észlelt mintát jelent. A problémakódok, a példaminták és a javasolt javítások teljes listáját lásd a kompatibilitási eseményekre vonatkozó hivatkozásban . - Frissítsd a csővezeték kódot, hogy eltávolítsuk a javasolt javítást követő észlelt mintákat, majd újra futtatd a csővezetéket.
- Ismételd, amíg egy sikeres frissítés nem küld ki több kompatibilitási eseményt. A csővezetéket automatikusan átlehet migrálni, és magad is engedélyezheted a környezeti verziót.
Egy környezeti verzió engedélyezése ugyanazokat a biztonsági ellenőrzéseket futtatja, hogy a Databricks automatikusan migrálja-e a pipeline-t, vagy te magad állítod be.environment_version Egy folyamat, amelynek kompatibilitási figyelmeztetései megoldatlanul, nem kerülnek környezeti verzióra, amíg a figyelmeztetések meg nem oldódnak. Ha a migráció nem fejezhető be biztonságosan, vagy bármilyen okból meghibásodik, az megáll, mielőtt bármilyen adatot írna, és a csővezeték tovább fut a korábbi futásidőn.
Ha egy frissítés ezen okból megáll, a pipeline eseménynaplója és a frissítési hibaüzenet leírja az okot és a megoldás lépéseit. Kövesd ezeket a lépéseket, és futtasd újra a csővezetéket, hogy befejezd a migrációt. Ha úgy gondolod, hogy egy kompatibilitási figyelmeztetés hamis pozitív, oldd meg a jelölt mintázatot, vagy vedd fel a kapcsolatot az Azure Databricks ügyfélszolgálatával.
Kompatibilitási események referenciája
Amikor a kompatibilitási vizsgálat egy csővezetéken fut, minden észlelt mintázat után egy BehaviorChangeInSparkConnectWARN eseményt bocsát ki a csővezeték eseménynaplójában . Amikor az előző sikeres frissítés bármilyen mintát észlelt, a csővezetéket nem migrálják környezeti verzióba, amíg a mintákat nem kezelik.
Minden esemény egyetlen problémakódot jelent, amely azonosítja az észlelt kódot. A kód kereséséhez keresse meg a Problémakódok táblában – minden sor a példamintát és a javasolt javítást tartalmazó kategóriaszakaszra mutat.
Eseményalakzat
BehaviorChangeInSparkConnect az események a szokásos folyamatesemény-sémát követik:
-
event_typeazbehavior_change_in_spark_connect. -
levelazWARN. -
detailsbehavior_change_in_spark_connectaz objektumot tartalmazza, amely egyetlenissuemezővel rendelkezik. A probléma értéke az alább felsorolt kódok egyike. -
messageaz észlelt minta ember által olvasható leírása.
Problémakódok
| Kategória | Problémakód | Description |
|---|---|---|
| Adatbázis- és katalógusmutációk | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Az alapértelmezett katalógus a DataFrame létrehozása után módosult. A meglévő DataFrame az új alapértelmezett katalógus használatával feloldhatja a táblákat. |
| Adatbázis- és katalógusmutációk | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE CATALOG egy folyamatdekorátor által dekorált függvényen kívül hívták. Az alapértelmezett katalógus váratlanul megváltozhat a későbbi műveletekhez. |
| Adatbázis- és katalógusmutációk | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Az alapértelmezett adatbázis a DataFrame létrehozása után módosult. A meglévő DataFrame az új alapértelmezett adatbázissal oldhatja fel a táblákat. |
| Adatbázis- és katalógusmutációk | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE DATABASE egy folyamatdekorátor által dekorált függvényen kívül hívták. Az alapértelmezett adatbázis váratlanul megváltozhat a későbbi műveletekhez. |
| Lelkes végrehajtás a folyamatfüggvények között | CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény egy ellenőrzőpont-parancsot hív meg. |
| Lelkes végrehajtás a folyamatfüggvények között | CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény lelkesen létrehoz egy DataFrame nézetet (createOrReplaceTempView vagy hasonlót). |
| Lelkes végrehajtás a folyamatfüggvények között | CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény létrehoz egy erőforrásprofilt. |
| Lelkes végrehajtás a folyamatfüggvények között | GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény hívásai spark.resources vagy egy kapcsolódó erőforrás API. |
| Lelkes végrehajtás a folyamatfüggvények között | MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény lelkesen MERGE INTO hajt végre egy céltáblát. |
| Lelkes végrehajtás a folyamatfüggvények között | ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény egy lelkes Spark ML-műveletet hajt végre. |
| Lelkes végrehajtás a folyamatfüggvények között | REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény regisztrál egy Python adatforrást. |
| Lelkes végrehajtás a folyamatfüggvények között | STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény egy aktív streamelési lekérdezési leírón működik. |
| Lelkes végrehajtás a folyamatfüggvények között | STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény regisztrál vagy eltávolít egy streamelési lekérdezés-figyelőt. |
| Lelkes végrehajtás a folyamatfüggvények között | STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény hívásokat indít spark.streams a streamelési lekérdezések kezeléséhez. |
| Lelkes végrehajtás a folyamatfüggvények között | WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény egy lelkes DataFrameWriterV2 műveletet hajt végre. |
| Lelkes végrehajtás a folyamatfüggvények között | WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény egy lelkes DataFrame.write műveletet hajt végre. |
| Lelkes végrehajtás a folyamatfüggvények között | WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
A folyamatfüggvény elindít egy streamelési lekérdezést (writeStream.start()). |
| Spark-konfigurációs mutációk | CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED |
spark.conf.set() vagy spark.conf.unset() egy folyamatdekorátor által dekorált funkcióban hívták meg. Ez környezeti verzió esetén nem támogatott. |
| Spark-konfigurációs mutációk | SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.set() egy folyamatdekorátor által dekorált függvényen kívül lett meghívva egy DataFrame létrehozása után. A konfiguráció módosítása hatással lehet a meglévő DataFrame-re a végrehajtáskor. |
| Spark-konfigurációs mutációk | UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.unset() egy folyamatdekorátor által dekorált függvényen kívül lett meghívva egy DataFrame létrehozása után. A konfiguráció módosítása hatással lehet a meglévő DataFrame-re a végrehajtáskor. |
| Ideiglenes nézetcsere | REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Egy globális ideiglenes nézet a létrehozása után váltott fel egy DataFrame-et, amely hivatkozik rá. A csere a meglévő DataFrame-ben is tükröződhet. |
| Ideiglenes nézetcsere | REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Az ideiglenes nézet egy dataFrame-hivatkozás létrehozása után váltott fel. A csere a meglévő DataFrame-ben is tükröződhet. |
| UDF- és UDTF-mutációk | OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Egy UDF ugyanazzal a névvel lett újra regisztrálva, miután létrejött egy DataFrame-hivatkozás. A meglévő DataFrame használhatja az új UDF-definíciót. |
| UDF- és UDTF-mutációk | OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Az UDTF ugyanazzal a névvel lett regisztrálva, miután létrejött egy DataFrame-hivatkozás. A meglévő DataFrame használhatja az új UDTF-definíciót. |
| UDF- és UDTF-mutációk | UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
A UDF egy globális Python változóra hivatkozik. Egy környezeti verzió esetén az UDF a változó értékét használja az UDF meghatározásakor, nem pedig a meghívási időpontban. |
| UDF- és UDTF-mutációk | UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
Az UDTF egy globális Python változóra hivatkozik. Egy környezeti verzió esetén az UDTF a változó értékét használja az UDTF meghatározásakor, nem a meghívás időpontjában. |
Adatbázis- és katalógusmutációk
Ezek a problémák akkor jelennek meg, ha a folyamatkód az alapértelmezett adatbázist vagy katalógust mutálja. A környezeti verzióval a mutáció előtt létrehozott DataFrame-ek feloldhatják a táblákat az új adatbázis vagy katalógus használatával.
Eseményt kiváltó példaminta:
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()
Környezeti verzió nélkül a df katalógusból oldja events fel a marketing problémát. Egy környezeti verzióval df a katalógusból oldja events fel a sales problémát.
Javasolt javítás: Teljes mértékben megfeleltetheti a táblaneveket, így a megoldás nem függ az alapértelmezett katalógustól vagy adatbázistól, és ne módosítsa az alapértelmezett katalógust vagy adatbázist a DataFrame létrehozása és használata között.
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-konfigurációs mutációk
Ezek a problémák akkor jelentkeznek, ha a folyamatkód a Spark-konfigurációt olyan módon mutálja, amely megváltoztathatja a DataFrame viselkedését egy környezeti verzióban.
Eseményt kiváltó példaminta:
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")
Környezeti verzió nélkül a cast a Conf értéket használja a DataFrame létrehozási idején. A környezeti verzióval a leadott adatok érvénytelen bemenetet használnak spark.sql.ansi.enabled=true , és sikertelenek lehetnek.
Javasolt javítás: A DataFrame létrehozása előtt állítsa be az összes szükséges Spark-konfigurációt a folyamatfájl tetején. Lekérdezésenkénti konfigurációhoz használja a folyamat configuration beállításait a folyamat specifikációjában.
Ideiglenes nézetcsere
Ezek a problémák akkor lépnek fel, ha a folyamatkód egy ideiglenes nézetet cserél le egy dataFrame-hivatkozás létrehozása után. Egy környezeti verzió esetén a meglévő DataFrame tükrözheti az új nézet tartalmát.
Eseményt kiváltó példaminta:
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
A környezet verziója nélkül a mytable következőt tartalmazza [(1, "Original Row")]: . A környezeti verzió tartalmazza a mytable következőt [(2, "Replaced Row")]: .
Javasolt javítás: Az egyes ideiglenes nézeteket egyetlen alkalommal hozza létre, és ne cserélje le. Ha több, kapcsolódó adatokat tartalmazó nézetre van szüksége, adjon meg mindegyiknek külön nevet.
UDF- és UDTF-mutációk
Ezek a problémák akkor jelentkeznek, ha a folyamatkód UDF-et vagy UDTF-et mutál olyan módon, amely megváltoztatja a környezeti verzió viselkedését.
Eseményt kiváltó példaminta:
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")))
A környezet verziója nélkül a my_mv következőt tartalmazza [("alex_b",)]: . A környezeti verzió tartalmazza a my_mv következőt [("alex_a",)]: .
aggested javítás: Értékek továbbítása az UDF-be argumentumként ahelyett, hogy Python globálisakból rögzítené őket, vagy állítsa be a globálist az UDF definiálása előtt, és ne módosítsa azt később.
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")))
Lelkes végrehajtás a folyamatfüggvények között
Ezek a problémák akkor jelennek meg, ha a folyamatkód egy folyamatdekorátor (@tablestb @materialized_view.) által dekorált függvényen belül hajt végre egy lelkes Spark-parancsot. A folyamatfüggvények várhatóan definiálnak és visszaadnak egy DataFrame-et; Az adatok írására, a streamelési lekérdezések kezelésére, az erőforrások regisztrálására vagy az ML-műveletek futtatására szolgáló lelkes parancsok nem engedélyezettek egy környezeti verziókészlettel rendelkező folyamatfüggvényben.
Javasolt javítás: Helyezze át a lelkes műveletet a folyamatfüggvényen kívülre, és adjon vissza egy DataFrame-et a folyamatfüggvényből. Az olyan mellékhatások, mint például a táblázatba írás vagy a streamelési lekérdezés indítása, kívül tartoznak a folyamatdefiníción; a folyamatmotor kezeli a folyamatfüggvény által visszaadott DataFrame materializálását.
Kompatibilitási események keresése az eseménynaplóban
A következő lekérdezés egy folyamat összes kompatibilitási eseményét adja vissza, a legutóbbi sorrendben:
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;
Események számlálása problémakód alapján a legutóbbi frissítések között:
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;
Az eseménynapló lekérdezéséről az eseménynapló lekérdezése című témakörben olvashat.
További források
- Konfiguráld a környezeti verziókat a pipeline-ekhez — funkcióáttekintés, automatikus migráció, és hogyan engedélyezd magad egy környezeti verziót.
- Folyamateseménynapló-séma – teljes folyamatesemény-séma.
- Folyamatesemény-napló – a folyamat eseménynaplójának lekérdezése.