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

Spark Structured Streaming作业遇java.io.FileNotFoundException求助

Fixing FileNotFoundException in Spark Structured Streaming (YARN Mode)

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/file2 on your driver node, only the executor running on the driver node can see those files. Executors on other nodes will throw a FileNotFoundException when 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).

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/

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/file2 on 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 ignoreMissingFiles option 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:52:23