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

如何将LGBM模型分发到PySpark Executor,分块处理数据集并转Pandas DF

如何在Spark Executor上分布式运行LGBM模型推理

完全可以将LGBM模型分发到Executor,让每个Executor独立处理自己负责的数据集分区,避免数据回传驱动端(Driver),大幅提升超大规模数据集的推理效率。下面是两种落地实现方案:

方法一:使用mapPartitions自定义分区处理

核心是通过广播序列化后的模型,让每个Executor仅加载一次模型,再对本地分区数据做Pandas转换与推理:

import pickle
import pandas as pd
from pyspark.sql import Row
from lightgbm import LGBMClassifier

# 假设已训练完成的LGBM模型
trained_model = LGBMClassifier()
# 序列化模型(若pickle兼容问题,可替换为joblib.dumps)
serialized_model = pickle.dumps(trained_model)
# 广播模型到所有Executor(每个Executor仅接收一次)
broadcast_model = spark.sparkContext.broadcast(serialized_model)

def process_partition(partition_data):
    # 反序列化模型(每个Executor仅执行一次)
    model = pickle.loads(broadcast_model.value)
    # 将分区数据转为Pandas DataFrame
    pd_df = pd.DataFrame(partition_data)
    # 运行模型推理
    predictions = model.predict(pd_df)
    # 拼接原数据与预测结果,返回Spark Row格式
    for idx, row in pd_df.iterrows():
        yield Row(**row.to_dict(), prediction=predictions[idx])

# 对目标分布式DataFrame应用分区处理,生成结果DataFrame
result_spark_df = target_spark_df.mapPartitions(process_partition).toDF()

方法二:使用Spark Pandas UDF(原生高效方案)

利用Spark官方提供的Pandas UDF,简化分区处理逻辑,代码更简洁:

import pickle
import pandas as pd
from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import IntegerType  # 根据实际预测结果类型调整

serialized_model = pickle.dumps(trained_model)
broadcast_model = spark.sparkContext.broadcast(serialized_model)

# 定义迭代器式Pandas UDF,每个Executor仅加载一次模型
@pandas_udf(IntegerType())
def lgbm_predict(iterator):
    model = pickle.loads(broadcast_model.value)
    for pd_df in iterator:
        yield model.predict(pd_df)

# 直接在Spark DataFrame上调用UDF生成预测列
result_spark_df = target_spark_df.withColumn(
    "prediction",
    lgbm_predict(*target_spark_df.columns)  # 传入所有特征列
)

关键注意事项

  • 版本一致性:确保Driver与所有Executor的LGBM、Pandas、Spark版本完全一致,避免序列化/反序列化失败。
  • 资源配置:根据分区数据大小调整spark.executor.memory,避免Pandas转换时内存溢出;同时合理设置分区数(spark.sql.shuffle.partitions),均衡Executor负载。
  • 数据兼容性:Spark DataFrame的列名、数据类型必须与模型训练时的输入完全匹配,否则会触发推理报错。

内容的提问来源于stack exchange,提问作者Vinícius Matheus Olivieri

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 03:58:30