You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将Spark训练的AlsModel部署到Spring Boot微服务以提供推荐

Great question! Your approach of training the ALS model in a separate Spark cluster, exporting it to a portable format, and loading it into a Spark-free Spring Boot service is totally feasible—and perfect for handling regular model updates. Let’s break down the best ways to implement this, including two solid options for model serialization, plus strategies for smooth updates.

实现无Spark依赖的Spring Boot ALS推荐服务

一、方案选型:PMML vs MLeap

You’ve mentioned PMML, which is a standard for model serialization—but there’s another tool called MLeap that’s purpose-built for Spark model deployment and often works better for ALS. Here’s a quick comparison:

  • PMML: A universal standard, but Spark’s ALS implementation has some compatibility quirks with PMML (especially around feature handling).
  • MLeap: Lightweight, Spark-native, and designed to run models without a Spark cluster. It has full support for ALS and better performance for real-time predictions.

We’ll cover both approaches below, starting with PMML (since you mentioned it) and then diving into the more robust MLeap option.

二、基于PMML的实现步骤

1. 训练ALS模型并导出PMML

First, in your Spark training pipeline, use the pmml-sparkml library to export your trained ALS model to a PMML file. Make sure your Spark and PMML library versions are compatible (e.g., Spark 3.x pairs with pmml-sparkml 1.5.x+).

Scala code example for exporting:

import org.jpmml.sparkml.PMMLBuilder
import java.io.FileOutputStream

// Assume you've already trained your ALS model (alsModel)
val pmml = new PMMLBuilder(spark.sparkContext, alsModel.getModel).build()

// Save to local file (then transfer to your service layer)
val fos = new FileOutputStream("als-recommender.pmml")
org.jpmml.model.SerializationUtil.serializePMML(pmml, fos)
fos.close()

2. 传输PMML文件到Spring Boot服务

Use your existing data transfer method (SCP, FTP, CI/CD pipeline, or database storage) to move the PMML file to your Spring Boot service’s accessible directory.

3. Spring Boot中加载PMML并生成推荐

Add the PMML evaluator dependency to your pom.xml (or build.gradle):

<dependency>
    <groupId>org.jpmml</groupId>
    <artifactId>pmml-evaluator</artifactId>
    <version>1.6.4</version>
</dependency>

Then implement a service to load the model and handle predictions:

import org.jpmml.evaluator.Evaluator;
import org.jpmml.evaluator.EvaluatorUtil;
import org.jpmml.model.PMMLUtil;
import org.dmg.pmml.PMML;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
import java.io.File;
import java.util.HashMap;
import java.util.Map;

@Service
public class PMMLRecommendationService {
    private Evaluator evaluator;

    @PostConstruct
    public void loadModel() throws Exception {
        // Load the PMML file from your service's directory
        PMML pmml = PMMLUtil.unmarshal(new File("als-recommender.pmml"));
        evaluator = EvaluatorUtil.createEvaluator(pmml);
    }

    public Double predictRating(Long userId, Long itemId) {
        // Match input field names to what you used during training
        Map<String, Object> input = new HashMap<>();
        input.put("user", userId);
        input.put("item", itemId);

        // Run prediction
        Map<String, ?> result = evaluator.evaluate(input);
        return (Double) result.get("prediction");
    }
}

三、更优选择:基于MLeap的实现

MLeap is the better choice for Spark ALS models because it avoids PMML’s compatibility issues and offers faster inference. Here’s how to implement it:

1. 训练ALS模型并导出MLeap Bundle

Add MLeap dependencies to your Spark training project, then export the model as a MLeap bundle (a zip file containing the model and schema):

import ml.combust.bundle.BundleFile
import ml.combust.mleap.spark.SparkSupport._
import org.apache.spark.ml.bundle.SparkBundleContext
import java.io.File

// Assume you have your trained ALS model (alsModel) and training data (trainingData)
val sbc = SparkBundleContext().withDataset(trainingData)

// Export to a zip bundle
val bundleFile = BundleFile(File.createTempFile("als-bundle", ".zip").toURI)
alsModel.writeBundle.save(bundleFile)(sbc).get

2. 传输Bundle到Spring Boot服务

Same as PMML: use your existing transfer method to move the zip bundle to your Spring Boot service.

3. Spring Boot中加载MLeap Bundle并预测

Add the MLeap runtime dependency to your Spring Boot project:

