Inferência de ONNX no Spark

ONNX (Open Neural Network Exchange) fornece um runtime portátil com otimização de hardware para modelos de machine learning. Ao converter um modelo no formato ONNX, você pode executar a inferência em lote no Spark com latência mais baixa e sem depender da estrutura de treinamento original no momento da previsão.

Neste artigo, você treinará um modelo LightGBM com SynapseML, converterá-o no formato ONNX e, em seguida, usará o modelo ONNX para executar a inferência no Spark em Microsoft Fabric.

Pré-requisitos

  • Anexe o notebook a um lakehouse. No lado esquerdo do bloco de anotações, selecione Adicionar para adicionar uma lakehouse existente ou criar uma.
  • Fabric Runtime 1.2 ou posterior.

Instalar pacotes necessários

Execute a célula a seguir no notebook para instalar os pacotes necessários. O pacote onnxmltools não está pré-instalado no runtime do Fabric.

%pip install onnxmltools --quiet

Após a conclusão da instalação, verifique se os pacotes estão disponíveis:

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

Note

O pacote lightgbm é pré-instalado no Fabric Runtime 1.2 e posterior. Você só precisa instalar onnxmltools.

Carregar os dados de exemplo

Carregue o conjunto de dados de previsão de falência do Armazenamento de Blobs do Azure público:

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

A tabela exibida inclui colunas como:

Falido? Indicador de Lucro Líquido Equidade para responsabilidade
0 1.0 0.0165
0 1.0 0.0208

Treinar um modelo LightGBM

Use o VectorAssembler para combinar colunas de características e, em seguida, treine um 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)

Verifique se o modelo foi treinado com êxito:

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

Converter o modelo em formato ONNX

Exporte o modelo treinado para um propulsor LightGBM e converta-o em 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))

Verifique se a conversão ONNX foi bem-sucedida:

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

A saída mostra o tamanho da carga do modelo ONNX em bytes (normalmente em torno de 800.000 bytes).

Importante

Use from onnxmltools.convert.common.data_types import FloatTensorType para a definição de tipo. O caminho from onnxconverter_common.data_types import FloatTensorType de importação mais antigo é incompatível com as versões atuais de onnxmltools.

Carregar e configurar o modelo ONNX

Carregue a carga útil do ONNX em uma instância do SynapseML ONNXModel e inspecione as entradas e saídas do modelo:

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

A saída lista os nós de entrada e saída do modelo.

Configure o modelo mapeando colunas de entrada e saída. O FeedDict mapeia os nomes de entrada do modelo ONNX para os nomes das colunas do DataFrame. O FetchDict mapeia os nomes de coluna de saída desejados para os nomes de saída do modelo ONNX:

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

Executar inferência

Crie dados de teste e transforme-os por meio do modelo ONNX:

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

Note

Como os dados de teste são gerados aleatoriamente, os valores de previsão não representam resultados reais. Esta seção demonstra que o modelo ONNX é executado corretamente no Spark.

A saída deve conter colunas para features, predictione probability:

Features previsão probabilidade
{"type":1,"values":[0.105... 0 {"0":0.835...
{"type":1,"values":[0.814... 0 {"0":0.658...

Verifique os resultados produzidos pela inferência:

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"

A saída confirma que todas as linhas de teste receberam pontuação e que o DataFrame de resultado contém as colunas features, prediction e probability.

Solução de problemas

Questão Cause Resolução
ModuleNotFoundError: No module named 'onnxmltools' O pacote não está pré-instalado em Fabric runtime. Execute %pip install onnxmltools --quiet e reinicie o kernel Python.
RuntimeError: Operator LgbmClassifier got an input with a wrong type Caminho de importação incorreto para FloatTensorType. Use from onnxmltools.convert.common.data_types import FloatTensorType em vez de importar de onnxconverter_common.data_types.
ModuleNotFoundError: No module named 'onnx.mapping' Versão 1.7.0 ou anterior de onnxmltools incompatível com o pacote atual onnx. Execute %pip install onnxmltools --upgrade --quiet para instalar uma versão compatível.
ONNX conversion returns empty payload Falha na extração de cadeia de caracteres do modelo de booster. Verifique se retorna model.getLightGBMBooster().modelStr().get() uma cadeia de caracteres não vazia antes da conversão.
Feature (Column_) appears more than one time durante model.fit() Colunas de um conjunto de dados com caracteres especiais geram nomes duplicados após a sanitização do LightGBM. Adicione .setDataTransferMode("bulk") à LightGBMClassifier configuração. O modo em massa usa a Seta do Apache e evita o problema de limpeza do nome da coluna.
AssertionError no SparkContext em ONNXModel() A sessão do Spark não é inicializada. Execute este código em um notebook do Fabric com um lakehouse anexado. A spark variável é pré-inicializada pelo runtime.

Limpar os recursos

Se você não precisar mais do DataFrame de teste armazenado em cache, cancele-o para liberar a memória do cluster:

testDf.unpersist()