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

并发运行Delta Lake作业写入不相交分区触发ProtocolChangedException

问题:并发写入Delta Lake不相交分区时遭遇ProtocolChangedException

我在AWS S3上并发运行两个独立的Apache Spark + Delta Lake作业实例,每个实例处理独立的Parquet文件,并写入Delta表的不相交分区。尽管文件和分区完全独立,仍遇到以下异常:

delta.exceptions.ProtocolChangedException: The protocol version of the Delta table has been changed by a concurrent update. This happens when multiple writers are writing to an empty directory. Creating the table ahead of time will avoid this conflict. Please try the operation again.
Conflicting commit: {"timestamp":1670947072124,"operation":"WRITE","operationParameters":{"mode":Overwrite,"partitionBy":["PartitionColumn"]},"isolationLevel":"Serializable","isBlindAppend":false,"operationMetrics":{"numFiles":"1","numOutputRows":"30915","numOutputBytes":"7007784"},"engineInfo":"Apache-Spark/3.3.1 Delta-Lake/1.2.1","txnId":"58a1c854-13a6-475e-bb0b-cfe2146441cc"}

我尝试过overwrite和append模式,均未能解决问题。理论上,并发写入不相交分区应该不会触发冲突。

作业代码片段如下:

from pyspark.sql import SparkSession
from pyspark.sql.functions import lit, current_date, col, when
from delta import *
from delta.tables import *
import time

spark = SparkSession.builder \
            .appName("PySparkLocal")\
            .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")\
            .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")\
            .config("spark.hadoop.fs.AbstractFileSystem.s3.impl","org.apache.hadoop.fs.s3a.S3AFileSystem")\
            .config("spark.delta.logStore.class", "org.apache.spark.sql.delta.storage.S3SingleDriverLogStore")\
            .config("spark.jars.packages", "io.delta:delta-core_2.12:1.2.1")\
            .config("spark.sql.sources.partitionOverwriteMode", "dynamic")\
            .config("spark.databricks.delta.schema.autoMerge.enabled", "true")\
            .config("spark.sql.parquet.compression.codec", "gzip")\
            .config("spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled", "true")\
            .getOrCreate()


INPUT_DATA_PATH = "s3://inputFile.parquet" # 并行作业使用同一文件夹下的不同输入文件
DELTATABLE_PATH = "s3://DeltaTable"

deltaData = spark.read.parquet(INPUT_DATA_PATH )
deltaData.write.format("delta").option("overwriteSchema","true").partitionBy("PartitionColumn").mode("overwrite").save(DELTATABLE_PATH)

解决方案

1. 预创建Delta表

异常提示明确指出:空目录下多写者会触发协议版本冲突。因此需要提前初始化Delta表,避免并发作业同时创建表结构。

方式一:通过DataFrame创建空表

# 读取样本文件获取数据schema
sample_df = spark.read.parquet("s3://sample-input.parquet")
# 创建空Delta表并指定分区列
sample_df.limit(0).write.format("delta")\
  .partitionBy("PartitionColumn")\
  .mode("ignore")\
  .save("s3://DeltaTable")

方式二:通过Spark SQL创建表

CREATE TABLE IF NOT EXISTS delta.`s3://DeltaTable` (
  -- 按实际数据schema定义列
  col1 STRING,
  col2 INT,
  PartitionColumn STRING
)
PARTITIONED BY (PartitionColumn)
USING delta

2. 优化写入配置与模式

  • 移除不必要的overwriteSchema=true:仅当需要修改表schema时添加该参数,每次写入都携带会增加冲突概率。
  • 使用S3DynamoDBLogStore替代S3SingleDriverLogStore:后者为单节点写日志,在S3并发场景下易引发冲突。修改Spark配置:
    .config("spark.delta.logStore.class", "org.apache.spark.sql.delta.storage.S3DynamoDBLogStore")
    
    注:使用该日志存储需提前创建DynamoDB表(默认表名为delta_logs,可通过spark.delta.logStore.s3dynamodb.tableName自定义)。
  • 配合partitionOverwriteMode=dynamic使用overwrite模式:确保仅覆盖目标分区而非整个表,减少锁竞争。

3. 使用DeltaTable API提升并发兼容性

相比DataFrame的write API,DeltaTable的merge/upsert API对并发场景支持更友好:

from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, DELTATABLE_PATH)
# 合并数据(若需覆盖分区可调整逻辑)
delta_table.alias("target")\
  .merge(
    deltaData.alias("source"),
    "target.PartitionColumn = source.PartitionColumn"
  )\
  .whenMatchedUpdateAll()\
  .whenNotMatchedInsertAll()\
  .execute()

若仅需追加数据到分区,可简化为:

deltaData.write.format("delta")\
  .mode("append")\
  .partitionBy("PartitionColumn")\
  .save(DELTATABLE_PATH)

内容的提问来源于stack exchange,提问作者KIRAN REDDY

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 22:40:44