Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Important
Versi lingkungan untuk pipa Lakeflow tersedia dalam Pratinjau Publik.
Alur dengan versi environment mengatur jalankan kode Python melalui Spark Connect. Halaman ini membahas apa yang tidak kompatibel, apa yang berperilaku berbeda, dan bagaimana Databricks memindai pipeline untuk pola yang terpengaruh.
Keterbatasan
Versi lingkungan belum kompatibel dengan semua fungsionalitas alur. Eksekusi alur dengan set versi lingkungan gagal jika kode Python alur melakukan salah satu hal berikut:
- Memutasi status sesi Spark di dalam fungsi yang dihiasi dengan dekorator alur. Contohnya termasuk
spark.conf.set(...),spark.sql("USE CATALOG ..."), dancreateOrReplaceTempView. - Menggunakan API PySpark yang tidak tersedia di Spark Connect, termasuk
SparkContext, ,RDDSQLContext, dan API Py4J apa pun. Lihat Apa yang didukung di Spark Connect.
Jika mengaktifkan versi lingkungan pada alur menyebabkannya gagal, menonaktifkan versi lingkungan mengembalikan alur ke status sebelumnya.
Perubahan perilaku
Spark Connect memiliki sejumlah kecil perbedaan perilaku dari runtime PySpark klasik. Lihat Spark Connect vs. Spark klasik untuk referensi lengkap. Pemindaian Kompatibilitas mendeteksi pola-pola ini sebelumnya dan memblokir migrasi hingga diatasi, sehingga Anda dapat menemukan dan memperbaikinya sebelum memengaruhi data produksi.
Dalam alur, situasi paling umum di mana perilaku mungkin berbeda adalah:
- Konstruksi DataFrame interleaved dan mutasi sesi
- UDF yang mereferensikan status Python yang dapat diubah
Konstruksi DataFrame interleaved dan mutasi sesi
Saat alur membuat DataFrame, lalu mengubah status sesi Spark (misalnya, mengubah katalog atau skema default, mengatur konfigurasi, mengganti tampilan sementara, atau mendaftarkan ulang UDF), lalu menggunakan DataFrame:
- Tanpa versi lingkungan, DataFrame menggunakan status sesi pra-mutasi .
- Dengan versi lingkungan, DataFrame menggunakan status sesi pascamutasi .
Contohnya:
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
Tanpa versi lingkungan, mytable berisi [(1, "Original Row")]. Dengan versi lingkungan, mytable berisi [(2, "Replaced Row")].
UDF yang mereferensikan status Python yang dapat diubah
Ketika UDF mereferensikan variabel global Python yang nilainya berubah setelah UDF ditentukan:
- Tanpa versi lingkungan, UDF menggunakan nilai terbaru variabel.
- Dengan versi lingkungan, UDF menggunakan nilai pada saat UDF ditentukan.
Contohnya:
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")))
Tanpa versi lingkungan, my_mv berisi [("alex_b",)]. Dengan versi lingkungan, my_mv berisi [("alex_a",)].
Jika alur bergantung pada salah satu pola, audit sebelum mengaktifkan versi lingkungan.
Pemindaian kompatibilitas
Pemindaian kompatibilitas menemukan pola kode dalam pipeline Anda yang akan menghasilkan hasil berbeda di bawah versi lingkungan, sehingga Anda dapat memperbaikinya sebelum pipeline dimigrasikan secara otomatis. Saat pemindaian diaktifkan pada alur:
- Setiap pembaruan mengeluarkan satu
BehaviorChangeInSparkConnectWARNevent dalam log event pipeline per pola yang terdeteksi. - Pipeline tidak dimigrasikan ke versi lingkungan, dan Anda tidak dapat mengaktifkannya sendiri sampai semua peringatan kompatibilitas dari pembaruan sebelumnya yang berhasil diatasi.
Pemeriksaan ini tidak berlaku untuk pipeline yang belum memiliki pembaruan sebelumnya, atau yang sudah memiliki versi lingkungan yang sudah diatur.
Mengaktifkan pemindaian pada alur
Anda dapat mengaktifkan pemindaian kompatibilitas dengan menambahkan pipelines.environmentVersion.enableCompatibilityScan konfigurasi alur. Anda dapat menambahkan konfigurasi melalui antarmuka pengguna editor alur atau dengan menambahkan entri ke konfigurasi alur JSON.
Melalui UI:
- Dari editor alur, klik Pengaturan.
- Temukan bagian Konfigurasi di pengaturan alur.
- Klik
Tambahkan konfigurasi.
- Masukkan
pipelines.environmentVersion.enableCompatibilityScansebagai kunci dantruesebagai nilai. - Simpan pengaturan alur.
Di alur JSON:
Tambahkan entri berikut ke configuration blok:
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Tinjau dan selesaikan peringatan kompatibilitas
Untuk menemukan dan membersihkan pola yang memblokir versi lingkungan pada pipeline Anda:
- Jalankan pipeline dalam mode dry run , lalu kueri log event pipeline untuk event
BehaviorChangeInSparkConnectWARN. Setiap peristiwa melaporkan satu pola yang terdeteksi. Lihat Referensi peristiwa kompatibilitas untuk daftar lengkap kode masalah, contoh pola, dan perbaikan yang disarankan. - Perbarui kode pipeline untuk menghapus pola yang terdeteksi setelah perbaikan yang disarankan, dan jalankan pipeline lagi.
- Ulangi sampai pembaruan yang berhasil tidak lagi menghasilkan acara kompatibilitas. Pipeline kemudian dapat dimigrasikan secara otomatis, dan Anda juga dapat mengaktifkan versi lingkungan sendiri.
Mengaktifkan versi lingkungan menjalankan pemeriksaan keamanan yang sama apakah Databricks memigrasikan pipeline secara otomatis atau Anda yang mengatur environment_version sendiri. Pipeline dengan peringatan kompatibilitas yang belum terselesaikan tidak akan berpindah ke versi lingkungan sampai peringatan tersebut diselesaikan. Jika migrasi tidak dapat diselesaikan dengan aman, atau gagal karena alasan apapun, migrasi berhenti sebelum menulis data apa pun dan pipeline terus berjalan pada runtime sebelumnya.
Ketika pembaruan berhenti karena salah satu alasan ini, log peristiwa pipeline dan pesan kesalahan pembaruan menjelaskan penyebab dan langkah-langkah untuk menyelesaikannya. Ikuti langkah-langkah tersebut dan jalankan pipeline lagi untuk menyelesaikan migrasi. Jika Anda percaya peringatan kompatibilitas adalah positif palsu, selesaikan pola yang ditandai atau hubungi dukungan Azure Databricks.
Referensi peristiwa kompatibilitas
Ketika pemindaian kompatibilitas berjalan pada pipeline, ia memancarkan satu BehaviorChangeInSparkConnectWARN event dalam log event pipeline per pola yang terdeteksi. Ketika pembaruan sebelumnya yang berhasil mendeteksi pola, pipeline tidak dimigrasikan ke versi lingkungan sampai pola-pola tersebut diatasi.
Setiap peristiwa melaporkan satu kode masalah yang mengidentifikasi apa yang terdeteksi. Untuk mencari kode, temukan di tabel Kode masalah — setiap baris ditautkan ke bagian kategori yang berisi pola contoh dan perbaikan yang disarankan .
Bentuk peristiwa
BehaviorChangeInSparkConnect peristiwa mengikuti skema log peristiwa alur standar:
-
event_typeadalahbehavior_change_in_spark_connect. -
leveladalahWARN. -
detailsbehavior_change_in_spark_connectberisi objek, yang memiliki satuissuebidang. Nilai masalah adalah salah satu kode yang tercantum di bawah ini. -
messageadalah deskripsi yang dapat dibaca manusia dari pola yang terdeteksi.
Kode masalah
| Kategori | Kode masalah | Description |
|---|---|---|
| Mutasi database dan katalog | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Katalog default diubah setelah DataFrame dibuat. DataFrame yang ada dapat mengatasi tabel menggunakan katalog default baru. |
| Mutasi database dan katalog | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE CATALOG dipanggil di luar fungsi yang dihiasi oleh dekorator alur. Katalog default dapat berubah secara tidak terduga untuk operasi berikutnya. |
| Mutasi database dan katalog | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Database default diubah setelah DataFrame dibuat. DataFrame yang ada dapat mengatasi tabel menggunakan database default baru. |
| Mutasi database dan katalog | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE DATABASE dipanggil di luar fungsi yang dihiasi oleh dekorator alur. Database default dapat berubah secara tidak terduga untuk operasi berikutnya. |
| Eksekusi bersemangat dalam fungsi alur | CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur memanggil perintah titik pemeriksaan. |
| Eksekusi bersemangat dalam fungsi alur | CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur dengan bersemangat membuat tampilan DataFrame (createOrReplaceTempView atau serupa). |
| Eksekusi bersemangat dalam fungsi alur | CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur membuat profil sumber daya. |
| Eksekusi bersemangat dalam fungsi alur | GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur memanggil spark.resources atau API sumber daya terkait. |
| Eksekusi bersemangat dalam fungsi alur | MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur melakukan bersemangat MERGE INTO pada tabel target. |
| Eksekusi bersemangat dalam fungsi alur | ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur melakukan operasi Spark ML yang bersemangat. |
| Eksekusi bersemangat dalam fungsi alur | REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur mendaftarkan sumber data Python. |
| Eksekusi bersemangat dalam fungsi alur | STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur beroperasi pada handel kueri streaming aktif. |
| Eksekusi bersemangat dalam fungsi alur | STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur mendaftarkan atau menghapus pendengar kueri streaming. |
| Eksekusi bersemangat dalam fungsi alur | STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur memanggil spark.streams untuk mengelola kueri streaming. |
| Eksekusi bersemangat dalam fungsi alur | WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur melakukan operasi yang bersemangat DataFrameWriterV2 . |
| Eksekusi bersemangat dalam fungsi alur | WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur melakukan operasi yang bersemangat DataFrame.write . |
| Eksekusi bersemangat dalam fungsi alur | WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Fungsi alur memulai kueri streaming (writeStream.start()). |
| Mutasi konfigurasi Spark | CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED |
spark.conf.set() atau spark.conf.unset() dipanggil di dalam fungsi yang dihiasi oleh dekorator alur. Ini tidak didukung dengan versi lingkungan. |
| Mutasi konfigurasi Spark | SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.set() dipanggil di luar fungsi yang dihiasi oleh dekorator alur setelah DataFrame dibuat. Perubahan konfigurasi dapat memengaruhi DataFrame yang ada pada waktu eksekusi. |
| Mutasi konfigurasi Spark | UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.unset() dipanggil di luar fungsi yang dihiasi oleh dekorator alur setelah DataFrame dibuat. Perubahan konfigurasi dapat memengaruhi DataFrame yang ada pada waktu eksekusi. |
| Penggantian tampilan sementara | REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Tampilan sementara global diganti setelah DataFrame mereferensikannya dibuat. Penggantian dapat tercermin dalam DataFrame yang ada. |
| Penggantian tampilan sementara | REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Tampilan sementara diganti setelah DataFrame mereferensikannya dibuat. Penggantian dapat tercermin dalam DataFrame yang ada. |
| Mutasi UDF dan UDTF | OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
UDF didaftarkan ulang dengan nama yang sama setelah DataFrame yang merujuknya dibuat. DataFrame yang ada dapat menggunakan definisi UDF baru. |
| Mutasi UDF dan UDTF | OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
UDTF didaftarkan ulang dengan nama yang sama setelah DataFrame yang mereferensikannya dibuat. DataFrame yang ada dapat menggunakan definisi UDTF baru. |
| Mutasi UDF dan UDTF | UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDF mereferensikan variabel Python global yang dapat diubah. Dengan versi lingkungan, UDF menggunakan nilai variabel pada saat UDF ditentukan, bukan pada waktu pemanggilan. |
| Mutasi UDF dan UDTF | UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDTF mereferensikan variabel Python global yang dapat diubah. Dengan versi lingkungan, UDTF menggunakan nilai variabel pada saat UDTF ditentukan, bukan pada waktu pemanggilan. |
Mutasi database dan katalog
Masalah ini dipancarkan saat kode alur memutasi database atau katalog default. Dengan versi lingkungan, DataFrames yang dibangun sebelum mutasi dapat menyelesaikan tabel menggunakan database atau katalog baru.
Contoh pola yang memicu peristiwa:
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()
Tanpa versi lingkungan, df diselesaikan events dari marketing katalog. Dengan versi lingkungan, df diselesaikan events dari sales katalog.
Perbaikan yang disarankan: Sepenuhnya memenuhi syarat nama tabel sehingga resolusi tidak bergantung pada katalog atau database default, dan hindari mengubah katalog atau database default antara pembuatan dan penggunaan DataFrame.
from pyspark import pipelines as dp
df = spark.read.table("marketing.default.events")
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Mutasi konfigurasi Spark
Masalah ini dipancarkan ketika kode alur mengubah konfigurasi Spark dengan cara yang dapat mengubah perilaku DataFrame di bawah versi lingkungan.
Contoh pola yang memicu peristiwa:
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")
Tanpa versi lingkungan, pemeran menggunakan nilai conf pada waktu pembuatan DataFrame. Dengan versi lingkungan, pemeran spark.sql.ansi.enabled=true menggunakan dan mungkin gagal pada input yang tidak valid.
Perbaikan yang disarankan: Atur semua konfigurasi Spark yang diperlukan di bagian atas file alur, sebelum DataFrame dibuat. Untuk konfigurasi per kueri, gunakan pengaturan alur configuration dalam spesifikasi alur.
Penggantian tampilan sementara
Masalah ini dipancarkan ketika kode alur menggantikan tampilan sementara setelah DataFrame yang merujuknya dibuat. Dengan versi lingkungan, DataFrame yang ada dapat mencerminkan konten tampilan baru.
Contoh pola yang memicu peristiwa:
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
Tanpa versi lingkungan, mytable berisi [(1, "Original Row")]. Dengan versi lingkungan, mytable berisi [(2, "Replaced Row")].
Perbaikan yang disarankan: Buat setiap tampilan sementara satu kali dan jangan ganti. Jika Anda memerlukan beberapa tampilan dengan data terkait, beri nama yang berbeda.
Mutasi UDF dan UDTF
Masalah ini dipancarkan ketika kode alur mengubah UDF atau UDTF dengan cara yang mengubah perilaku di bawah versi lingkungan.
Contoh pola yang memicu peristiwa:
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")))
Tanpa versi lingkungan, my_mv berisi [("alex_b",)]. Dengan versi lingkungan, my_mv berisi [("alex_a",)].
Perbaiki: Meneruskan nilai ke dalam UDF sebagai argumen alih-alih menangkapnya dari Python global, atau atur global sebelum menentukan UDF dan jangan bermutasi setelahnya.
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")))
Eksekusi bersemangat dalam fungsi alur
Masalah ini dipancarkan ketika kode alur melakukan perintah Spark yang bersemangat di dalam fungsi yang dihiasi oleh dekorator alur (@table, @materialized_view, dll.). Fungsi alur diharapkan untuk menentukan dan mengembalikan DataFrame; Perintah bersemangat yang menulis data, mengelola kueri streaming, mendaftarkan sumber daya, atau menjalankan operasi ML tidak diizinkan di dalam fungsi alur dengan set versi lingkungan.
Perbaikan yang disarankan: Pindahkan operasi bersemangat di luar fungsi alur dan kembalikan DataFrame dari fungsi alur sebagai gantinya. Efek samping seperti menulis ke tabel atau memulai kueri streaming berada di luar definisi alur; mesin alur menangani materialisasi DataFrame yang dikembalikan oleh fungsi alur.
Menemukan peristiwa kompatibilitas di log peristiwa
Kueri berikut mengembalikan semua peristiwa kompatibilitas untuk alur, diurutkan terbaru terlebih dahulu:
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;
Untuk menghitung peristiwa berdasarkan kode masalah di seluruh pembaruan terbaru:
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;
Untuk cara mengkueri log peristiwa, lihat Mengkueri log peristiwa.
Sumber daya tambahan
- Konfigurasikan versi lingkungan untuk pipeline — gambaran fitur, migrasi otomatis, dan cara mengaktifkan versi lingkungan sendiri.
- Skema log peristiwa alur — skema log peristiwa alur penuh.
- Log peristiwa alur — cara mengkueri log peristiwa alur.