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.
Mesin eksekusi asli di Microsoft Fabric sekarang mendukung Python fungsi yang ditentukan pengguna (UDF), UDF Scala, dan jenis data kompleks (array, peta, dan struct). Kemampuan ini memungkinkan Anda menulis aplikasi Spark ekspresif tanpa mengorbankan performa.
dukungan UDF Python
Python adalah salah satu bahasa yang paling populer dalam rekayasa data dan ilmu data. Secara historis, Python UDF menyebabkan overhead yang signifikan pada Spark akibat biaya serialisasi antara JVM dan proses pekerja Python. Mesin eksekusi asli meminimalkan transisi mahal ini, memungkinkan eksekusi yang lebih cepat tanpa perubahan kode.
Cara kerja UDF Python di mesin eksekusi asli
Dalam model eksekusi Spark konvensional, eksekusi UDF Python melibatkan:
- Konversi data dari format internal Spark.
- Serialisasi dan pemindahan ke proses pekerja Python.
- Python eksekusi UDF.
- Serialisasi hasil kembali ke JVM.
- Spark melanjutkan eksekusi.
Gerakan lintas runtime ini menciptakan biaya serialisasi/deserialisasi, inefisiensi CPU, dan alur eksekusi kolom rusak. Mesin eksekusi asli mengurangi overhead ini dengan mengoptimalkan jalur transfer data dan mempertahankan pemrosesan vektorisasi jika memungkinkan.
Jenis UDF Python yang didukung
Mesin eksekusi asli mendukung:
-
UDF skalar: Fungsi Python per baris yang didaftarkan dengan
udf(). -
UDF Tervektorisasi (Pandas): Fungsi yang dianotasi dengan
@pandas_udfyang beroperasi pada batch data menggunakan Apache Arrow untuk transfer data yang efisien.
UDF yang divektorisasi memberikan peningkatan kinerja terbesar karena secara alami selaras dengan model pemrosesan kolumnar dari mesin eksekusi native.
Contoh: Vektorisasi Python UDF
import pandas as pd
from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType
@pandas_udf(DoubleType())
def calculate_discount(price: pd.Series, rate: pd.Series) -> pd.Series:
return price * (1 - rate)
df = spark.table("sales.transactions")
result = df.withColumn("discounted_price", calculate_discount(df.price, df.discount_rate))
result.show()
Tidak ada konfigurasi tambahan yang diperlukan selain mengaktifkan mesin eksekusi asli. UDF Python yang sudah ada secara otomatis mendapatkan manfaat.
Dukungan UDF Scala
Mesin eksekusi asli juga mempercepat UDF Scala. Karena UDF Scala berjalan secara asli di JVM, mesin dapat membongkar operasi yang didukung ke jalur eksekusi C++ yang divektorisasi sambil menjaga evaluasi Scala UDF tetap efisien dalam runtime yang sama.
Contoh: Scala UDF
import org.apache.spark.sql.functions.udf
val toUpperCase = udf((s: String) => s.toUpperCase)
val df = spark.table("catalog.customers")
val result = df.withColumn("name_upper", toUpperCase(df("name")))
result.show()
UDF Scala yang beroperasi pada jenis data yang didukung dipercepat tanpa perubahan kode saat mesin eksekusi asli diaktifkan.
Dukungan jenis data kompleks
Arsitektur lakehouse modern bergantung pada data semi terstruktur dan berlapis. Mesin eksekusi asli sekarang memberikan dukungan yang dioptimalkan untuk:
| Jenis data | Deskripsi | Contoh kasus penggunaan |
|---|---|---|
| Array | Pengumpulan elemen yang diurutkan | Tag acara, kategori produk |
| Peta | Pasangan kunci-nilai | Properti konfigurasi, metadata |
| Struktur | Bidang bernama dengan jenis yang berbeda | Catatan pelanggan berlapis, objek alamat |
Operasi yang didukung untuk jenis kompleks
Mesin eksekusi asli mempercepat operasi umum pada jenis data yang kompleks:
- Fungsi array:
explode,array_contains,size,flatten,transform - Fungsi peta:
map_keys,map_values,element_at - Akses struct: akses field dengan notasi titik,
getField - Kombinasi bertingkat: array dari struct, map dengan nilai berupa array
Contoh: Bekerja dengan array dan struktur
from pyspark.sql.functions import explode, col, size
# Read data with nested schema
df = spark.table("events.telemetry")
# Operations on arrays - accelerated by native engine
result = (df
.filter(size(col("tags")) > 0)
.select(
col("event_id"),
col("metadata.source"), # Struct field access
explode(col("tags")).alias("tag")
)
)
result.show()
Contoh: Bekerja dengan peta
from pyspark.sql.functions import map_keys, map_values, col
df = spark.table("config.settings")
# Map operations - accelerated by native engine
result = (df
.select(
col("setting_id"),
map_keys(col("properties")).alias("keys"),
map_values(col("properties")).alias("values")
)
)
result.show()
Hasil performa
Tolok ukur internal menunjukkan peningkatan signifikan di seluruh beban kerja yang menggunakan Python UDF dan jenis data kompleks:
| Tipe beban kerja | Peningkatan performa |
|---|---|
| UDF Python yang tervektorisasi | Hingga 5,76x lebih cepat |
| Python UDF skalar | Hingga 1,08x lebih cepat |
| TPC-DS end-to-end (dengan tipe kompleks) | Hingga 2,35x lebih cepat |
Peningkatan ini berasal dari berkurangnya beban serialisasi, meningkatnya vektorisasi, dan eksekusi kolumnar secara menyeluruh.
Manfaat untuk pola lakehouse tingkat lanjut
Akselerasi jenis data kompleks sangat penting untuk:
- Pengoptimalan Z-ORDER: Kolom berlapis berpartisipasi dalam tata letak data yang dioptimalkan.
- Pengklusteran cairan: Kolom jenis kompleks mendapat manfaat dari pengklusteran tanpa meratakan.
- Analitik semi-terstruktur: Payload JSON dan eventstream tetap bersarang untuk kueri alami.
- Arsitektur berbasis peristiwa: Data telemetri dan IoT mempertahankan struktur hierarkisnya.
Alih-alih meratakan data atau merestrukturisasi alur untuk performa, bekerja secara alami dengan skema kompleks sambil mempertahankan efisiensi eksekusi yang tinggi.
Aktifkan Fitur
Dukungan untuk Python UDF, Scala UDF, dan tipe data kompleks tersedia ketika mesin eksekusi asli diaktifkan. Tidak diperlukan konfigurasi tambahan.
Untuk mengaktifkan mesin eksekusi asli, lihat Mesin eksekusi asli untuk Fabric Data Engineering.
Prerequisites
- Runtime 1.3 (Apache Spark 3.5) atau Runtime 2.0 (Apache Spark 4.0).
- Mesin eksekusi native diaktifkan pada tingkat lingkungan, notebook, atau definisi tugas Spark.
Keterbatasan
- Tidak semua pustaka Python didukung dalam jalur vektorisasi. Pustaka yang memerlukan serialisasi objek Python sembarang mungkin tetap dapat memicu fallback.
- Jenis kompleks yang sangat berlapis (misalnya, array peta struktur) mungkin kembali ke mesin JVM untuk operasi tertentu.
- Mode ANSI tidak didukung dengan mesin eksekusi asli.