Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Important
Lakeflow pipelines için ortam sürümleri Public Preview bölümündedir.
environment version kümesine sahip işlem hatları Python kodu Spark Connect aracılığıyla çalıştırır. Bu sayfa, neyin uyumsuz olduğunu, neyin farklı davrandığını ve Databricks'in etkilenen kalıpları bulmak için bir boru hattını nasıl taradığını ele alır.
Sınırlamalar
Ortam sürümleri henüz tüm işlem hattı işlevleriyle uyumlu değildir. İşlem hattının Python kodu aşağıdakilerden birini yaparsa ortam sürümü kümesiyle işlem hattı çalıştırması başarısız olur:
- İşlem hatları dekoratörüyle dekore edilmiş bir işlevin içinde Spark oturum durumunu sessize alır. Örnek olarak
spark.conf.set(...),spark.sql("USE CATALOG ...")vecreateOrReplaceTempViewverilebilir. - ,
SparkContext,RDDve tüm Py4J API'leri dahil olmak üzereSQLContextSpark Connect'te kullanılamayan PySpark API'lerini kullanır. Bkz . Spark Connect'te desteklenenler.
bir işlem hattında ortam sürümünün etkinleştirilmesi başarısız olmasına neden oluyorsa, ortam sürümünün devre dışı bırakılması işlem hattını önceki durumuna döndürür.
Davranış değişiklikleri
Spark Connect'in klasik PySpark çalışma zamanından az sayıda davranış farkı vardır. Tam başvuru için bkz . Spark Connect ve klasik Spark . Uyumluluk taraması bu kalıpları önceden tespit eder ve taşınmayı engeller, böylece üretim verilerini etkilemeden önce bulup düzeltebilirsiniz.
İşlem hattında, davranışın farklı olabileceği en yaygın durumlar şunlardır:
- Interleaved DataFrame yapısı ve oturum mutasyonu
- Değiştirilebilir Python durumuna başvuran UDFs
Interleaved DataFrame yapısı ve oturum mutasyonu
İşlem hattı bir DataFrame oluştururken Spark oturum durumunu sessize alır (örneğin, varsayılan kataloğu veya şemayı değiştirir, yapılandırma ayarlar, geçici görünümü değiştirir veya bir UDF'yi yeniden kaydeder), ardından DataFrame'i kullanır:
- Ortam sürümü olmadan DataFrame , mutasyon öncesi oturum durumunu kullanır.
- Ortam sürümüyle DataFrame , mutasyon sonrası oturum durumunu kullanır.
Örneğin:
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
Ortam sürümü olmadan içerir mytable[(1, "Original Row")]. Ortam sürümü ile içerir mytable[(2, "Replaced Row")].
Değiştirilebilir Python durumuna başvuran UDF'ler
UDF, UDF tanımlandıktan sonra değeri değişen bir Python genel değişkenine başvurduğunda:
- Ortam sürümü olmadan UDF, değişkenin en son değerini kullanır.
- Ortam sürümüyle UDF, UDF tanımlandığı sırada değerini kullanır.
Örneğin:
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")))
Ortam sürümü olmadan içerir my_mv[("alex_b",)]. Ortam sürümü ile içerir my_mv[("alex_a",)].
İşlem hattı iki desenden birini kullanıyorsa, ortam sürümünü etkinleştirmeden önce bunu denetleyin.
Uyumluluk taraması
Uyumluluk taraması, boru hattınızda bir ortam sürümünde farklı sonuçlar üretecek kod kalıplarını bulur, böylece bir boru hattı otomatik olarak taşınmadan önce bunları düzeltebilirsiniz. İşlem hattında tarama etkinleştirildiğinde:
- Her güncelleme, tespit edilen desen başına pipeline olay günlüğünde bir
BehaviorChangeInSparkConnectWARNolay yayımlar. - Boru hattı bir ortam sürümüne taşınmaz ve önceki başarılı güncellemeden gelen tüm uyumluluk uyarıları ele alınana kadar kendiniz etkinleştiremezsiniz.
Bu kontrol, önceki güncellemesi olmayan veya zaten bir ortam sürüm seti olan bir pipeline için geçerli değildir.
İşlem hattında taramayı etkinleştirme
İşlem hattı yapılandırmasını ekleyerek uyumluluk taramasını pipelines.environmentVersion.enableCompatibilityScan etkinleştirebilirsiniz. İşlem hattı düzenleyicisi kullanıcı arabirimi aracılığıyla veya işlem hattı yapılandırması JSON'sine bir giriş ekleyerek yapılandırma ekleyebilirsiniz.
Kullanıcı arabirimi aracılığıyla:
- İşlem hattı düzenleyicisinde Ayarlar'a tıklayın.
- İşlem hattı ayarlarında Yapılandırma bölümünü bulun.
-
Yapılandırma ekleyin.
- Anahtar olarak ve
pipelines.environmentVersion.enableCompatibilityScandeğer olarak girintrue. - İşlem hattı ayarlarını kaydedin.
İşlem hattı JSON'unda:
Aşağıdaki girdiyi bloğuna configuration ekleyin:
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Uyumluluk uyarılarını gözden geçirin ve çözün
Bir ortam sürümünü pipeline'ınızda engelleyen kalıpları bulmak ve temizlemek için:
- Pipeline'ı kuru çalıştırma modunda çalıştırın, ardından pipeline olay kaydında olaylar
BehaviorChangeInSparkConnectWARNiçin sorgulayın. Her olay algılanan bir desen bildirir. Sorun kodlarının, örnek desenlerin ve önerilen düzeltmelerin tam listesi için uyumluluk olayları başvurusuna bakın. - Önerilen düzeltmeden sonra tespit edilen kalıpları kaldırmak için boru hattı kodunu güncelleyin ve boru hattını tekrar çalıştırın.
- Başarılı bir güncelleme artık uyumluluk olayı vermeyene kadar tekrarlayın. Pipeline otomatik olarak taşınabilir ve ayrıca bir ortam versiyonunu kendiniz etkinleştirebilirsiniz.
Bir ortam sürümünü etkinleştirmek aynı güvenliği çalıştırır, Databricks'in pipeline'ı otomatik olarak mı yoksa kendiniz ayarladığını environment_version mı kontrol eder. Çözülmemiş uyumluluk uyarılarına sahip bir boru hattı, uyarılar çözülene kadar bir ortam sürümüne geçmez. Eğer geçiş güvenli bir şekilde tamamlanamazsa veya herhangi bir nedenle başarısız olursa, veri yazmadan durur ve pipeline önceki çalışma zamanında çalışmaya devam eder.
Bu nedenlerden biri nedeniyle bir güncelleme durduğunda, pipeline olay kaydı ve güncelleme hata mesajı nedeni ve çözüm adımlarını açıklar. Bu adımları takip edin ve geçişi tamamlamak için pipeline'ı tekrar çalıştırın. Uyumluluk uyarısının yanlış pozitif olduğunu düşünüyorsanız, işaretlenen deseni çözün veya Azure Databricks destek ile iletişime geçin.
Uyumluluk olayları başvurusu
Uyumluluk taraması bir boru hattı üzerinde çalıştığında, tespit edilen desen için boru hattı olay günlüğünde bir BehaviorChangeInSparkConnectWARN olay yayımlar. Önceki başarılı güncelleme herhangi bir desen tespit ettiğinde, pipeline desenler ele alınana kadar bir ortam sürümüne taşınmaz.
Her olay, algılananları tanımlayan tek bir sorun kodu bildirir. Bir kodu aramak için Sorun kodları tablosunda bulabilirsiniz; her satır örnek desen ve önerilen düzeltmeyi içeren kategori bölümüne bağlanır.
Olay şekli
BehaviorChangeInSparkConnect olaylar standart işlem hattı olay günlüğü şemasını izler:
-
event_type,behavior_change_in_spark_connect'e eşittir. -
level,WARN'e eşittir. -
detailsbehavior_change_in_spark_connecttek birissuealanı olan nesnesini içerir. Sorun değeri, aşağıda listelenen kodlardan biridir. -
messagealgılanan desenin okunabilir bir açıklamasıdır.
Sorun kodları
| Kategori | Sorun kodu | Description |
|---|---|---|
| Veritabanı ve katalog mutasyonları | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
DataFrame oluşturulduktan sonra varsayılan katalog değiştirildi. Mevcut DataFrame, yeni varsayılan kataloğu kullanarak tabloları çözebilir. |
| Veritabanı ve katalog mutasyonları | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE CATALOG bir işlem hattı dekoratörü tarafından dekore edilmiş bir işlevin dışında çağrıldı. Sonraki işlemler için varsayılan katalog beklenmedik şekilde değişebilir. |
| Veritabanı ve katalog mutasyonları | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Bir DataFrame oluşturulduktan sonra varsayılan veritabanı değiştirildi. Mevcut DataFrame, yeni varsayılan veritabanını kullanarak tabloları çözümleyebilmektedir. |
| Veritabanı ve katalog mutasyonları | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE DATABASE bir işlem hattı dekoratörü tarafından dekore edilmiş bir işlevin dışında çağrıldı. Sonraki işlemler için varsayılan veritabanı beklenmedik şekilde değişebilir. |
| Akış işlevleri içinde istekli yürütme | CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Flow işlevi bir denetim noktası komutu çağırır. |
| Akış işlevleri içinde istekli yürütme | CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi hevesle bir DataFrame görünümü (createOrReplaceTempView veya benzeri) oluşturur. |
| Akış işlevleri içinde istekli yürütme | CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi bir kaynak profili oluşturur. |
| Akış işlevleri içinde istekli yürütme | GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi çağrıları veya ilgili bir kaynak API'sini çağırır spark.resources . |
| Akış işlevleri içinde istekli yürütme | MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi, hedef tabloda istekli MERGE INTO bir performans gösterir. |
| Akış işlevleri içinde istekli yürütme | ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi, istekli bir Spark ML işlemi gerçekleştirir. |
| Akış işlevleri içinde istekli yürütme | REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Flow işlevi bir Python veri kaynağı kaydeder. |
| Akış işlevleri içinde istekli yürütme | STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi etkin bir akış sorgu tutamacı üzerinde çalışır. |
| Akış işlevleri içinde istekli yürütme | STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi bir akış sorgusu dinleyicisini kaydeder veya kaldırır. |
| Akış işlevleri içinde istekli yürütme | STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi akış sorgularını yönetmek için çağırır spark.streams . |
| Akış işlevleri içinde istekli yürütme | WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi, istekli DataFrameWriterV2 bir işlem gerçekleştirir. |
| Akış işlevleri içinde istekli yürütme | WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi, istekli DataFrame.write bir işlem gerçekleştirir. |
| Akış işlevleri içinde istekli yürütme | WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Akış işlevi bir akış sorgusu (writeStream.start()) başlatır. |
| Spark yapılandırma mutasyonları | CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED |
spark.conf.set() veya spark.conf.unset() bir işlem hattı dekoratörü tarafından dekore edilmiş bir işlevin içinde çağrıldı. Bu, ortam sürümüyle desteklenmez. |
| Spark yapılandırma mutasyonları | SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.set() bir DataFrame oluşturulduktan sonra işlem hattı dekoratörü tarafından dekore edilmiş bir işlevin dışında çağrıldı. Yapılandırma değişikliği, yürütme sırasında mevcut DataFrame'i etkileyebilir. |
| Spark yapılandırma mutasyonları | UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.unset() bir DataFrame oluşturulduktan sonra işlem hattı dekoratörü tarafından dekore edilmiş bir işlevin dışında çağrıldı. Yapılandırma değişikliği, yürütme sırasında mevcut DataFrame'i etkileyebilir. |
| Geçici görünüm değiştirmeleri | REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
DataFrame'e başvuran bir genel geçici görünüm oluşturulduktan sonra değiştirildi. Değiştirme, mevcut DataFrame'e yansıtılabilir. |
| Geçici görünüm değiştirmeleri | REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
DataFrame'e başvuruda bulunan bir geçici görünüm oluşturulduktan sonra değiştirildi. Değiştirme, mevcut DataFrame'e yansıtılabilir. |
| UDF ve UDTF mutasyonları | OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Bir UDF, oluşturulduğuna başvuran bir DataFrame'den sonra aynı adla yeniden kaydedildi. Mevcut DataFrame yeni UDF tanımını kullanabilir. |
| UDF ve UDTF mutasyonları | OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
UDTF, oluşturulduğuna başvuran bir DataFrame'in ardından aynı adla yeniden kaydedildi. Mevcut DataFrame yeni UDTF tanımını kullanabilir. |
| UDF ve UDTF mutasyonları | UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDF, genel bir değiştirilebilir Python değişkenine başvurur. Ortam sürümüyle UDF, çağırma zamanında değil, UDF tanımlandığı sırada değişkenin değerini kullanır. |
| UDF ve UDTF mutasyonları | UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDTF, genel bir değiştirilebilir Python değişkenine başvurur. Ortam sürümüyle UDTF, çağırma zamanında değil, UDTF tanımlandığı sırada değişkenin değerini kullanır. |
Veritabanı ve katalog mutasyonları
İşlem hattı kodu varsayılan veritabanını veya kataloğu sessize aldığınızda bu sorunlar yayılır. Ortam sürümüyle, mutasyondan önce inşa edilen DataFrame'ler yeni veritabanını veya kataloğu kullanarak tabloları çözebilir.
Bir olayı tetikleyen örnek desen:
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()
Ortam sürümü olmadan, df katalogdan events çözümlermarketing. Ortam sürümüyle, df katalogdan events çözümlersales.
Önerilen düzeltme: Çözümlemenin varsayılan kataloğa veya veritabanına bağımlı olmaması için tablo adlarını tam olarak niteleyin ve DataFrame oluşturma ve kullanma arasında varsayılan kataloğu veya veritabanını değiştirmekten kaçının.
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 yapılandırma mutasyonları
İşlem hattı kodu Spark yapılandırmasını ortam sürümü altında DataFrame davranışını değiştirebilecek şekilde sessize aldığınızda bu sorunlar yayılır.
Bir olayı tetikleyen örnek desen:
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")
Ortam sürümü olmadan, atama DataFrame oluşturma zamanında conf değerini kullanır. Ortam sürümüyle, atama işlemi kullanır spark.sql.ansi.enabled=true ve geçersiz girişte başarısız olabilir.
Önerilen düzeltme: Herhangi bir DataFrame oluşturulmadan önce işlem hattı dosyasının en üstünde tüm gerekli Spark yapılandırmalarını ayarlayın. Sorgu başına yapılandırma için işlem hattı belirtimindeki configuration işlem hattının ayarını kullanın.
Geçici görünüm değiştirmeleri
Bu sorunlar, işlem hattı kodu oluşturulduktan sonra dataframe'e başvuran geçici bir görünümün yerini alırsa ortaya çıkar. Ortam sürümüyle, mevcut DataFrame yeni görünüm içeriğini yansıtabilir.
Bir olayı tetikleyen örnek desen:
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
Ortam sürümü olmadan içerir mytable[(1, "Original Row")]. Ortam sürümü ile içerir mytable[(2, "Replaced Row")].
Önerilen düzeltme: Her geçici görünümü tek bir kez oluşturun ve değiştirmeyin. İlgili verilerle birden çok görünüme ihtiyacınız varsa, her birine ayrı bir ad verin.
UDF ve UDTF mutasyonları
bu sorunlar, işlem hattı kodu bir UDF veya UDTF'yi ortam sürümü altında davranışı değiştirecek şekilde sessize alırsa yayılır.
Bir olayı tetikleyen örnek desen:
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")))
Ortam sürümü olmadan içerir my_mv[("alex_b",)]. Ortam sürümü ile içerir my_mv[("alex_a",)].
Suggested fix: Değerleri Python genel değerlerden yakalamak yerine UDF'ye bağımsız değişken olarak geçirin veya UDF'yi tanımlamadan önce geneli ayarlayın ve daha sonra sessize alma.
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")))
Akış işlevleri içinde istekli yürütme
İşlem hattı kodu, bir işlem hattı dekoratörü (@table, vb.) tarafından dekore edilmiş bir işlev içinde istekli bir Spark komutu gerçekleştirdiğinde bu @materialized_viewsorunlar yayılır. Akış işlevlerinin bir DataFrame tanımlaması ve döndürmesi beklenir; ortam sürümü ayarlanmış bir akış işlevi içinde veri yazan, akış sorgularını yöneten, kaynakları kaydeden veya ML işlemlerini çalıştıran istekli komutlara izin verilmez.
Önerilen düzeltme: İstekli işlemi akış işlevinin dışına taşıyın ve bunun yerine akış işlevinden bir DataFrame döndürür. Tabloya yazma veya akış sorgusu başlatma gibi yan etkiler işlem hattı tanımının dışındadır; işlem hattı altyapısı, akış işlevi tarafından döndürülen DataFrame'in gerçekleştirilmesini işler.
Olay günlüğünde uyumluluk olaylarını bulma
Aşağıdaki sorgu, bir işlem hattı için en son sıralanmış tüm uyumluluk olaylarını döndürür:
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;
Son güncelleştirmelerde sorun koduna göre olayları saymak için:
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;
Olay günlüğünü sorgulamak için bkz. Olay günlüğünü sorgulama.
Ek kaynaklar
- Pipeline'lar için ortam sürümlerini yapılandırın — özellikler genel büzeri, otomatik göç ve bir ortam sürümünü kendiniz nasıl etkinleştireceğiniz.
- İşlem hattı olay günlüğü şeması — tam işlem hattı olay günlüğü şeması.
- İşlem hattı olay günlüğü — işlem hattı olay günlüğünü sorgulama.