大数据领域Spark的机器学习模型部署:从实验室到生产的全链路实践

引言:模型训练好了,然后呢?

三年前,我带着团队用Spark ML训练了一个用户流失(Churn)预测模型——基于100万条用户行为数据,结合性别、合约类型、在网时长、月消费等特征,用逻辑回归算法达到了85%的准确率。当我们兴奋地拿着结果找业务团队时,却被泼了一盆冷水:

“这个模型能实时告诉我,现在正在咨询的用户会不会流失吗?”
“每天新增的10万条数据,怎么快速得到预测结果?”
“如果模型迭代了,怎么保证线上服务不会出错?”

那时我才意识到:模型训练只是开始,部署才是把数据价值转化为业务价值的关键。Spark ML的模型天生是为分布式批处理设计的,但生产环境需要低延迟的在线服务、高吞吐量的批处理、实时的流计算——这些需求和Spark的“原生属性”存在根本矛盾。

这篇文章,我会结合15年的大数据工程经验,从核心挑战、基础概念、场景实战、性能优化、监控运维五个维度,手把手教你把Spark ML模型落地生产。无论是离线批处理、在线API服务还是实时流计算,你都能找到可复制的解决方案。

一、Spark ML模型部署的4大核心挑战

在讲具体方案前,我们得先搞清楚:为什么Spark模型部署这么难?

挑战1:分布式模型与在线服务的矛盾

Spark ML的模型(如PipelineModel)依赖SparkContext——这是一个重量级的分布式计算入口,启动时需要加载集群配置、连接资源管理器(YARN/K8s),耗时5~30秒。而在线服务需要毫秒级响应(比如100ms以内),直接用SparkContext做预测,延迟会高到无法接受。

挑战2:模型序列化与跨平台兼容

Spark原生的save()方法会把模型保存为Parquet格式+元数据,只能用Spark加载。如果你的业务系统是Python Flask、Java Spring Boot或Go服务,想直接用这个模型,门都没有——因为它们无法解析Spark的序列化格式。

挑战3:性能与延迟的平衡

  • 离线批处理需要高吞吐量(比如每小时处理1000万条数据),但Spark的Shuffle操作会成为瓶颈;
  • 在线服务需要低延迟(比如单请求<100ms),但Spark的JVM overhead会拖慢速度;
  • 实时流计算需要低延迟+高可用(比如处理Kafka的每秒1万条消息),但模型更新会导致流任务重启。

挑战4:模型版本管理与监控

生产环境中,模型会不断迭代(比如每周更新一次),你需要解决:

  • 如何管理不同版本的模型?(比如“Production版”“Staging版”“Dev版”)
  • 如何快速回滚有问题的模型?
  • 如何监控模型的性能(比如准确率下降)和服务的健康状态(比如QPS骤降)?

二、基础概念:Spark ML模型与序列化

要解决部署问题,先得理解Spark ML的核心组件和模型序列化的本质。

1. Spark ML的核心组件

Spark ML的API是管道化(Pipeline)设计的,核心组件包括:

组件作用方法
Transformer转换数据(如特征编码、归一化)transform()
Estimator训练模型(如逻辑回归、随机森林)fit()
Pipeline将多个Transformer和一个Estimator组合成“流水线”fit() → PipelineModel
PipelineModel训练好的流水线模型(包含所有Transformer的参数和Estimator的模型)transform()

举个例子,一个用户流失预测的Pipeline结构如下:

graph TD
    A[原始数据] --> B[StringIndexer(性别转索引)]
    B --> C[OneHotEncoder(索引转独热编码)]
    C --> D[VectorAssembler(组装特征向量)]
    D --> E[LogisticRegression(训练分类器)]
    E --> F[PipelineModel(最终模型)]

PipelineModel是我们要部署的核心——它包含了从“原始数据”到“预测结果”的所有步骤。

2. 模型序列化的3种方式

序列化是将模型转化为可存储/传输格式的过程,Spark ML支持以下3种方式:

方式1:Spark原生格式(Parquet+元数据)
  • 用法:model.save("hdfs://path/to/model")
  • 原理:用Parquet存储模型参数(如逻辑回归的系数),用JSON存储元数据(如Pipeline的 stages 结构)。
  • 缺点:只能用Spark加载,无法跨语言/框架。
  • 适用场景:纯Spark生态的离线批处理。
