PySpark读写同表遇UNSUPPORTED_OVERWRITE.TABLE错误求助
Docker化Spark集群PySpark覆盖写入Parquet表失败问题
核心问题
执行SCD1更新后,将DataFrame写回原Parquet表时触发报错:
AnalysisException: [UNSUPPORTED_OVERWRITE.TABLE] Can't overwrite the target that is also being read from
场景流程:
- 将初始数据写入指定路径的Parquet表
- 读取该表进行数据更新
- 尝试覆盖写入原表时失败
已尝试无效方案:
- 对读取后的DataFrame执行
cache()+count()后再写入 - 手动删除底层文件后写入
复现代码
from pyspark.sql import SparkSession import pyspark.sql.functions as F # 连接Docker化Spark集群 spark = SparkSession.builder.remote("sc://localhost:15002").getOrCreate() # 配置参数 path = "/data/user_table" database_name = "DB_test" table_name = "test_users" # 初始数据写入 df_ = spark.createDataFrame( [(0, "Tom"), (0, "Alice")], ["age", "name"]) df_.write.mode('overwrite') \ .option("path", path) \ .saveAsTable(f"{database_name}.{table_name}") # 读取原表并更新 df = spark.table(f"{database_name}.{table_name}") df1 = df.withColumn('test', F.lit(1)) # 覆盖写入原表(触发报错) df1.write.mode('overwrite') \ .option("path", path) \ .saveAsTable(f"{database_name}.{table_name}")
集群配置信息
目录结构
./Spark-Docker docker-compose.yml Dockerfile entrypoint.sh
docker-compose.yml
name: spark-cluster services: spark-master: image: my_spark_p3.13 build: context: . dockerfile: ./Dockerfile container_name: spark-master hostname: spark-master environment: - PYTHON_VERSION=python3.13 - SPARK_MODE=master - SPARK_MASTER_WEBUI_PORT=8080 - SPARK_MASTER_PORT=7077 - SPARK_SUBMIT_OPTIONS=--packages io.delta:delta-spark_2.12:3.2.1,io.delta:delta-storage_2.12:3.2.0 --conf "spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension" --conf "spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog" - SPARK_MASTER_HOST=spark-master - PYSPARK_PYTHON=/usr/bin/python3 - PYSPARK_DRIVER_PYTHON=/usr/bin/python3 ports: - 8080:8080 - 7077:7077 networks: - spark-network volumes: - ./data:/data:rw entrypoint: - "bash" - "-c" - "/opt/spark/sbin/start-master.sh && tail -f /dev/null" spark-connect: image: my_spark_p3.13 container_name: spark-connect hostname: spark-connect ports: - "4040:4040" - "15002:15002" networks: - spark-network depends_on: - spark-master volumes: - ./jars/spark-connect_2.12-3.5.1.jar:/opt/spark/jars/spark-connect_2.12-3.5.1.jar - ./jars/delta-spark_2.12-3.2.0.jar:/opt/spark/jars/delta-spark_2.12-3.2.0.jar - ./jars/delta-storage-3.2.0.jar:/opt/spark/jars/delta-storage-3.2.0.jar - ./data:/data:rw command: - "bash" - "-c" - '/opt/spark/sbin/start-connect-server.sh --jars /opt/spark/jars/spark-connect_2.12-3.5.1.jar,/opt/spark/jars/delta-spark_2.12-3.2.0.jar,delta-storage-3.2.0.jar --conf "spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension" --conf "spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog" && tail -f /dev/null' spark-worker: image: my_spark_p3.13 container_name: spark-worker hostname: spark-worker environment: - SPARK_MODE=worker - SPARK_MASTER=spark://spark-master:7077 - SPARK_WORKER_CORES=8 - SPARK_WORKER_MEMORY=16G - SPARK_WORKER_WEBUI_PORT=8081 - PYSPARK_PYTHON=/usr/bin/python3 - PYSPARK_DRIVER_PYTHON=/usr/bin/python3 ports: - 8081:8081 networks: - spark-network volumes: - ./data:/data:rw depends_on: - spark-master entrypoint: - "bash" - "-c" - "/opt/spark/sbin/start-worker.sh spark://spark-master:7077 && tail -f /dev/null" networks: spark-network:
使用的Jar包
- delta-storage-3.2.0.jar
- delta-spark_2.12-3.2.0.jar
- spark-connect_2.12-3.5.1.jar
Dockerfile
FROM eclipse-temurin:11-jre-focal ARG spark_uid=185 RUN groupadd --system --gid=${spark_uid} spark && \ useradd --system --uid=${spark_uid} --gid=spark spark RUN set -ex; \ apt-get update; RUN apt-get install -y software-properties-common; RUN add-apt-repository ppa:deadsnakes/ppa; RUN apt-get install -y python3.13 gnupg2 wget bash tini libc6 libpam-modules krb5-user libnss3 procps net-tools gosu libnss-wrapper; \ mkdir -p /opt/spark; \ mkdir /opt/spark/python; \ mkdir -p /opt/spark/examples; \ mkdir -p /opt/spark/work-dir; \ chmod g+w /opt/spark/work-dir; \ touch /opt/spark/RELEASE; \ chown -R spark:spark /opt/spark; \ echo "auth required pam_wheel.so use_uid" >> /etc/pam.d/su; \ rm -rf /var/lib/apt/lists/* # Install Apache Spark ENV SPARK_TGZ_URL=https://archive.apache.org/dist/spark/spark-3.5.3/spark-3.5.3-bin-hadoop3.tgz \ SPARK_TGZ_ASC_URL=https://archive.apache.org/dist/spark/spark-3.5.3/spark-3.5.3-bin-hadoop3.tgz.asc \ GPG_KEY=0A2D660358B6F6F8071FD16F6606986CF5A8447C RUN set -ex; \ export SPARK_TMP="$(mktemp -d)"; \ cd $SPARK_TMP; \ wget -nv -O spark.tgz "$SPARK_TGZ_URL"; \ wget -nv -O spark.tgz.asc "$SPARK_TGZ_ASC_URL"; \ export GNUPGHOME="$(mktemp -d)"; \ gpg --batch --keyserver hkps://keys.openpgp.org --recv-key "$GPG_KEY" || \ gpg --batch --keyserver hkps://keyserver.ubuntu.com --recv-keys "$GPG_KEY"; \ gpg --batch --verify spark.tgz.asc spark.tgz; \ gpgconf --kill all; \ rm -rf "$GNUPGHOME" spark.tgz.asc; \ \ tar -xf spark.tgz --strip-components=1; \ chown -R spark:spark .; \ mv jars /opt/spark/; \ mv RELEASE /opt/spark/; \ mv bin /opt/spark/; \ mv sbin /opt/spark/; \ mv kubernetes/dockerfiles/spark/decom.sh /opt/; \ mv examples /opt/spark/; \ ln -s "$(basename /opt/spark/examples/jars/spark-examples_*.jar)" /opt/spark/examples/jars/spark-examples.jar; \ mv kubernetes/tests /opt/spark/; \ mv data /opt/spark/; \ mv python/pyspark /opt/spark/python/pyspark/; \ mv python/lib /opt/spark/python/lib/; \ mv R /opt/spark/; \ chmod a+x /opt/decom.sh; \ cd ..; \ rm -rf "$SPARK_TMP"; COPY entrypoint.sh /opt/entrypoint.sh RUN chmod 777 /opt/entrypoint.sh ENV SPARK_HOME /opt/spark WORKDIR /opt/spark/work-dir RUN mv /usr/bin/python3.13 /usr/bin/python3 USER spark ENTRYPOINT [ "/opt/entrypoint.sh" ]
entrypoint
使用apache/spark-docker 3.5.3版本的entrypoint脚本
解决方案
1. 临时路径中转(推荐生产使用)
通过临时路径将读取和写入操作解耦,避免同表读写冲突:
# 读取原表后先写入临时路径 df.write.mode('overwrite').parquet("/data/temp_user_table") # 从临时路径读取数据 temp_df = spark.read.parquet("/data/temp_user_table") # 执行数据更新 df1 = temp_df.withColumn('test', F.lit(1)) # 覆盖写入原表 df1.write.mode('overwrite') \ .option("path", path) \ .saveAsTable(f"{database_name}.{table_name}") # 清理临时路径(可选,根据需求保留) import os import shutil if os.path.exists("/data/temp_user_table"): shutil.rmtree("/data/temp_user_table")
2. 使用CTAS语句覆盖
通过SQL的CREATE OR REPLACE TABLE语句绕开DataFrame API的限制:
# 将更新后的DataFrame注册为临时视图 df1.createOrReplaceTempView("updated_users_view") # 执行CTAS覆盖原表 spark.sql(f""" CREATE OR REPLACE TABLE {database_name}.{table_name} LOCATION '{path}' AS SELECT * FROM updated_users_view """)
3. 统一Delta版本并调整配置
当前集群中spark-master使用delta-spark_2.12:3.2.1,spark-connect使用3.2.0,版本不一致可能引发问题,统一为3.2.1:
- 更新spark-connect的command参数中的delta-spark版本为3.2.1
- 若需放宽限制(仅测试环境使用),添加以下配置:
--conf "spark.sql.delta.allowOverwriteSchema=true" --conf "spark.sql.delta.allowConcurrentWrites=true"
4. 强制触发数据加载(小数据量场景适用)
通过collect()强制完成读取操作后再执行写入:
df = spark.table(f"{database_name}.{table_name}") # 强制触发数据加载到内存 df.collect() # 执行更新并写入 df1 = df.withColumn('test', F.lit(1)) df1.write.mode('overwrite') \ .option("path", path) \ .saveAsTable(f"{database_name}.{table_name}")
内容的提问来源于stack exchange,提问作者Ndndhjx Bznznz
相关产品推荐
相关产品推荐

