Spark集群读取CSV正常但JSON报文件不存在的原因排查
集群配置(Docker Compose)
version: "3.3" services: spark-master: image: ${IMAGE_VERSION} ports: - "9090:8080" - "7077:7077" volumes: - ./apps:/home/developer/pyspark-apps - ./data:/opt/spark-data - ./jars:/home/developer/jars environment: - SPARK_LOCAL_IP=spark-master - SPARK_WORKLOAD=master - SPARK_WORKER_CORES=2 - SPARK_WORKER_MEMORY=1G - SPARK_DRIVER_MEMORY=1G - SPARK_EXECUTOR_MEMORY=2G spark-worker-a: image: ${IMAGE_VERSION} ports: - "9091:8080" - "7000:7000" depends_on: - spark-master environment: - SPARK_MASTER=spark://spark-master:7077 - SPARK_WORKER_CORES=2 - SPARK_WORKER_MEMORY=2G - SPARK_DRIVER_MEMORY=2G - SPARK_EXECUTOR_MEMORY=2G - SPARK_WORKLOAD=worker - SPARK_LOCAL_IP=spark-worker-a volumes: - ./apps:/home/developer/pyspark-apps - ./data:/opt/spark-data - ./jars:/home/developer/jars spark-worker-b: image: ${IMAGE_VERSION} ports: - "9092:8080" - "7001:7000" depends_on: - spark-master environment: - SPARK_MASTER=spark://spark-master:7077 - SPARK_WORKER_CORES=2 - SPARK_WORKER_MEMORY=2G - SPARK_DRIVER_MEMORY=2G - SPARK_EXECUTOR_MEMORY=2G - SPARK_WORKLOAD=worker - SPARK_LOCAL_IP=spark-worker-b volumes: - ./apps:/home/developer/pyspark-apps - ./data:/opt/spark-data - ./jars:/home/developer/jars
CSV文件读取正常场景
本地Jupyter Lab运行以下代码,成功读取CSV文件:
import pyspark import pyspark.sql.functions as F import os import pandas as pd from pyspark import SparkConf, SparkContext from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName('big_file_processor') \ .config("spark.sql.execution.arrow.pyspark.enabled", "true") \ .getOrCreate() df = spark.read.csv('data/big-data.csv').withColumnRenamed("_c0","fake_name") df.show()
返回结果:
+---------+ |fake_name| +---------+ |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| |fake_data| +---------+
only showing top 20 rows
JSON文件读取报错场景
执行以下读取JSON文件的代码:
schema = StructType([ StructField("name", StringType(), True), StructField("doc_list", ArrayType(MapType(StringType(),StringType(),True),True), True), ]) df = spark.read.json('data/test.json', schema=schema) df.show()
出现错误:
Py4JJavaError: An error occurred while calling o123.showString.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 5.0 > failed 4 times, most recent failure: Lost task 0.3 in stage 5.0 (TID 19) (172.19.0.4 executor 0): java.io.FileNotFoundException:
File file:/home/andre/Projects/Spark/spark-standalone/apps/dataframes/data/test.json does not exist
报错原因分析
路径解析不匹配:
你在Jupyter中使用相对路径data/test.json,本地Driver进程(Jupyter)会将其解析为本地文件系统的绝对路径/home/andre/Projects/Spark/spark-standalone/apps/dataframes/data/test.json。但Spark Standalone的Executor运行在Docker容器内,容器中data卷的映射路径是/opt/spark-data,而非本地路径。当Spark将JSON读取任务分发到Executor时,Executor会在自身容器的文件系统中查找上述本地绝对路径,自然找不到对应文件,抛出异常。CSV读取成功的特殊逻辑:
CSV文件读取成功是因为文件体积较小,Spark直接在Driver端完成了读取和处理,没有触发分布式任务(无需Executor参与),因此没有暴露路径不匹配的问题。而JSON文件读取需要进行Schema校验或分片处理,必须触发Executor执行任务,从而暴露了路径问题。
解决方法
- 直接使用容器内的绝对路径读取文件:修改代码中的路径为
/opt/spark-data/test.json,示例:df = spark.read.json('/opt/spark-data/test.json', schema=schema) - 或者采用分布式文件系统(如HDFS)存储文件,确保Driver和Executor能通过统一路径访问文件。
内容的提问来源于stack exchange,提问作者Andre Carneiro

