启用Arrow转换Spark DataFrame至Pandas时遇Py4JError问题
问题:启用Arrow加速Spark转Pandas时出现Py4JError错误
我尝试启用以下两个配置项来加速Spark DataFrame转Pandas DataFrame:
'spark.sql.execution.arrow.pyspark.enabled' 'spark.sql.execution.arrow.pyspark.fallback.enabled'
但执行转换时出现如下错误:
File /opt/conda/envs/python385/lib/python3.8/site-packages/pyspark/sql/pandas/conversion.py:108, in PandasConversionMixin.toPandas(self) 106 # Rename columns to avoid duplicated column names. 107 tmp_column_names = ['col_{}'.format(i) for i in range(len(self.columns))] --> 108 self_destruct = self.sql_ctx._conf.arrowPySparkSelfDestructEnabled() 109 batches = self.toDF(*tmp_column_names)._collect_as_arrow( 110 split_batches=self_destruct) 111 if len(batches) > 0: Py4JError: An error occurred while calling o1723.arrowPySparkSelfDestructEnabled. Trace: py4j.Py4JException: Method arrowPySparkSelfDestructEnabled([]) does not exist
已通过conda-forge安装pyarrow,求解决办法?
解决方案
- 核心原因:PySpark与集群端Spark版本不兼容。
arrowPySparkSelfDestructEnabled是较新版本PySpark新增的配置方法,若集群Scala端Spark版本低于本地PySpark版本,就会出现找不到该方法的错误。 - 具体修复步骤:
- 核对本地PySpark版本与集群Spark版本,确保二者完全一致(比如集群是Spark 3.2.1,本地PySpark也必须是3.2.1)
- 卸载不匹配的PySpark,重新安装对应版本:
pip uninstall pyspark -y pip install pyspark==<你的集群Spark版本号> - 验证pyarrow版本兼容性:不同PySpark版本对应不同的pyarrow版本要求,比如Spark 3.2.x适配pyarrow 6.0.0,Spark 3.3.x适配pyarrow 7.0.0及以上,可通过PySpark官方文档确认对应关系
- 重启Spark会话,重新配置Arrow参数后执行转换操作
内容的提问来源于stack exchange,提问作者Bhaskar
相关产品推荐
相关产品推荐