方式2:MLeap(跨平台的Spark模型序列化)

MLeap是专为Spark ML设计的跨平台序列化框架,支持将PipelineModel转成MLeap Bundle,可以用Java、Scala、Python的Runtime加载,不需要Spark依赖。

  • 原理:用Protocol Buffers(PB)序列化模型结构,用Arrow格式存储数据,保证跨平台兼容性。
  • 优点:轻量(加载时间<1秒)、跨语言、支持所有Spark ML组件。
  • 适用场景:在线服务、跨语言部署。
方式3:ONNX(开放神经网络交换格式)

ONNX是工业级的模型交换格式,支持Spark ML、TensorFlow、PyTorch等框架。

  • 用法:用onnxmltools将Spark模型转成ONNX格式。
  • 优点:通用性强,支持更多框架。
  • 缺点:对Spark ML的支持不够完善(比如部分Transformer无法转换)。
  • 适用场景:多框架联合部署(如Spark特征工程+PyTorch模型)。

总结对比:

方式跨平台轻量性Spark兼容性适用场景
Spark原生❌❌✅纯Spark离线批处理
MLeap✅✅✅在线服务、跨语言部署
ONNX✅✅⚠️多框架联合部署

三、Spark ML模型部署的3大场景与实战

根据业务需求,Spark模型的部署场景主要分为离线批处理、在线服务、实时流计算。我们逐个拆解,每个场景都给出可运行的代码和踩坑指南。

场景1:离线批处理部署——Spark Submit + MLflow

适用场景:每天/每小时处理大规模历史数据(如用户画像更新、风险评分批量计算)。
核心需求:高吞吐量、模型版本管理、结果可追溯。

1. 技术选型
  • 模型训练:PySpark + Spark ML
  • 版本管理:MLflow(跟踪模型版本、参数、指标)
  • 执行方式:Spark Submit(提交到YARN/K8s集群)
2. 实战步骤

步骤1:训练模型并注册到MLflow
MLflow是Databricks开源的模型管理工具,能帮你跟踪训练过程中的参数、指标、模型文件,并将模型注册到“模型仓库”(Model Registry)。

# 1. 初始化环境
from pyspark.sql import SparkSession
from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler
from pyspark.ml.classification import LogisticRegression
from pyspark.ml import Pipeline
import mlflow
import mlflow.spark

spark = SparkSession.builder.appName("ChurnTrain").getOrCreate()

# 2. 加载数据
data = spark.read.csv("hdfs://cluster/churn_data.csv", header=True, inferSchema=True)

# 3. 定义特征工程流水线
categorical_cols = ["gender", "contract_type"]
numeric_cols = ["tenure", "monthly_charges"]
label_col = "churn"

# 字符串转索引(StringIndexer)
string_indexers = [StringIndexer(inputCol=col, outputCol=f"{col}_index") for col in categorical_cols]
# 索引转独热编码(OneHotEncoder)
onehot_encoders = [OneHotEncoder(inputCol=f"{col}_index", outputCol=f"{col}_onehot") for col in categorical_cols]
# 组装特征向量(VectorAssembler)
assembler = VectorAssembler(inputCols=[f"{col}_onehot" for col in categorical_cols] + numeric_cols, outputCol="features")
# 逻辑回归分类器
lr = LogisticRegression(featuresCol="features", labelCol=label_col, maxIter=100, regParam=0.01)

# 4. 构建Pipeline并训练
pipeline = Pipeline(stages=string_indexers + onehot_encoders + [assembler, lr])
train_data, test_data = data.randomSplit([0.8, 0.2], seed=42)
model = pipeline.fit(train_data)

# 5. 用MLflow跟踪训练过程
with mlflow.start_run(run_name="churn_train_v1"):
    # 记录参数
    mlflow.log_param("maxIter", lr.getMaxIter())
    mlflow.log_param("regParam", lr.getRegParam())
    # 记录指标(测试集准确率)
    predictions = model.transform(test_data)
    accuracy = predictions.filter(predictions.churn == predictions.prediction).count() / test_data.count()
    mlflow.log_metric("accuracy", accuracy)
    # 保存并注册模型到MLflow
    mlflow.spark.log_model(
        model=model,
        artifact_path="churn_model",
        registered_model_name="churn_prediction"  # 模型仓库中的名称
    )

