Python UDF, Scala UDF, dan jenis data kompleks di mesin eksekusi asli

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:

  1. Konversi data dari format internal Spark.
  2. Serialisasi dan pemindahan ke proses pekerja Python.
  3. Python eksekusi UDF.
  4. Serialisasi hasil kembali ke JVM.
  5. 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_udf yang 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

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.