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

Spark Connect Docker容器写入Parquet失败的配置问询

Spark Connect 容器写入Parquet文件失败的解决方法

问题背景

Spark Connect(Apache和Bitnami镜像)与Jupyter容器通信正常,可成功创建并显示DataFrame,但执行df.write.save()写入Parquet文件时始终失败,报错无法创建临时目录。

现有配置

docker-compose.yml

services:
  spark-connect:
    image: apache/spark:latest
    container_name: spark-connect
    environment:
      - SPARK_NO_DAEMONIZE=yes
    ports:
      - '4040:4040'
      - "15002:15002"
    command: /opt/spark/sbin/start-connect-server.sh --packages org.apache.spark:spark-connect_2.12:3.5.1 --conf spark.driver.extraJavaOptions="-Divy.cache.dir=/tmp -Divy.home=/tmp" 
    networks:
      - arr-network

  spark-connect-bitnami:
    image: docker.io/bitnami/spark:latest
    container_name: spark-connect-bitnami
    ports:
      - '4041:4040'
      - "15003:15002"
    command: start-connect-server.sh --packages org.apache.spark:spark-connect_2.12:3.5.1 
    networks:
      - arr-network

  jupyter:
    build:
      context: .
      dockerfile: Dockerfile_jupyter
    container_name: jupyter
    volumes:
      - ./jupyter:/home/jovyan/arr
      - ${HOME}/.config:/home/jovyan/.config
      # - ../arr-hive:/home/spark/spark-warehouse
    ports: 
      - "8888:8888"
    command: start-notebook.py --ip='*' --NotebookApp.token='' --NotebookApp.password='' 
    networks:
      - arr-network

networks:
  arr-network:
    driver: bridge

Jupyter Dockerfile

FROM quay.io/jupyter/base-notebook

RUN mamba install --yes pandas pyspark[connect] grpcio grpcio-status black jupyterlab_code_formatter && \
    mamba clean --all -f -y && \
    fix-permissions "${CONDA_DIR}" && \
    fix-permissions "/home/${NB_USER}"

USER root

ARG spark_uid=185

RUN groupadd --system --gid=${spark_uid} spark && \
    useradd --system --uid=${spark_uid} --gid=spark spark --create-home && \
    usermod -a -G users spark && \
    mkdir -m 775 /home/spark/spark-warehouse && \
    chown spark /home/spark/spark-warehouse && \
    chgrp users /home/spark/spark-warehouse && \
    echo "spark ALL=(ALL) NOPASSWD: ALL" > /etc/sudoers && \
    chmod 0440 /etc/sudoers 

USER spark

测试代码

from pyspark.sql import SparkSession
from pathlib import Path

# directory owned by spark (uid: 185)
save_dir = Path.home().joinpath("../spark/spark-warehouse")
with open(save_dir.joinpath("touch_this"), "w") as f:
    f.write("touched")

spark = (
    SparkSession.builder.appName("arr")
    .remote("sc://spark-connect-bitnami:15002")
    .config("spark.sql.warehouse.dir", save_dir)
    .enableHiveSupport()
    .getOrCreate()
)

columns = ["id","name"]
data = [(1,"Sarah"),(2,"Maria")]
df = spark.createDataFrame(data).toDF(*columns)
df.show()

报错信息

---------------------------------------------------------------------------
SparkConnectGrpcException                 Traceback (most recent call last)
Cell In[2], line 1
----> 1 df.write.save(f'{save_dir.joinpath("test.parquet")}')

File /opt/conda/lib/python3.11/site-packages/pyspark/sql/connect/readwriter.py:601, in DataFrameWriter.save(self, path, format, mode, partitionBy, **options)
    599     self.format(format)
    600 self._write.path = path
--> 601 self._spark.client.execute_command(self._write.command(self._spark.client))

File /opt/conda/lib/python3.11/site-packages/pyspark/sql/connect/client/core.py:982, in SparkConnectClient.execute_command(self, command)
    980     req.user_context.user_id = self._user_id
    981 req.plan.command.CopyFrom(command)
--> 982 data, _, _, _, properties = self._execute_and_fetch(req)
    983 if data is not None:
    984     return (data.to_pandas(), properties)

File /opt/conda/lib/python3.11/site-packages/pyspark/sql/connect/client/core.py:1283, in SparkConnectClient._execute_and_fetch(self, req, self_destruct)
   1280 schema: Optional[StructType] = None
   1281 properties: Dict[str, Any] = {}
-> 1283 for response in self._execute_and_fetch_as_iterator(req):
   1284     if isinstance(response, StructType):
   1285         schema = response

File /opt/conda/lib/python3.11/site-packages/pyspark/sql/connect/client/core.py:1264, in SparkConnectClient._execute_and_fetch_as_iterator(self, req)
   1262                     yield from handle_response(b)
   1263 except Exception as error:
-> 1264     self._handle_error(error)

File /opt/conda/lib/python3.11/site-packages/pyspark/sql/connect/client/core.py:1503, in SparkConnectClient._handle_error(self, error)
   1490 """
   1491 Handle errors that occur during RPC calls.
   1492 
   (...)
   1500 Throws the appropriate internal Python exception.
   1501 """
   1502 if isinstance(error, grpc.RpcError):
-> 1503     self._handle_rpc_error(error)
   1504 elif isinstance(error, ValueError):
   1505     if "Cannot invoke RPC" in str(error) and "closed" in str(error):

