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
相关产品推荐
相关产品推荐

