能否在不转换格式的情况下对pandas dataframe应用PySpark pipeline?
问题解答
能否直接在pandas DataFrame上应用PySpark Pipeline?
不行。PySpark Pipeline的所有组件(缺失值插补器、机器学习模型等)均基于Spark分布式DataFrame抽象设计,底层依赖Spark的RDD执行引擎和分布式数据结构,完全不兼容单节点内存存储的pandas DataFrame,不做格式转换无法直接调用。
高效生成预测的实现方案
最优方案是借助PyArrow加速pandas和Spark DataFrame的互转,整个流程几乎没有额外的序列化开销,操作步骤如下:
- 第一步:初始化SparkSession时开启PyArrow优化配置,大幅降低跨框架数据转换的性能损耗
- 第二步:将pandas DataFrame转换为Spark DataFrame
- 第三步:直接调用已训练好的PySpark PipelineModel的
transform方法生成预测结果 - 第四步:将带预测结果的Spark DataFrame转回pandas DataFrame即可
代码示例
from pyspark.sql import SparkSession from pyspark.ml import PipelineModel # 初始化SparkSession并开启PyArrow优化 spark = SparkSession.builder \ .appName("pandas_inference") \ .config("spark.sql.execution.arrow.pyspark.enabled", "true") \ .getOrCreate() # 加载你已训练好的PySpark Pipeline模型 trained_pipeline = PipelineModel.load("你的模型存储路径") # 你的待预测pandas DataFrame,变量名为pd_df # pd_df = 你的pandas数据 # pandas转Spark DataFrame spark_input_df = spark.createDataFrame(pd_df) # 应用pipeline生成预测 spark_pred_df = trained_pipeline.transform(spark_input_df) # 预测结果转回pandas DataFrame pd_pred_df = spark_pred_df.toPandas()
注意事项
- 仅需转换数据结构,无需重新拟合Pipeline:已训练完成的
PipelineModel中已经存储了训练阶段计算好的缺失值填充规则、标准化参数、模型权重等信息,直接调用transform即可 - 若待预测的pandas DataFrame体积超过单节点内存,建议先将数据存储为Parquet格式文件,再用Spark读取文件后执行预测,避免单节点内存溢出
- 待预测的pandas DataFrame的列名、数据类型必须和Pipeline训练时的输入数据完全一致,否则会触发transform阶段的格式校验报错
内容的提问来源于stack exchange,提问作者Prabhat Mangina
相关产品推荐
相关产品推荐