File /opt/conda/lib/python3.11/site-packages/pyspark/sql/connect/client/core.py:1539, in SparkConnectClient._handle_rpc_error(self, rpc_error)
   1537             info = error_details_pb2.ErrorInfo()
   1538             d.Unpack(info)
-> 1539             raise convert_exception(info, status.message) from None
   1541     raise SparkConnectGrpcException(status.message) from None
   1542 else:

SparkConnectGrpcException: (org.apache.spark.SparkException) Job aborted due to stage failure: Task 0 in stage 7.0 failed 1 times, most recent failure: Lost task 0.0 in stage 7.0 (TID 16) (6d6e48933499 executor driver): java.io.IOException: Mkdirs failed to create file:/home/spark/spark-warehouse/test.parquet/_temporary/0/_temporary/attempt_20240806134346555678479421024381_0007_m_000000_16 (exists=false, cwd=file:/opt/bitnami/spark)
    at org.apache.hadoop.fs.ChecksumFileSystem.create(ChecksumFileSystem.java:515)
    at org.apache.hadoop.fs.ChecksumFileSystem.create(ChecksumFileSystem.java:500)
    at org.apache.hadoop.fs.FileSystem.create(FileSystem.java:1195)
    at org.apache.hadoop.fs.FileSystem.create(FileSystem.java:1175)
    at org.apache.parquet.hadoop.util.HadoopOutputFile.create(HadoopOutputFile.java:74)
    at org.apache.parquet.hadoop.ParquetFileWriter.<init>(ParquetFileWriter.java:347)
    at org.apache.parquet.hadoop.ParquetFileWriter.<init>(ParquetFileWriter.java:314)
    at org.apache.parquet.hadoop.ParquetOutputFormat.getRecordWriter(ParquetOutputFormat.java:484)
    at org.apache.parquet.hadoop.ParquetOutputFormat.getRecordWriter(ParquetOutputFormat.java:422)
    at org.apache.parquet.hadoop.ParquetOutputFormat.getRecordWriter(ParquetOutputFormat.java:411)
    at org.apache.spark.sql.execution.datasources.parquet.ParquetOutputWriter.<init>(ParquetOutputWriter.scala:36)
    at org.apache.spark.sql.execution.datasources.parquet.ParquetUtils$$anon$1.newInstance(ParquetUtils.scala:490)
    at org.apache.spark.sql.execution.datasources.SingleDirectoryDataWriter.newOutputWriter(FileFormatDataWriter.scala:161)
    at org.apache.spark.sql.execution.datasources.SingleDirectoryDataWriter.<init>(FileFormatDataWriter.scala:146)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$.executeTask(FileFormatWriter.scala:389)
    at org.apache.spark.sql.execution.datasources.WriteFilesExec.$anonfun$doExecuteWrite$1(WriteFiles.scala:100)
    at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
    at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
    at or...

已尝试的操作

  • 统一Spark Connect容器与Jupyter容器的UID
  • 为目标目录设置ACL权限
  • 调整目录工作权限

解决方案

1. 配置共享存储卷

Spark Connect的Driver/Executor运行在自身容器内,写入操作由Spark容器进程执行,因此需要让Spark容器和Jupyter容器共享同一个物理存储目录:

  • 修改docker-compose.yml,给两个Spark服务添加共享卷:
    spark-connect:
      # ...其他配置
      volumes:
        - ./spark-warehouse:/home/spark/spark-warehouse
    spark-connect-bitnami:
      # ...其他配置
      volumes:
        - ./spark-warehouse:/home/spark/spark-warehouse
    
  • 同时启用Jupyter容器中注释的卷配置:
    jupyter:
      # ...其他配置
      volumes:
        # ...其他卷
        - ./spark-warehouse:/home/spark/spark-warehouse
    
  • 宿主机提前创建目录并设置权限:
    mkdir ./spark-warehouse
    sudo chown 185:185 ./spark-warehouse
    sudo chmod 775 ./spark-warehouse
    

2. 修正Spark Session路径配置

避免使用相对路径,直接指定容器内的绝对共享路径,防止不同容器中路径解析不一致:

save_dir = "/home/spark/spark-warehouse"

3. 验证Spark容器内权限

进入Spark容器检查挂载目录的权限:

docker exec -it spark-connect-bitnami bash
ls -ld /home/spark/spark-warehouse

确保所有者为spark用户(UID 185),且具备读写权限。

4. 可选:使用分布式存储(适配Golang ETL场景)

如果后续用Golang做ETL,建议使用HDFS、MinIO(S3兼容)等分布式存储替代本地文件系统,避免跨容器文件访问问题:

  • 在Spark Connect启动命令中添加Hadoop配置:
    --conf spark.hadoop.fs.defaultFS=hdfs://hadoop-namenode:9000
    
  • 部署对应的分布式存储集群,确保Spark容器和Golang应用都能访问该存储服务。

关键原因

Spark Connect采用客户端-服务端架构:Jupyter/Golang是客户端仅提交任务,实际的数据写入操作由Spark容器内的Driver/Executor执行。之前的错误是因为客户端指定的路径仅在Jupyter容器内存在,Spark容器内无对应目录或权限不足。

内容的提问来源于stack exchange,提问作者Rob Raymond

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 20:38:10