Pandas2.0.0转PySpark DataFrame遇AttributeError:无iteritems属性怎么办?
解决Spark 3.3.2与Pandas 2.0.0兼容问题:spark.createDataFrame()报错AttributeError
问题描述
使用spark.createDataFrame()将Pandas DataFrame转换为PySpark DataFrame时,在Pandas 2.0.0及以上版本会触发以下错误:
AttributeError Traceback (most recent call last) File <command-2209449931455530>:64 61 df_train_test_p.loc[df_train_test_p.is_train=='N','preds']=preds_test 63 # save the original table and predictions into spark dataframe ---> 64 df_test = spark.createDataFrame(df_train_test_p.loc[df_train_test_p.is_train=='N']) 65 df_results = df_results.union(df_test) 67 # saving all relevant data File /databricks/spark/python/pyspark/instrumentation_utils.py:48, in _wrap_function.<locals>.wrapper(*args, **kwargs) 46 start = time.perf_counter() 47 try: ---> 48 res = func(*args, **kwargs) 49 logger.log_success( 50 module_name, class_name, function_name, time.perf_counter() - start, signature 51 ) 52 return res File /databricks/spark/python/pyspark/sql/session.py:1211, in SparkSession.createDataFrame(self, data, schema, samplingRatio, verifySchema) 1207 data = pd.DataFrame(data, columns=column_names) 1209 if has_pandas and isinstance(data, pd.DataFrame): 1210 # Create a DataFrame from pandas DataFrame. -> 1211 return super(SparkSession, self).createDataFrame( # type: ignore[call-overload] 1212 data, schema, samplingRatio, verifySchema 1213 ) 1214 return self._create_dataframe( 1215 data, schema, samplingRatio, verifySchema # type: ignore[arg-type] 1216 ) File /databricks/spark/python/pyspark/sql/pandas/conversion.py:478, in SparkConversionMixin.createDataFrame(self, data, schema, samplingRatio, verifySchema) 476 warn(msg) 477 raise --> 478 converted_data = self._convert_from_pandas(data, schema, timezone) 479 return self._create_dataframe(converted_data, schema, samplingRatio, verifySchema) File /databricks/spark/python/pyspark/sql/pandas/conversion.py:516, in SparkConversionMixin._convert_from_pandas(self, pdf, schema, timezone) 514 else: 515 should_localize = not is_timestamp_ntz_preferred() --> 516 for column, series in pdf.iteritems(): 517 s = series 518 if should_localize and is_datetime64tz_dtype(s.dtype) and s.dt.tz is not None: File /local_disk0/.ephemeral_nfs/envs/pythonEnv-fefe10af-04b7-4277-b395-2f16b77bd90b/lib/python3.9/site-packages/pandas/core/generic.py:5981, in NDFrame.__getattr__(self, name) 5974 if ( 5975 name not in self._internal_names_set 5976 and name not in self._metadata 5977 and name not in self._accessors 5978 and self._info_axis._can_hold_identifiers_and_holds_name(name) 5979 ): 5980 return self[name] -> 5981 return object.__getattribute__(self, name) AttributeError: 'DataFrame' object has no attribute 'iteritems'
可复现问题的示例代码:
import pandas as pd from pyspark.sql import SparkSession # 创建示例Pandas DataFrame data = {'name': ['John', 'Mike', 'Sara', 'Adam'], 'age': [25, 30, 18, 40]} df_pandas = pd.DataFrame(data) # 转换为PySpark DataFrame spark = SparkSession.builder.appName('pandasToSpark').getOrCreate() df_spark = spark.createDataFrame(df_pandas) # 展示结果 df_spark.show()
原因分析
Pandas 2.0.0正式移除了iteritems()方法,统一使用items()替代;而Spark 3.3.2的pyspark.sql.pandas.conversion模块中仍在调用pdf.iteritems()遍历DataFrame列,导致属性查找失败。Spark官方在3.4.0及以上版本已经修复了该问题,将代码中的iteritems()替换为items()。
解决方案
方法1:升级Spark版本到3.4.0及以上
这是最彻底的解决方案,升级后Spark会自动兼容Pandas 2.0.0+的API变更,无需修改现有代码即可正常使用spark.createDataFrame()。
方法2:临时兼容处理(不推荐长期使用)
在转换前给Pandas DataFrame添加iteritems的别名,指向items()方法,临时绕过兼容性问题:
import pandas as pd from pyspark.sql import SparkSession # 添加临时别名,兼容Spark 3.3.2的调用 pd.DataFrame.iteritems = pd.DataFrame.items data = {'name': ['John', 'Mike', 'Sara', 'Adam'], 'age': [25, 30, 18, 40]} df_pandas = pd.DataFrame(data) spark = SparkSession.builder.appName('pandasToSpark').getOrCreate() df_spark = spark.createDataFrame(df_pandas) df_spark.show()
⚠️ 注意:该方法属于临时workaround,可能与其他依赖Pandas的代码产生冲突,仅建议在无法升级Spark时短期使用。
方法3:使用替代转换方式
可以绕过spark.createDataFrame()对Pandas API的直接调用,改用其他方式转换:
方式A:转换为列表后创建DataFrame
import pandas as pd from pyspark.sql import SparkSession data = {'name': ['John', 'Mike', 'Sara', 'Adam'], 'age': [25, 30, 18, 40]} df_pandas = pd.DataFrame(data) spark = SparkSession.builder.appName('pandasToSpark').getOrCreate() # 将Pandas DataFrame转换为列表,指定列名作为schema df_spark = spark.createDataFrame(df_pandas.values.tolist(), schema=df_pandas.columns.tolist()) df_spark.show()
方式B:使用PySpark Pandas API(如果环境支持)
import pandas as pd import pyspark.pandas as ps from pyspark.sql import SparkSession data = {'name': ['John', 'Mike', 'Sara', 'Adam'], 'age': [25, 30, 18, 40]} df_pandas = pd.DataFrame(data) spark = SparkSession.builder.appName('pandasToSpark').getOrCreate() # 先转为PySpark Pandas DataFrame,再转换为普通PySpark DataFrame df_spark = ps.from_pandas(df_pandas).to_spark() df_spark.show()
内容的提问来源于stack exchange,提问作者shiftyscales
相关产品推荐
相关产品推荐

