Spark本地调用RDD.distinct().collect()报PySparkRuntimeError属性错误
本地无集群环境下使用PySpark 3.4.2(Python 3.10.12)编写脚本,目标是提取数据中license键的唯一值。单独执行licenses_collected = rdd.map(lambda x: x["license"]).collect()或unique_licenses = rdd.map(lambda x: x["license"]).distinct()均正常,但执行unique_licenses_collected = rdd.map(lambda x: x["license"]).distinct().collect()时抛出以下错误:
AttributeError: Can't get attribute 'PySparkRuntimeError' on <module 'pyspark.errors.exceptions.base' from '/opt/spark/python/lib/pyspark.zip/pyspark/errors/exceptions/base.py'>
脚本代码如下:
from pyspark.sql import SparkSession spark = SparkSession.builder.master("local").\ appName("prima applicazione spark").\ config("spark.ui.port", "4040").\ getOrCreate() df_spark = spark.read.json(path_json) rdd = df_spark.rdd rdd_collected = rdd.collect() """Get the first two records""" for line in rdd.take(2): print(line) print("\n\n") """ Get the name of the licenses""" licenses_collected = rdd.map(lambda x: x["license"]).collect() unique_licenses = rdd.map(lambda x: x["license"]).distinct() unique_licenses_collected = rdd.map(lambda x: x["license"]).distinct().collect()
1. 确保PySpark版本完全匹配
该错误常因pip安装的PySpark包与本地Spark二进制文件版本不一致导致。执行以下命令检查并修复:
# 查看当前pip安装的PySpark版本 pip show pyspark # 卸载现有版本并重装指定版本 pip uninstall pyspark -y pip install pyspark==3.4.2
2. 统一Python环境路径
在脚本开头添加环境变量,指定Spark Worker使用与脚本相同的Python解释器(替换为你的Python3.10实际路径):
import os # 替换为你的Python3.10执行路径 os.environ["PYSPARK_PYTHON"] = "/usr/bin/python3.10" os.environ["PYSPARK_DRIVER_PYTHON"] = "/usr/bin/python3.10" from pyspark.sql import SparkSession # 后续SparkSession创建代码不变
3. 改用DataFrame API替代RDD
既然已经读取生成了DataFrame,无需转换为RDD操作,DataFrame API更高效且兼容性更好:
# 直接用DataFrame提取唯一license值 unique_licenses_collected = df_spark.select("license").distinct().collect()
4. 清理Spark缓存
若存在缓存数据干扰,可在创建SparkSession后清理缓存:
spark = SparkSession.builder.master("local").\ appName("prima applicazione spark").\ config("spark.ui.port", "4040").\ getOrCreate() # 清理缓存 spark.catalog.clearCache()
内容的提问来源于stack exchange,提问作者Tms91

