如何将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
相关产品推荐
相关产品推荐

