使用range创建PySpark DataFrame触发两类报错的原因咨询
使用range生成序列创建PySpark DataFrame的报错修复
两段代码的错误点
- 第一段pandas转Spark的代码:
- pandas DataFrame列赋值不规范:
"CSIRO Adjusted Sea Level"列传入单个浮点数0.0而非等长序列,虽然pandas会自动广播填充,但属于隐式操作容易触发兼容问题。 - 抛出的
Can't get attribute '_fill_function'是版本兼容问题:Spark 3.2.1内置的cloudpickle版本和当前环境安装的pandas版本不匹配,序列化pandas对象时失败。
- pandas DataFrame列赋值不规范:
- 第二段列表构造的代码:
已经构造好了包含38个元组的行列表lista,传入createDataFrame时额外套了一层方括号,等价于传入了1行数据,这行数据本身是长度为38的列表,但指定的表结构只有2个字段,因此触发长度不匹配报错。
可直接运行的正确写法
不需要依赖pandas,直接用列表构造DataFrame即可,完全绕开序列化兼容问题:
# 生成行数据:每个元组对应DataFrame的一行 data_list = [(year, 0.0) for year in range(2013, 2051)] # 直接传入行数据列表,不要额外嵌套方括号,指定字段名和类型 df = spark.createDataFrame( data_list, schema="Year INT, `CSIRO Adjusted Sea Level` DOUBLE" ) # 验证结果 df.show(5)
如果一定要走pandas转Spark的路径,先修正pandas的列赋值逻辑,若仍报序列化错误,将pandas版本降级到1.5.x适配Spark 3.2.1即可:
import pandas as pd year_range = list(range(2013, 2051)) pdf = pd.DataFrame( { "Year": year_range, # 传入和Year列等长的序列,不依赖pandas自动广播 "CSIRO Adjusted Sea Level": [0.0] * len(year_range) } ) df_pyspark = spark.createDataFrame(pdf) df_pyspark.show(5)
内容的提问来源于stack exchange,提问作者enas dyo
相关产品推荐
相关产品推荐

