A környezet verziókompatibilitása

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 ...")és createOrReplaceTempView.
  • Olyan PySpark API-kat használ, amelyek nem érhetők el a Spark Connectben, beleértve SparkContexta RDDpy4J SQLContextAPI-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 BehaviorChangeInSparkConnectWARN esemé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:

  1. A folyamatszerkesztőben kattintson a Beállítások gombra.
  2. Keresse meg a Konfiguráció szakaszt a folyamatbeállítások között.
  3. Kattintson a Plusz ikonra.Konfiguráció hozzáadása.
  4. Adja meg pipelines.environmentVersion.enableCompatibilityScan kulcsként és true értékként.
  5. 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:

  1. 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 .
  2. 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.
  3. 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_type az behavior_change_in_spark_connect.
  • level az WARN.
  • details behavior_change_in_spark_connect az objektumot tartalmazza, amely egyetlen issue mezővel rendelkezik. A probléma értéke az alább felsorolt kódok egyike.
  • message az é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