ONNX-slutning på Spark

ONNX (Open Neural Network Exchange) tilbyr en portabel, maskinvareoptimalisert kjøretid for maskinlæringsmodeller. Ved å konvertere en modell til ONNX-format kan du kjøre batch-inferens på Spark med lavere latens og uten å være avhengig av det opprinnelige treningsrammeverket ved prediksjonstidspunktet.

I denne artikkelen trener du en LightGBM-modell med SynapseML, konverterer den til ONNX-format, og bruker deretter ONNX-modellen for å utføre inferenser på Spark i Microsoft Fabric.

Forutsetning

  • Få et Microsoft Fabric-abonnement. Eller registrer deg for en gratis prøveversjon av Microsoft Fabric.

  • Logg på Microsoft Fabric.

  • Bytt til Fabric ved å bruke erfaringsbryteren nederst til venstre på hjemmesiden din.

    Skjermbilde som viser valget av Fabric i menyen for opplevelsesbytter.

  • Legg notatblokken til et lakehouse. På venstre side av notatboken din velger du Legg til for å legge til et eksisterende innsjøhus eller opprette et.
  • Fabric Runtime 1.2 eller nyere.

Installere nødvendige pakker

Kjør neste celle i notatboken din for å installere de nødvendige pakkene. onnxmltools-pakken er ikke forhåndsinstallert i Fabric-kjøretiden.

%pip install onnxmltools --quiet

Etter at installasjonen er fullført, sjekk at pakkene er tilgjengelige:

import onnxmltools
import lightgbm
print(f"onnxmltools version: {onnxmltools.__version__}")
print(f"lightgbm version: {lightgbm.__version__}")

Bemerkning

lightgbm-pakken er forhåndsinstallert i Fabric Runtime 1.2 og nyere. Du trenger bare å installere onnxmltools.

Laste inn eksempeldataene

Load the bankruptcy prediction dataset from public Azure Blob Storage:

df = (
    spark.read.format("csv")
    .option("header", True)
    .option("inferSchema", True)
    .load(
        "wasbs://publicwasb@mmlspark.blob.core.windows.net/company_bankruptcy_prediction_data.csv"
    )
)

print(f"Rows: {df.count()}, Columns: {len(df.columns)}")
display(df.limit(5))

Den viste tabellen inkluderer kolonner som:

Konkurs? Nettoinntektsflagg Egenkapital til ansvar
0 1.0 0.0165
0 1.0 0.0208

Tren en LightGBM-modell

Bruk for VectorAssembler å kombinere funksjonskolonner, og tren deretter en LightGBMClassifier:

from pyspark.ml.feature import VectorAssembler
from synapse.ml.lightgbm import LightGBMClassifier

feature_cols = df.columns[1:]
featurizer = VectorAssembler(inputCols=feature_cols, outputCol="features")

train_data = featurizer.transform(df)["Bankrupt?", "features"]

model = (
    LightGBMClassifier(featuresCol="features", labelCol="Bankrupt?")
    .setDataTransferMode("bulk")
    .setEarlyStoppingRound(300)
    .setLambdaL1(0.5)
    .setNumIterations(1000)
    .setNumThreads(-1)
    .setMaxDeltaStep(0.5)
    .setNumLeaves(31)
    .setMaxDepth(-1)
    .setBaggingFraction(0.7)
    .setFeatureFraction(0.7)
    .setBaggingFreq(2)
    .setObjective("binary")
    .setIsUnbalance(True)
    .setMinSumHessianInLeaf(20)
    .setMinGainToSplit(0.01)
)

model = model.fit(train_data)

Verifiser at modellen er trent vellykket:

print(f"Model type: {type(model).__name__}")
print(f"Number of features: {len(feature_cols)}")

Konvertere modellen til ONNX-format

Eksporter den trente modellen til en LightGBM-booster, og konverter den deretter til ONNX:

import lightgbm as lgb
from typing import Union
from lightgbm import Booster, LGBMClassifier
from onnxmltools.convert import convert_lightgbm
from onnxmltools.convert.common.data_types import FloatTensorType


def convert_to_onnx(lgbm_model: Union[LGBMClassifier, Booster], input_size: int) -> bytes:
    initial_types = [("input", FloatTensorType([-1, input_size]))]
    onnx_model = convert_lightgbm(
        lgbm_model, initial_types=initial_types, target_opset=13
    )
    return onnx_model.SerializeToString()


booster_model_str = model.getLightGBMBooster().modelStr().get()
booster = lgb.Booster(model_str=booster_model_str)
model_payload_ml = convert_to_onnx(booster, len(feature_cols))

Verifiser at ONNX-konverteringen lyktes:

