Python UDF, Scala UDF a zložité dátové typy v natívnom výkonnom engine

Natívny výkonný engine v Microsoft Fabric teraz podporuje Python používateľsky definované funkcie (UDF), Scala UDF a zložité dátové typy (polia, mapy a štruktúry). Tieto schopnosti vám umožňujú písať expresívne Spark aplikácie bez straty výkonu.

Podpora Python UDF

Python je jeden z najpopulárnejších jazykov v dátovom inžinierstve a dátovej vede. Historicky prinášali Python UDF značné režijné náklady v Sparku kvôli nákladom na serializáciu medzi JVM a Python worker procesmi. Natívny výkonný engine minimalizuje tieto nákladné prechody, čo umožňuje rýchlejšie vykonávanie bez zmien kódu.

Ako fungujú Python UDF v natívnom výkonnom engine

V konvenčnom modeli vykonávania Spark zahŕňa vykonávanie Python UDF:

  1. Konverzia dát z interného formátu Sparku.
  2. Serializácia a prenos do pracovných procesov v Python.
  3. Python UDF spustenie.
  4. Serializácia výsledkov späť do JVM.
  5. Spark pokračuje v poprave.

Tento pohyb naprieč behom spôsobuje náklady na serializáciu/deserializáciu, neefektívnosť CPU a nefunkčné stĺpcové vykonávacie pipeline. Natívny výkonný engine znižuje túto režijnú záťaž optimalizáciou cesty prenosu dát a udržiavaním vektorizovaného spracovania, kde je to možné.

Podporované typy UDF v Python

Natívny vykonávací engine podporuje:

  • Skalárne UDFs: Funkcie Python po riadkoch registrované v udf().
  • Vektorizované (Pandas) UDF: Funkcie @pandas_udf vybavené týmto procesorom pracujú na dávkach dát pomocou Apache Arrow pre efektívny prenos.

Vektorizované UDF dosahujú najväčšie výkonnostné zisky, pretože sa prirodzene zosúladia so stĺpcovým modelom natívneho výkonného enginu.

Príklad: Vektorizovaný 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()

Nie je potrebná žiadna ďalšia konfigurácia okrem povolenia natívneho vykonávacieho enginu. Existujúce Python UDF automaticky profitujú.

Podpora Scala UDF

Natívny vykonávací engine tiež urýchľuje Scala UDF. Keďže Scala UDF bežia natívne v JVM, engine môže preniesť podporované operácie na vektorizovanú C++ vykonávaciu cestu a zároveň udržiavať efektívne vyhodnocovanie Scala UDF v rovnakom runtime.

Príklad: 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()

Scala UDF, ktoré pracujú na podporovaných dátových typoch, sa zrýchľujú bez zmien kódu, keď je povolený natívny výkonný engine.

Podpora komplexných dátových typov

Moderné architektúry jazerných domov závisia od polostruktúrovaných a vnorených dát. Natívny výkonný engine teraz poskytuje optimalizovanú podporu pre:

Typ údajov Description Príklad použitia
Pole Usporiadaná kolekcia prvkov Tagy udalostí, kategórie produktov
Mapa Páry kľúč-hodnota Konfiguračné vlastnosti, metadáta
Struct Pomenované polia s rôznymi typmi Vnorené zákaznícke záznamy, adresné objekty

Podporované operácie pre zložité typy

Natívny výkonný engine urýchľuje bežné operácie na zložitých dátových typoch:

  • Funkcie poľa: explode, , , array_contains, sizeflattentransform
  • Mapovacie funkcie: map_keys, map_values, element_at
  • Prístup k štruktúre: Prístup k polí s dot notáciou, getField
  • Vnorené kombinácie: Polia štruktúr, mapy s hodnotami polí

Príklad: Práca s poliami a štruktúrami

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()

Príklad: Práca s mapami

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()

Výkonnostné výsledky

Interné benchmarkovanie ukazuje významné zlepšenia naprieč pracovnými záťažami využívajúcimi Python UDF a zložité dátové typy:

Typ pracovnej záťaže Zlepšenie výkonu
Vektorizované Python UDF Až 5,76x rýchlejšie
Skalárne Python UDF Až 1,08x rýchlejšie
TPC-DS end-to-end (s komplexnými typmi) Až 2,35x rýchlejšie

Tieto zisky sú výsledkom znížených režijných nákladov na serializáciu, zlepšenej vektorizácie a end-to-end stĺpcového vykonávania.

Výhody pokročilých vzorov jazerných domov

Zrýchlenie komplexných dátových typov je obzvlášť dôležité pre:

  • Optimalizácia Z-ORDER: Vnorené stĺpce sa podieľajú na optimalizovanom rozložení dát.
  • Kvapalné zhlukovanie: Komplexné typy stĺpcov profitujú z zhlukovania bez sploštenia.
  • Polostruktúrovaná analytika: JSON payloady a eventstreamy zostávajú vnorené pre prirodzené dotazovanie.
  • Architektúry riadené udalosťami: Telemetria a IoT dáta si zachovávajú svoju hierarchickú štruktúru.

Namiesto splošťovania dát alebo reštrukturalizácie pipeline pre výkon pracujte prirodzene so zložitými schémami pri zachovaní vysokej efektivity vykonávania.

Aktivácia funkcie

Python UDF, Scala UDF a podpora zložitých dátových typov sú dostupné, keď je povolený natívny výkonný engine. Nie je potrebná žiadna ďalšia konfigurácia.

Pre povolenie natívneho vykonávacieho enginu pozri Natívny výkonný engine pre Fabric Data Engineering.

Požiadavky

Obmedzenia

  • Nie všetky Python knižnice sú podporované v rámci vektorizovanej cesty. Knižnice, ktoré vyžadujú ľubovoľnú Python objektovú serializáciu, môžu stále vyvolať záložný systém.
  • Hlboko vnorené komplexné typy (napríklad polia máp štruktúr) sa môžu pri určitých operáciách vrátiť k JVM enginu.
  • ANSI režim nie je podporovaný natívnym vykonávacím enginom.