<dependency>
    <groupId>ml.combust.mleap</groupId>
    <artifactId>mleap-runtime_2.12</artifactId>
    <version>0.20.0</version>
</dependency>

Implement the recommendation service:

import ml.combust.bundle.BundleFile;
import ml.combust.mleap.runtime.frame.DefaultLeapFrame;
import ml.combust.mleap.runtime.frame.LeapFrame;
import ml.combust.mleap.core.types.StructType;
import ml.combust.mleap.runtime.transformer.Transformer;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
import java.io.File;
import java.util.Arrays;
import java.util.List;

@Service
public class MLeapRecommendationService {
    private Transformer model;
    private StructType inputSchema;

    @PostConstruct
    public void loadModel() throws Exception {
        // Load the MLeap bundle zip file
        BundleFile bundleFile = BundleFile(new File("als-bundle.zip").toURI);
        model = bundleFile.loadMleapBundle().get().root();
        inputSchema = model.inputSchema();
    }

    public Double predictRating(Long userId, Long itemId) {
        // Create input data matching the training schema
        List<Object> inputRow = Arrays.asList(userId, itemId);
        LeapFrame inputFrame = new DefaultLeapFrame(inputSchema, Arrays.asList(inputRow));

        // Run prediction
        LeapFrame resultFrame = model.transform(inputFrame).get();
        // Extract the prediction value (adjust index if needed)
        return resultFrame.select("prediction").getDouble(0, 0).get();
    }
}

四、定期更新模型的策略

To handle regular model updates without downtime, use one of these strategies:

1. 定时拉取更新

Use Spring’s @Scheduled annotation to periodically check for a new model file (or version flag in a database). Load the new model atomically to avoid breaking in-flight requests:

import org.jpmml.evaluator.Evaluator;
import java.util.concurrent.atomic.AtomicReference;
import org.springframework.scheduling.annotation.Scheduled;

@Service
public class UpdatableRecommendationService {
    private AtomicReference<Evaluator> evaluatorRef = new AtomicReference<>();

    @PostConstruct
    public void init() {
        loadLatestModel();
    }

    // Check for updates every hour (adjust interval as needed)
    @Scheduled(fixedRate = 3600000)
    public void checkForUpdates() {
        loadLatestModel();
    }

    private void loadLatestModel() {
        try {
            PMML pmml = PMMLUtil.unmarshal(new File("latest-als.pmml"));
            Evaluator newEvaluator = EvaluatorUtil.createEvaluator(pmml);
            // Atomically replace the old model
            evaluatorRef.set(newEvaluator);
        } catch (Exception e) {
            // Log error and keep using the old model
            e.printStackTrace();
        }
    }

    public Double predictRating(Long userId, Long itemId) {
        Evaluator currentEvaluator = evaluatorRef.get();
        // Run prediction with the current model
        Map<String, Object> input = new HashMap<>();
        input.put("user", userId);
        input.put("item", itemId);
        Map<String, ?> result = currentEvaluator.evaluate(input);
        return (Double) result.get("prediction");
    }
}

2. 事件触发更新

Have your Spark training pipeline call a Spring Boot API endpoint when a new model is ready. The service then immediately loads the new model:

@RestController
@RequestMapping("/model")
public class ModelUpdateController {
    private final UpdatableRecommendationService recommendationService;

    public ModelUpdateController(UpdatableRecommendationService recommendationService) {
        this.recommendationService = recommendationService;
    }

    @PostMapping("/update")
    public ResponseEntity<String> triggerUpdate() {
        recommendationService.loadLatestModel();
        return ResponseEntity.ok("Model updated successfully");
    }
}

3. 蓝绿部署

For large models where loading takes time, deploy two instances of your Spring Boot service: one using the old model, one with the new model. Switch traffic to the new instance once it’s fully loaded, then shut down the old one.

五、关键注意事项

  • Version Compatibility: Ensure your Spark training version matches the PMML/MLeap library version (e.g., Spark 3.x needs MLeap 0.20.x+).
  • Field Name Consistency: The input field names in your Spring Boot service must exactly match what you used during Spark training (e.g., if you called the user ID field user_id in training, don’t use userId in the service).
  • Performance: For high-concurrency scenarios, cache frequent predictions or use a thread pool to handle prediction requests. MLeap generally outperforms PMML in inference speed.

内容的提问来源于stack exchange,提问作者Kaojo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:17:38