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
相关产品推荐
相关产品推荐