spark.stop()

步骤2:提交批处理任务
用Spark Submit提交任务,加载MLflow中的生产版本模型,处理新增数据并写入Hive。

# batch_prediction.py
from pyspark.sql import SparkSession
import mlflow.spark

spark = SparkSession.builder.appName("ChurnBatchPred").getOrCreate()

# 1. 加载新增数据(比如前一天的用户数据)
new_data = spark.read.csv("hdfs://cluster/new_churn_data.csv", header=True, inferSchema=True)

# 2. 加载MLflow中的生产版本模型
# 模型URI格式:models:/<模型名称>/<版本标签>
model_uri = "models:/churn_prediction/Production"
model = mlflow.spark.load_model(model_uri)

# 3. 执行预测
predictions = model.transform(new_data)

# 4. 写入Hive表(供业务团队查询)
predictions.write.mode("overwrite").saveAsTable("churn_db.churn_predictions")

spark.stop()

提交命令(YARN集群):

spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --executor-memory 8g \
  --num-executors 10 \
  batch_prediction.py
3. 关键踩坑点
  • 模型版本管理:用MLflow的Production标签标记稳定版本,避免误加载开发中的模型;
  • 数据Schema一致性:新增数据的Schema必须和训练数据一致(比如特征列名、类型),否则transform()会报错;
  • 资源优化:根据数据量调整executor-memory和num-executors,避免OOM或资源闲置。

场景2:在线服务部署——FastAPI + MLeap

适用场景:需要低延迟响应的在线服务(如客服实时查用户流失风险、APP实时推荐)。
核心需求:低延迟(<100ms)、高并发、跨语言兼容。

1. 技术选型
  • 模型转换:MLeap(将Spark模型转成跨平台Bundle)
  • 服务框架:FastAPI(高性能Python Web框架,支持异步)
  • 部署方式:Docker(封装环境,避免依赖冲突)
2. 实战步骤

步骤1:将Spark模型转成MLeap Bundle
MLeap需要一个样例数据来推断模型的输入Schema,所以训练时要保存一份样例数据。

# export_mleap.py
from pyspark.sql import SparkSession
from pyspark.ml import PipelineModel
import mleap.pyspark
from mleap.pyspark.spark_support import SimpleSparkSerializer

# 初始化SparkSession(需启用MLeap扩展)
spark = SparkSession.builder.appName("MLeapExport") \
    .config("spark.jars.packages", "ml.combust.mleap:mleap-spark_2.12:0.21.0") \
    .getOrCreate()
spark.withExtensions(mleap.pyspark.extensions.mlleap_extension())

# 加载训练好的Spark模型
model = PipelineModel.load("hdfs://cluster/churn_model")

# 加载样例数据(训练集的前1条)
sample_data = spark.read.csv("hdfs://cluster/train_sample.csv", header=True, inferSchema=True)

# 导出MLeap Bundle(格式:jar:file:/path/to/bundle.zip)
SimpleSparkSerializer().serializeToBundle(
    model=model,
    bundle_path="jar:file:/tmp/churn_model.zip",
    sample_input=sample_data  # 样例数据用于推断Schema
)

spark.stop()

步骤2:用FastAPI搭建在线服务
MLeap的Python Runtime支持直接加载Bundle,并用Pandas DataFrame做输入——这比Spark的DataFrame轻量得多。

# main.py
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel, Field
from mleap.runtime import SimpleServable
import pandas as pd

# 初始化FastAPI
app = FastAPI(title="用户流失预测API", version="1.0")

# 加载MLeap模型(只加载一次,缓存到内存)
model = SimpleServable.read("file:///tmp/churn_model.zip")

# 定义请求体(和训练时的特征顺序一致)
class ChurnFeatures(BaseModel):
    gender: str = Field(example="Female", description="性别:Male/Female")
    contract_type: str = Field(example="Month-to-month", description="合约类型:Month-to-month/One year/Two year")
    tenure: int = Field(example=12, description="在网时长(月)")
    monthly_charges: float = Field(example=50.0, description="月消费(元)")

