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 alur Lakeflow berada di Beta.
Alur dengan versi environment mengatur jalankan kode Python melalui Spark Connect. Halaman ini mencakup apa yang tidak kompatibel, apa yang berperilaku berbeda, cara memindai alur untuk pola yang terpengaruh, dan cara memigrasikan alur yang ada.
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 ini sebelumnya dan memblokir pengaktifan hingga ditangani, 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 membantu Anda menemukan pola kode di alur Anda yang akan menghasilkan hasil yang berbeda di bawah versi lingkungan, sebelum Anda mengaktifkannya. Pemindaian adalah keikutsertaan. Saat pemindaian diaktifkan pada alur:
- Setiap eksekusi alur memancarkan satu
BehaviorChangeInSparkConnectWARNperistiwa dalam log peristiwa alur per pola yang terdeteksi. - Anda tidak dapat mengaktifkan versi lingkungan pada alur hingga Anda mengatasi semua peringatan kompatibilitas dari pembaruan yang berhasil sebelumnya.
Jika pemindaian tidak diaktifkan, tidak ada peristiwa yang dikeluarkan dan environment_version pengaktifan tidak diblokir. Databricks merekomendasikan untuk mengaktifkan pemindaian dan menyelesaikan pola apa pun yang terdeteksi sebelum mengaktifkan versi lingkungan pada alur.
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"
}
Alur kerja yang direkomendasikan
- Aktifkan pemindaian pada alur.
- Memicu eksekusi alur.
-
Mengkueri log peristiwa alur untuk
BehaviorChangeInSparkConnectWARNperistiwa. Lihat Referensi peristiwa kompatibilitas untuk daftar lengkap kode masalah, contoh pola, dan perbaikan yang disarankan. - Perbarui kode alur untuk menghapus pola yang terdeteksi dan jalankan alur lagi hingga tidak ada lagi peristiwa yang dikeluarkan.
- Tambahkan
environment_versionke alur menggunakan salah satu metode di Mengaktifkan versi lingkungan pada alur.
Jika Anda yakin peringatan kompatibilitas adalah positif palsu dan tetap ingin mengaktifkannya environment_version , hapus pipelines.environmentVersion.enableCompatibilityScan entri dari konfigurasi alur untuk melewati pemeriksaan. (Mengatur nilai ke false tidak diizinkan — Anda harus menghapus entri sepenuhnya.)
Pemeriksaan preflight tidak berjalan pada alur yang tidak memiliki pembaruan sebelumnya, atau pada alur yang sudah memiliki set versi lingkungan.
Memigrasikan alur yang ada ke versi lingkungan
Untuk memigrasikan alur yang sudah ada yang belum menggunakan versi lingkungan, ikuti alur kerja end-to-end ini. Ini memandu Anda menemukan pola kode yang mungkin berperilaku berbeda di bawah Spark Connect, memperbaikinya, dan meluncurkan versi lingkungan dengan aman.
Aktifkan pemindaian kompatibilitas pada alur. Aktifkan pemindaian pada alur seperti yang dijelaskan dalam Pemindaian kompatibilitas. Inilah yang menyebabkan pola yang terdeteksi muncul di log peristiwa dan apa yang memungkinkan pemeriksaan preflight yang melindungi upaya pengaktifan Anda.
Memicu eksekusi alur dan meninjau peristiwa kompatibilitas. Memicu pembaruan alur normal. Setelah berhasil diselesaikan, kueri log peristiwa alur untuk
BehaviorChangeInSparkConnectWARNperistiwa. Setiap peristiwa melaporkan satu pola yang terdeteksi. Lihat Referensi peristiwa kompatibilitas untuk daftar lengkap kode masalah, contoh pola, dan perbaikan yang disarankan.Perbarui kode alur Anda untuk mengatasi pola yang terdeteksi. Untuk setiap pola yang terdeteksi, perbarui kode alur Anda dengan mengikuti perbaikan yang disarankan. Setelah setiap perubahan, picu pembaruan alur lain dan verifikasi peristiwa yang sesuai tidak lagi muncul. Ulangi hingga log peristiwa tidak lagi menampilkan peristiwa kompatibilitas apa pun untuk pembaruan yang berhasil.
Aktifkan versi lingkungan pada alur. Setelah pembaruan terbaru yang berhasil tidak memiliki peristiwa kompatibilitas, tambahkan
environment_versionke alur menggunakan UI, API, atau bundel seperti yang dijelaskan di Mengaktifkan versi lingkungan pada alur. Pembaruan berikutnya berjalan dengan Spark Connect dan versi bahasa Python yang disematkan dan pustaka yang telah diinstal sebelumnya.Jika pembaruan gagal karena peringatan kompatibilitas masih ada, hilangkan
environment_version, kembali ke langkah 2, dan atasi peringatan yang tersisa sebelum mencoba lagi.Verifikasi migrasi. Setelah pembaruan pertama dengan versi lingkungan selesai, verifikasi:
- Peristiwa
create_updatedalam log peristiwa menunjukkanenvironment_versiondiatur ke nilai yang diharapkan. - Alur menghasilkan data yang diharapkan dan tidak ada peristiwa kesalahan baru yang muncul.
- Tabel hilir pemeriksaan spot untuk setiap perbedaan perilaku halang yang dijelaskan dalam Perubahan perilaku.
- Peristiwa
Rollback
Jika alur salah tingkah setelah migrasi, hapus environment_version dari pengaturan alur. Pembaruan berikutnya berjalan dengan konfigurasi runtime Python sebelumnya. Gunakan eksekusi rolled-back untuk men-debug, lalu ulangi migrasi dari langkah 2 setelah Anda mengidentifikasi dan memperbaiki masalah.
Referensi peristiwa kompatibilitas
Saat pemindaian kompatibilitas diaktifkan pada alur, pemindaian tersebut memancarkan satu BehaviorChangeInSparkConnectWARN peristiwa dalam log peristiwa alur per pola yang terdeteksi. Ketika pemindaian diaktifkan dan pembaruan yang berhasil sebelumnya mendeteksi pola apa pun, alur juga memblokir environment_version pengaktifan hingga pola ditangani.
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 alur — gambaran umum fitur, cara mengaktifkan versi lingkungan.
- Skema log peristiwa alur — skema log peristiwa alur penuh.
- Log peristiwa alur — cara mengkueri log peristiwa alur.