Spark Structured Streaming作业遇java.io.FileNotFoundException求助
Let's break down your problem and fix it step by step—this is a common pitfall when running Spark Streaming on YARN with local file paths.
The Root Cause
You're running your job in YARN cluster mode (--master yarn), but using a file:// path to read local CSV files. Here's why this causes random errors:
- In YARN mode, Spark executors run on different nodes across your cluster. Each executor only has access to its own node's local filesystem.
- When you add files to
/home/pardeep/file2on your driver node, only the executor running on the driver node can see those files. Executors on other nodes will throw aFileNotFoundExceptionwhen they try to process those files. - This explains the randomness: sometimes the driver's executor picks up the file (no error), other times a remote executor tries to read it (error).
Recommended Solution: Move Source Files to HDFS
The most reliable fix is to use HDFS as your source storage, since it's a shared filesystem accessible to all nodes in the cluster.
Step 1: Upload Your Source Directory to HDFS
First, create a directory in HDFS and move your existing CSV files there:
hdfs dfs -mkdir -p /user/pardeep/csv_source hdfs dfs -put /home/pardeep/file2/* /user/pardeep/csv_source/
Step 2: Update Your Streaming Code
Modify your load path to use the HDFS URL instead of the local file path, and add a safety option to skip transient missing files:
from pyspark.sql.types import StructType sch = StructType(...) spark.readStream \ .format("csv") \ .schema(sch) \ .option("header", True) \ .option("delimiter", ',') \ .option("ignoreMissingFiles", True) # Skip files that can't be found temporarily .load("hdfs:///user/pardeep/csv_source") # Use HDFS path here .writeStream \ .format("parquet") \ .outputMode("append") \ .trigger(processingTime='10 seconds') \ .option("path","<Hdfs_path>") \ .option("checkpointLocation","<Checkpoint Loc>") \ .start()
Step 3: Adjust Your File Movement Job
Update the job that adds new files to use HDFS atomic moves. Instead of moving files to the local directory, upload them to an HDFS temp directory first, then move them to the source directory (HDFS mv is atomic, so Spark won't see partial files):
# Example: Upload file to temp directory first hdfs dfs -put /path/to/new_file.csv /user/pardeep/tmp/ # Atomically move to source directory hdfs dfs -mv /user/pardeep/tmp/new_file.csv /user/pardeep/csv_source/
Alternative: Using Local Files (Not Recommended)
If you absolutely must use local files, you need to ensure the source directory is shared across all cluster nodes:
- Mount a shared storage system (like NFS) to
/home/pardeep/file2on every node in your cluster. This way, all executors can access the same files. - Note: This adds operational overhead (maintaining the shared storage) and is less reliable than HDFS for streaming workloads.
Key Optimizations to Prevent Future Errors
- Avoid partial files: Never write directly to the streaming source directory. Always use a temp directory and atomic moves to ensure files are fully written before Spark sees them.
- Ignore missing files: The
ignoreMissingFilesoption tells Spark to skip files that can't be found, preventing your entire streaming job from failing due to transient file system issues.
内容的提问来源于stack exchange,提问作者pardeep garg