# 定义预测接口
@app.post("/predict", summary="预测用户流失风险")
async def predict_churn(features: ChurnFeatures):
    try:
        # 将请求体转换为Pandas DataFrame(MLeap支持Pandas输入)
        df = pd.DataFrame([features.dict()])
        # 执行预测
        predictions = model.transform(df)
        # 提取结果(prediction:0=不流失,1=流失;probability:流失概率)
        result = {
            "user_id": "test_user_123",  # 实际应从请求中获取用户ID
            "churn_prediction": int(predictions["prediction"].iloc[0]),
            "churn_probability": float(predictions["probability"].iloc[0][1])  # 正类(流失)的概率
        }
        return result
    except Exception as e:
        raise HTTPException(status_code=500, detail=f"预测失败:{str(e)}")

# 运行服务:uvicorn main:app --reload --host 0.0.0.0 --port 8000
3. 性能优化
  • 模型缓存:将模型加载到全局变量,避免每次请求重新加载(加载时间约0.5秒);
  • 异步接口:用async def定义接口,支持高并发(FastAPI的异步性能比同步高3~5倍);
  • 特征预处理前移:将StringIndexer、OneHotEncoder等预处理步骤提前到数据写入数据库前,避免在服务中重复计算(比如用户注册时就将性别转成索引);
  • 缓存频繁请求:用Redis缓存高频用户的特征和预测结果(比如10分钟内重复请求的用户,直接返回缓存结果)。
4. Docker部署(可选)

为了避免环境依赖问题,用Docker封装服务:

# Dockerfile
FROM python:3.9-slim

# 设置工作目录
WORKDIR /app

# 安装依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 复制代码和模型
COPY main.py .
COPY churn_model.zip /tmp/churn_model.zip

# 暴露端口
EXPOSE 8000

# 运行服务
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]

requirements.txt:

fastapi==0.95.0
uvicorn==0.22.0
pydantic==1.10.7
pandas==1.5.3
mleap==0.21.0

场景3:实时流处理部署——Structured Streaming + Kafka + MLflow

适用场景:处理实时数据(如Kafka中的用户行为流、IoT设备数据),需要实时预测(比如实时推荐、实时风险预警)。
核心需求:低延迟(<1秒)、高可用(无停机更新模型)、Exactly-Once语义。

1. 技术选型
  • 流处理框架:Spark Structured Streaming(支持SQL-like操作,Exactly-Once)
  • 消息队列:Kafka(高吞吐量、低延迟)
  • 模型更新:MLflow(定期检查模型版本,热更新)
2. 实战步骤

步骤1:定义流处理Pipeline
Structured Streaming将流数据视为“无限增长的DataFrame”,支持用readStream读取Kafka数据,用writeStream写入下游(如Elasticsearch、Redis)。

# streaming_prediction.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col
from pyspark.sql.types import StructType, StringType, IntegerType, FloatType
import mlflow.spark
import time

# 初始化SparkSession
spark = SparkSession.builder.appName("ChurnStreaming") \
    .config("spark.sql.streaming.checkpointLocation", "/tmp/checkpoint") \
    .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.4.0,ml.combust.mleap:mleap-spark_2.12:0.21.0") \
    .getOrCreate()

# Kafka配置
kafka_bootstrap_servers = "kafka:9092"
kafka_topic = "user_behavior"

# 定义Kafka消息的Schema(需和生产者发送的格式一致)
schema = StructType() \
    .add("user_id", StringType()) \
    .add("gender", StringType()) \
    .add("contract_type", StringType()) \
    .add("tenure", IntegerType()) \
    .add("monthly_charges", FloatType())

# 1. 读取Kafka流数据
stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", kafka_bootstrap_servers) \
    .option("subscribe", kafka_topic) \
    .option("startingOffsets", "latest") \
    .load() \
    # 将Kafka的value字段(JSON字符串)解析为StructType
    .select(from_json(col("value").cast("string"), schema).alias("data")) \
    .select("data.*")

# 2. 加载MLflow模型(定期检查新版本)
def load_latest_model():
    """加载MLflow中的生产版本模型"""
    model_uri = "models:/churn_prediction/Production"
    return mlflow.spark.load_model(model_uri)

# 初始加载模型
current_model = load_latest_model()

