PyCharm运行PySpark时出现“Python worker did not connect back in time”错误
嘿,我太懂你现在的困扰啦!作为刚上手Spark Python单元测试的新手,这种“基础代码能跑,一用转换就崩”的情况真的很常见,咱们结合你的环境(Windows+PyCharm+Anaconda Python3.5+Spark2.1.0)来一步步排查解决~
先梳理核心问题
你提到直接并行化列表后执行collect(哦对了,你写的output=rdd.collect其实少了括号,应该是rdd.collect(),不过你说能运行,估计是打印了方法对象)没问题,但一用转换操作(比如map/reduceByKey/flatMap这类)就出问题,大概率是这几个方向的原因:
1. Windows环境必备的WinUtils配置缺失
Spark依赖Hadoop的WinUtils工具来处理文件系统相关操作,Windows系统默认没有这个,执行转换操作时需要启动Executor,这时候就会因为找不到WinUtils报错。这是Windows跑Spark的头号坑!
解决步骤:
- 下载和你的Spark版本匹配的WinUtils(Spark2.1.0对应Hadoop2.7.x版本的WinUtils)
- 新建一个文件夹(比如
C:\hadoop),在里面再建bin子文件夹,把WinUtils.exe放进去 - 配置系统环境变量:
- 新增
HADOOP_HOME变量,值设为C:\hadoop - 把
%HADOOP_HOME%\bin添加到系统PATH里
- 新增
- 重启PyCharm让环境变量生效
2. SparkContext初始化不规范
如果你的SparkContext没有正确设置本地模式,执行转换操作时无法启动本地计算集群,就会报错。要确保初始化代码是这样的:
from pyspark import SparkContext, SparkConf # 强制设置为本地模式,[*]表示用所有可用CPU核心 conf = SparkConf().setAppName("WordCountTest").setMaster("local[*]") sc = SparkContext(conf=conf)
不要省略setMaster("local[*]"),否则Spark会尝试连接集群,而你本地没有集群的话就会失败。
3. 依赖包版本冲突(比如Py4J)
Spark和Python之间靠Py4J通信,Spark2.1.0对应的Py4J版本是0.10.4,如果你的Anaconda环境里Py4J版本不对,就会出现兼容性问题。
检查版本的命令:
pip show py4j
如果版本不符,卸载重装:
pip uninstall py4j -y pip install py4j==0.10.4
4. 转换操作的语法错误
有时候可能是你写的转换逻辑有问题,比如lambda表达式写错、引用了Spark无法序列化的对象(比如自定义类的实例,Spark2.1.0对Python的序列化支持有限)。
给你一个能正常运行的WordCount完整示例,你可以对比着看:
from pyspark import SparkContext, SparkConf # 初始化SparkContext conf = SparkConf().setAppName("WordCountTest").setMaster("local[*]") sc = SparkContext(conf=conf) # 测试数据 mylist = ["the", "earth", "revolves", "around", "sun", "the", "sun"] rdd = sc.parallelize(mylist) # 转换+行动操作 word_counts_rdd = rdd.map(lambda word: (word, 1)).reduceByKey(lambda a, b: a + b) output = word_counts_rdd.collect() # 打印结果 for word, count in output: print(f"{word}: {count}") # 记得关闭SparkContext,避免资源占用 sc.stop()
按照上面的步骤排查,应该就能解决你的问题啦!如果还有具体的报错信息,也可以针对性地再调整~
内容的提问来源于stack exchange,提问作者harpreet kaur