print(f"ONNX model payload size: {len(model_payload_ml)} bytes")
assert len(model_payload_ml) > 0, "ONNX conversion failed: empty payload"

Utdataene viser ONNX-modellens nyttelaststørrelse i bytes (typisk rundt 800 000 byte).

Important

Bruk from onnxmltools.convert.common.data_types import FloatTensorType for typedefinisjonen. Den eldre importstien from onnxconverter_common.data_types import FloatTensorType er inkompatibel med nåværende versjoner av onnxmltools.

Last inn og konfigurer ONNX-modellen

Last inn ONNX-nyttelasten i en SynapseML ONNXModel og inspiser modellens innganger og utganger:

from synapse.ml.onnx import ONNXModel

onnx_ml = ONNXModel().setModelPayload(model_payload_ml)

print("Model inputs:" + str(onnx_ml.getModelInputs()))
print("Model outputs:" + str(onnx_ml.getModelOutputs()))

Utgangen lister modellens inn- og utgangsnoder.

Konfigurer modellen ved å kartlegge inn- og utgangskolonner. Den FeedDict mapper ONNX-modellens inndatanavn til DataFrame-kolonnenavn. Navnene FetchDict på ønskede utdatakolonner til ONNX-modellens utdatanavn:

onnx_ml = (
    onnx_ml.setDeviceType("CPU")
    .setFeedDict({"input": "features"})
    .setFetchDict({"probability": "probabilities", "prediction": "label"})
    .setMiniBatchSize(5000)
)

Run-inferens

Lag testdata og transformer dem gjennom ONNX-modellen:

from pyspark.ml.feature import VectorAssembler
import pandas as pd
import numpy as np

n = 10000
m = 95
test = np.random.rand(n, m)
testPdf = pd.DataFrame(test)
cols = list(map(str, testPdf.columns))
testDf = spark.createDataFrame(testPdf)
testDf = testDf.repartition(4)
testDf = (
    VectorAssembler()
    .setInputCols(cols)
    .setOutputCol("features")
    .transform(testDf)
    .drop(*cols)
    .cache()
)

display(onnx_ml.transform(testDf))

Bemerkning

Fordi testdataene genereres tilfeldig, representerer ikke prediksjonsverdiene virkelige resultater. Denne delen demonstrerer at ONNX-modellen kjører korrekt på Spark.

Utdataene skal inneholde kolonner for features, prediction, og probability:

Funksjoner prediksjon sannsynlighet
{"type":1,"values":[0.105... 0 {"0":0.835...
{"type":1,"values":[0.814... 0 {"0":0.658...

Verifiser resultatene som har gitt slutningen:

results = onnx_ml.transform(testDf)
print(f"Result count: {results.count()}")
print(f"Output columns: {results.columns}")
assert "prediction" in results.columns, "Missing prediction column"
assert "probability" in results.columns, "Missing probability column"

Utdataene bekrefter at alle testrader ble scoret, og resultatdatarammen inneholder features, prediction, og probability kolonnene.

Feilsøking

Problem Årsak Løsning
ModuleNotFoundError: No module named 'onnxmltools' Pakken er ikke forhåndsinstallert i Fabric-runtime. Kjør %pip install onnxmltools --quiet og start Python-kjernen på nytt.
RuntimeError: Operator LgbmClassifier got an input with a wrong type Feil importvei for FloatTensorType. Bruk from onnxmltools.convert.common.data_types import FloatTensorType i stedet for å importere fra onnxconverter_common.data_types.
ModuleNotFoundError: No module named 'onnx.mapping' Inkompatibel onnxmltools versjon 1.7.0 eller tidligere med nåværende onnx pakke. Kjør %pip install onnxmltools --upgrade --quiet for å installere en kompatibel versjon.
ONNX conversion returns empty payload Booster-modellens strengutvinning mislyktes. Sjekk at det model.getLightGBMBooster().modelStr().get() returnerer en ikke-tom streng før konvertering.
Feature (Column_) appears more than one time Under model.fit() Datasettkolonner med spesialtegn produserer dupliserte navn etter LightGBM-rensing. Legg til i .setDataTransferMode("bulk") konfigurasjonen LightGBMClassifier . Bulk-modus bruker Apache Arrow og unngår problemet med kolonnenavn-sanering.
AssertionError på SparkContext i ONNXModel() Spark-økten er ikke initialisert. Kjør denne koden i en Fabric-notatbok med et innsjøhus festet. Variabelen spark er forhåndsinitialisert av kjøretiden.

Rydd opp ressurser

Hvis du ikke lenger trenger den bufrede test-DataFrame, kan du avlagre den for å frigjøre klyngeminne:

testDf.unpersist()