# 3. 定义foreachBatch函数(处理每个微批)
def process_batch(batch_df, batch_id):
    """
    batch_df:当前微批的DataFrame
    batch_id:微批的唯一ID
    """
    global current_model
    # 每10个微批检查一次模型更新(避免频繁请求MLflow)
    if batch_id % 10 == 0:
        new_model = load_latest_model()
        # 比较模型版本(假设model有version属性,实际需从MLflow元数据获取)
        if new_model.version != current_model.version:
            current_model = new_model
            print(f"更新模型到版本:{current_model.version}")
    
    # 4. 执行预测
    predictions = current_model.transform(batch_df)
    
    # 5. 写入Elasticsearch(供实时监控)
    predictions.write \
        .format("org.elasticsearch.spark.sql") \
        .option("es.nodes", "elasticsearch:9200") \
        .option("es.index", "churn_stream") \
        .option("es.write.operation", "upsert") \
        .mode("append") \
        .save()

# 6. 启动流查询
query = stream_df.writeStream \
    .foreachBatch(process_batch) \
    .trigger(processingTime="10 seconds")  # 每10秒处理一个微批
    .start()

# 等待流查询终止
query.awaitTermination()
3. 关键设计点
  • 模型热更新:用foreachBatch定期检查MLflow的模型版本,不需要重启流任务;
  • Exactly-Once语义:通过checkpointLocation保存流的状态,即使任务重启,也能从上次中断的地方继续;
  • 微批大小调整:用trigger(processingTime="10 seconds")控制微批间隔,平衡延迟和吞吐量(间隔越小,延迟越低,但资源消耗越高)。
4. 测试与验证
  • 用Kafka生产者发送测试数据:
    kafka-console-producer.sh --broker-list kafka:9092 --topic user_behavior
    # 发送JSON数据:{"user_id":"test_1","gender":"Female","contract_type":"Month-to-month","tenure":12,"monthly_charges":50.0}
    
  • 用Kibana查询Elasticsearch中的预测结果,验证是否实时写入。

四、性能优化:从实验室到生产的关键

无论是哪种场景,性能都是生产环境的“生命线”。以下是针对不同场景的优化技巧:

1. 离线批处理优化

  • 减少Shuffle:Shuffle是Spark的性能瓶颈,尽量用 repartition()代替coalesce(),或通过谓词下推(Predicate Pushdown)减少数据传输;
  • 使用Vectorized操作:Spark 3.0+支持Vectorized Parquet读取,比传统的Row-based读取快2~3倍,需开启spark.sql.parquet.enableVectorizedReader=true;
  • 动态资源分配:开启spark.dynamicAllocation.enabled=true,让Spark自动调整executor数量,适应数据量波动。

2. 在线服务优化

  • 用轻量Runtime:优先选择MLeap或ONNX Runtime,避免使用Spark的SparkContext;
  • 特征预处理前移:将StringIndexer、OneHotEncoder等步骤提前到数据写入数据库前,减少服务中的计算量;
  • 异步请求:用FastAPI的async def或Flask的async扩展,支持高并发;
  • 缓存:用Redis缓存高频用户的特征和预测结果,减少重复计算。

3. 实时流处理优化

  • 调整微批大小:根据数据吞吐量调整processingTime(比如每秒1万条数据,设置processingTime="5 seconds");
  • 状态管理:用flatMapGroupsWithState或mapGroupsWithState保存中间状态(比如用户的历史行为),避免重复计算;
  • 并行度调整:设置spark.sql.shuffle.partitions(默认200),根据executor数量调整(比如10个executor,设置为100)。

五、模型监控与运维:生产环境的“安全绳”

模型部署后,你需要监控服务的健康状态和模型的性能衰减,避免“模型上线即死亡”。

1. 服务监控:Prometheus + Grafana

  • 工具:用starlette-prometheus中间件收集FastAPI的 metrics(如QPS、延迟、错误率),用Prometheus存储,Grafana可视化。
  • 代码示例(FastAPI):
    from fastapi import FastAPI
    from starlette_prometheus import PrometheusMiddleware, metrics
    
    app = FastAPI()
    app.add_middleware(PrometheusMiddleware)
    app.add_route("/metrics", metrics)
    
  • Grafana Dashboard:监控以下指标:
    • http_requests_total:总请求数;
    • http_request_duration_seconds_bucket:请求延迟分布;
    • http_request_exceptions_total:错误请求数。

