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

AWS PySpark集群跨节点添加文件方法及读取报错解决

解决Spark读取本地文件时的「Input Path不存在」问题

Hey there! I see you're new to Spark and hitting a super common snag when trying to read a local file from your master node—let's get this sorted out for you right away.

First, let's zero in on the root cause from your error trace: the key line is org.apache.hadoop.mapred.InvalidInputException: Input path does not exist: file:/home/ec2-user/PR_DATA_35.csv. This happens because Spark's worker nodes don't have access to the file that's only stored on your master instance. You’re totally right about the two main fixes—using HDFS or copying the file to all cluster nodes—so here are the exact commands for both approaches:

方案1:将文件上传到HDFS(推荐生产环境使用)

HDFS is designed for distributed storage, so this is the most scalable and maintainable option long-term.

  • First, create a directory in HDFS (skip this step if the directory already exists):
    hdfs dfs -mkdir -p /user/ec2-user
    
  • Upload your local file to the HDFS directory:
    hdfs dfs -put /home/ec2-user/PR_DATA_35.csv /user/ec2-user/
    
  • Now you can read the file in Spark using the HDFS path:
    # Use the absolute HDFS path if your default filesystem is configured
    rdd = sc.textFile("/user/ec2-user/PR_DATA_35.csv")
    # Or specify the namenode address explicitly if needed
    # rdd = sc.textFile("hdfs://your-namenode-host:9000/user/ec2-user/PR_DATA_35.csv")
    

方案2:复制文件到所有集群节点(适合测试/小文件场景)

If you're just testing and don't want to set up HDFS yet, copying the file to every worker node works too.

  • If you have a list of worker nodes saved in a file (e.g., worker_nodes.txt, with one hostname/IP per line), use a loop to copy the file to each node:
    while read worker_node; do
      scp /home/ec2-user/PR_DATA_35.csv ec2-user@$worker_node:/home/ec2-user/
    done < worker_nodes.txt
    
  • For a small number of workers, you can copy manually one by one:
    scp /home/ec2-user/PR_DATA_35.csv ec2-user@worker1:/home/ec2-user/
    scp /home/ec2-user/PR_DATA_35.csv ec2-user@worker2:/home/ec2-user/
    # Repeat this command for every worker node in your cluster
    

Quick Verification Tip

Before running your Spark code, double-check the file is accessible:

  • For HDFS: Run hdfs dfs -ls /user/ec2-user/ to confirm the file exists.
  • For worker nodes: SSH into one worker and run ls /home/ec2-user/ to verify the file is there.

Here’s the full error trace you provided for reference:

--------------------------------------------------------------------------- Py4JJavaError Traceback (most recent call
last) in ()
----> 1 ncols = rdd.first().features.size # number of columns (no class) of the dataset

/home/ec2-user/spark/python/pyspark/rdd.pyc in first(self) 1359
ValueError: RDD is empty 1360 """
-> 1361 rs = self.take(1) 1362 if rs: 1363 return rs[0]

/home/ec2-user/spark/python/pyspark/rdd.pyc in take(self, num) 1311
""" 1312 items = []
-> 1313 totalParts = self.getNumPartitions() 1314 partsScanned = 0 1315

/home/ec2-user/spark/python/pyspark/rdd.pyc in getNumPartitions(self)
2438 2439 def getNumPartitions(self):
-> 2440 return self._prev_jrdd.partitions().size() 2441 2442 @property

/home/ec2-user/spark/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py
in call(self, *args) 1131 answer =
self.gateway_client.send_command(command) 1132 return_value
= get_return_value(
-> 1133 answer, self.gateway_client, self.target_id, self.name) 1134 1135 for temp_arg in temp_args:

/home/ec2-user/spark/python/pyspark/sql/utils.pyc in deco(*a, **kw)
61 def deco(*a, **kw):
62 try:
---> 63 return f(*a, **kw)
64 except py4j.protocol.Py4JJavaError as e:
65 s = e.java_exception.toString()

/home/ec2-user/spark/python/lib/py4j-0.10.4-src.zip/py4j/protocol.py
in get_return_value(answer, gateway_client, target_id, name)
317 raise Py4JJavaError(
318 "An error occurred while calling {0}{1}{2}.\n".
--> 319 format(target_id, ".", name), value)
320 else:
321 raise Py4JError(

Py4JJavaError: An error occurred while calling o122.partitions. :
org.apache.hadoop.mapred.InvalidInputException: Input path does not
exist: file:/home/ec2-user/PR_DATA_35.csv at
org.apache.hadoop.mapred.FileInputFormat.singleThreadedListStatus(FileInputFormat.java:285)
at
org.apache.hadoop.mapred.FileInputFormat.listStatus(FileInputFormat.java:228)
at
org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:313)
at org.apache.spark.rdd.HadoopRDD.getPartitions(HadoopRDD.scala:194)
at
org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:252)
at
org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:250)
at scala.Option.getOrElse(Option.scala:121) at
org.apache.spark.rdd.RDD.partitions(RDD.scala:250) at
org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35)
at
org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:252)
at
org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:250)
at scala.Option.getOrElse(Option.scala:121) at
org.apache.spark.rdd.RDD.partitions(RDD.scala:250) at
org.apache.spark.api.java.JavaRDDLike$class.partitions(JavaRDDLike.scala:61)
at
org.apache.spark.api.java.AbstractJavaRDDLike.partitions(JavaRDDLike.scala:45)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at
sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at
sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498) at
py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at
py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at
py4j.Gateway.invoke(Gateway.java:280) at
py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
at py4j.commands.CallCommand.execute(CallCommand.java:79) at
py4j.GatewayConnection.run(GatewayConnection.java:214) at
java.lang.Thread.run(Thread.java:748)

内容的提问来源于stack exchange,提问作者Rajat Handa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:24:05