You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark使用saveAsTextFile保存RDD至文本文件失败求助

Spark RDD saveAsTextFile 抛出Py4JJavaError问题

我正在编写Spark程序,从airports.text文件读取机场数据,筛选出位于美国的机场,将其名称与城市名称输出至out/airports_in_usa.text文件,但调用saveAsTextFile保存RDD对象时出现Py4JJavaError错误,无法保存查看RDD数据。已创建包含Utils.COMMA_DELIMITER的.py模块,并将其所在目录添加至PYTHONPATH以实现导入。

代码实现

from pyspark import SparkContext, SparkConf
import Utils     # 自定义模块

# 定义分割函数
def splitComma(line: str):    
    splits = Utils.COMMA_DELIMITER.split(line)    
    return "{}, {}".format(splits[1], splits[2])

conf = SparkConf().setAppName("airports").setMaster("local[*]")   
sc = SparkContext(conf = conf)

# 读取数据集生成RDD
airports = sc.textFile("airports.text")

# 筛选美国的机场
airportsInUSA = airports.filter(lambda line : Utils.COMMA_DELIMITER.split(line)[3] == "\"United States\"")  

# 映射提取名称和城市
airportsNameAndCityNames = airportsInUSA.map(splitComma)

# 保存结果
airportsNameAndCityNames.saveAsTextFile("out/airports_in_usa.text")

错误信息

Py4JJavaError                             Traceback (most recent call last)
c:\Users\user\Documents\Python Programming\PySpark_for_Big_Data_3.ipynb Cell 18 line 2
      1 # Save output in a new text file
----> 2 airportsNameAndCityNames.saveAsTextFile("out/airports_in_usa.text")
      3 # displays an error

File C:\spark\spark-3.4.2-bin-hadoop3\python\pyspark\rdd.py:3406, in RDD.saveAsTextFile(self, path, compressionCodecClass)
   3404     keyed._jrdd.map(self.ctx._jvm.BytesToString()).saveAsTextFile(path, compressionCodec)
   3405 else:
-> 3406     keyed._jrdd.map(self.ctx._jvm.BytesToString()).saveAsTextFile(path)

File C:\spark\spark-3.4.2-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\java_gateway.py:1322, in JavaMember.__call__(self, *args)
   1316 command = proto.CALL_COMMAND_NAME +\
   1317     self.command_header +\
   1318     args_command +\
   1319     proto.END_COMMAND_PART
   1321 answer = self.gateway_client.send_command(command)
-> 1322 return_value = get_return_value(
   1323     answer, self.gateway_client, self.target_id, self.name)
   1325 for temp_arg in temp_args:
   1326     if hasattr(temp_arg, "_detach"):

File C:\spark\spark-3.4.2-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\protocol.py:326, in get_return_value(answer, gateway_client, target_id, name)
    324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client)
...
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
    at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
    at java.base/java.lang.Thread.run(Thread.java:842)
Output is truncated. View as a scrollable element or open in a text editor. Adjust cell output settings...

软件配置

  • Python 3.11.8
  • PySpark 3.4.2
  • Spark 3.4.2
  • Java 17
  • VS Code IDE

内容的提问来源于stack exchange,提问作者Uche Kalu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 00:27:09