2. 模型性能监控:Evidently AI + MLflow

  • 数据漂移检测:用Evidently AI比较在线数据和训练数据的特征分布(比如性别比例、月消费均值),如果差异超过阈值(如0.5),触发报警。
    import pandas as pd
    from evidently.dashboard import Dashboard
    from evidently.tabs import DataDriftTab
    
    # 加载训练数据(基准)
    train_data = pd.read_csv("train_data.csv")
    # 加载在线数据(最近1小时)
    online_data = pd.read_csv("online_data.csv")
    
    # 生成数据漂移报告
    dashboard = Dashboard(tabs=[DataDriftTab()])
    dashboard.calculate(train_data, online_data)
    dashboard.save("data_drift_report.html")
    
  • 模型性能监控:用MLflow Tracking记录在线预测的准确率、F1-score等指标,定期对比训练时的指标,如果下降超过5%,触发模型重新训练。

3. 报警与自动化

  • 报警渠道:用Alertmanager将Prometheus的报警发送到Slack、邮件或企业微信;
  • 自动化运维:用Airflow或Prefect定时运行数据漂移检测脚本,若检测到漂移,自动触发模型重新训练和部署。

六、工具与资源推荐

  • 模型序列化:MLeap(Spark专属)、ONNX(多框架);
  • 部署框架:FastAPI(Python)、Spring Boot(Java)、Gin(Go);
  • 版本管理:MLflow(模型)、DVC(数据);
  • 监控工具:Prometheus(服务)、Grafana(可视化)、Evidently AI(数据漂移);
  • 流处理:Spark Structured Streaming、Flink;
  • 资源:
    • MLeap文档:https://mleap-docs.combust.ml/
    • MLflow文档:https://mlflow.org/docs/latest/index.html
    • FastAPI文档:https://fastapi.tiangolo.com/

七、未来趋势:Spark ML模型部署的下一站

随着大模型(LLM)和Serverless的普及,Spark ML模型的部署将向更自动化、更弹性、更集成的方向发展:

1. Serverless部署

用AWS Lambda、阿里云函数计算部署模型服务,按需付费,自动伸缩,无需管理服务器。比如用FastAPI + MLeap打包成Lambda函数,处理高并发请求。

2. 模型即服务(MaaS)

用SageMaker、MLflow Serving等平台,直接部署Spark模型,自动处理负载均衡、缩放、监控。比如SageMaker的SparkMLModel类型,支持直接加载Spark原生模型。

3. 实时特征平台

用Feast、Tecton等实时特征平台,将Spark处理的特征存储为在线特征表,模型服务直接读取在线特征,避免重复计算。比如用户的在网时长,由Spark批处理每天更新到Feast,模型服务实时读取。

4. LLM与Spark的结合

用Spark处理大规模文本特征(如TF-IDF、Word2Vec),然后将特征喂给LLM(如ChatGPT、Llama)做预测。比如用户评论的情感分析:用Spark提取评论的关键词特征,再用LLM预测情感倾向。

八、总结:从训练到部署的闭环

Spark ML模型的部署,本质是将分布式的、离线的模型,适配到生产环境的各种场景。关键是:

  1. 选对场景:根据业务需求选择离线批处理、在线服务或实时流计算;
  2. 选对工具:用MLeap解决跨平台问题,用MLflow解决版本管理,用FastAPI解决在线服务;
  3. 注重优化:从资源、代码、架构三个层面优化性能;
  4. 监控运维:用工具监控服务和模型的健康状态,避免“黑盒”运行。

最后,记住:模型部署不是终点,而是数据价值的起点。只有让模型真正运行在生产环境中,解决业务问题,才能体现机器学习的价值。

如果你在部署过程中遇到问题,欢迎在评论区留言——我会第一时间解答!

附录:代码仓库
所有示例代码已上传至GitHub:https://github.com/your-name/spark-ml-deployment-demo
包含:

  • 模型训练与MLflow注册代码;
  • MLeap模型导出代码;
  • FastAPI在线服务代码;
  • Structured Streaming流处理代码;
  • Dockerfile与部署脚本。
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