并发运行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配置:
注:使用该日志存储需提前创建DynamoDB表(默认表名为.config("spark.delta.logStore.class", "org.apache.spark.sql.delta.storage.S3DynamoDBLogStore")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
相关产品推荐
相关产品推荐

