Kompatibilitas versi lingkungan

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 ..."), dan createOrReplaceTempView.
  • 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

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 BehaviorChangeInSparkConnectWARN peristiwa 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:

  1. Dari editor alur, klik Pengaturan.
  2. Temukan bagian Konfigurasi di pengaturan alur.
  3. Klik ikon Plus.Tambahkan konfigurasi.
  4. Masukkan pipelines.environmentVersion.enableCompatibilityScan sebagai kunci dan true sebagai nilai.
  5. Simpan pengaturan alur.

Di alur JSON:

Tambahkan entri berikut ke configuration blok:

"configuration": {
  "pipelines.environmentVersion.enableCompatibilityScan": "true"
}
  1. Aktifkan pemindaian pada alur.
  2. Memicu eksekusi alur.
  3. Mengkueri log peristiwa alur untuk BehaviorChangeInSparkConnectWARN peristiwa. Lihat Referensi peristiwa kompatibilitas untuk daftar lengkap kode masalah, contoh pola, dan perbaikan yang disarankan.
  4. Perbarui kode alur untuk menghapus pola yang terdeteksi dan jalankan alur lagi hingga tidak ada lagi peristiwa yang dikeluarkan.
  5. Tambahkan environment_version ke 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.

  1. 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.

  2. Memicu eksekusi alur dan meninjau peristiwa kompatibilitas. Memicu pembaruan alur normal. Setelah berhasil diselesaikan, kueri log peristiwa alur untuk BehaviorChangeInSparkConnectWARN peristiwa. Setiap peristiwa melaporkan satu pola yang terdeteksi. Lihat Referensi peristiwa kompatibilitas untuk daftar lengkap kode masalah, contoh pola, dan perbaikan yang disarankan.

  3. 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.

  4. Aktifkan versi lingkungan pada alur. Setelah pembaruan terbaru yang berhasil tidak memiliki peristiwa kompatibilitas, tambahkan environment_version ke 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.

  5. Verifikasi migrasi. Setelah pembaruan pertama dengan versi lingkungan selesai, verifikasi:

    • Peristiwa create_update dalam log peristiwa menunjukkan environment_version diatur 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.

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_type adalah behavior_change_in_spark_connect.
  • level adalah WARN.
  • details behavior_change_in_spark_connect berisi objek, yang memiliki satu issue bidang. Nilai masalah adalah salah satu kode yang tercantum di bawah ini.
  • message adalah 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