Apache Iceberg合并耗时过长及多product_id并发执行优化咨询
问题描述
我有一个AWS Glue任务,尝试将数据合并到按product_id分区的Apache Iceberg表中,期望针对不同product_id通过AWS Glue任务执行并发合并操作。
表结构示例
product_id, name, ... , user_id
合并查询逻辑
合并条件:
existing_data.product_id = '{here_product_id}' AND new_data.product_id = existing_data.product_id AND existing_data.user_id = new_data.user_id
执行合并的代码:
merge_sql = f""" MERGE INTO glue_catalog.default.{APACHE_ICEBERG_PREFIX}{target_path} existing_data USING td new_data ON existing_data.product_id = '{here_product_id}' AND new_data.product_id = existing_data.product_id AND existing_data.user_id = new_data.user_id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * """ spark.sql(merge_sql)
测试结果
实际测试发现,即便使用分区,合并耗时依然很长,且针对不同product_id并发执行时,耗时会进一步增加,具体数据如下:
| 运行任务数 | 导出数据行数 | 执行耗时[s] |
|---|---|---|
| 1 | 5k | 395 |
| 4 | 5k | 450, 420, 423, 452 |
| 8 | 5k | 630, 628, 672, 677, 695, 613, 631, 641 |
| 1 | 10k | 432 |
| 4 | 10k | 628, 508, 597, 619 |
| 8 | 10k | 840, 861, 809, 846, 882, 876, 887, 861 |
我也曾尝试Delta table格式,遇到相同问题:针对同一张表不同product_id并行执行时,合并耗时大幅增加,疑似AWS Glue任务在合并时存在锁等待情况。
疑问
- 是否可以优化单任务合并耗时(5k行耗时395s过长)?
- 如何优化不同
product_id并发执行的合并耗时?
解决方案
一、单任务合并耗时优化
1. 优化合并查询的分区过滤逻辑
当前ON条件的分区过滤可能未被正确下推,建议提前过滤源表数据,减少参与合并的数据量:
# 先过滤出目标product_id的源数据 filtered_td = td.filter(f"product_id = '{here_product_id}'") filtered_td.createOrReplaceTempView("filtered_td") merge_sql = f""" MERGE INTO glue_catalog.default.{APACHE_ICEBERG_PREFIX}{target_path} existing_data USING filtered_td new_data ON new_data.product_id = existing_data.product_id AND existing_data.user_id = new_data.user_id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * """
同时通过DESCRIBE EXTENDED命令验证Iceberg表的分区列是否被Glue Catalog正确识别。
2. 调整Glue任务资源与Spark配置
- 升级Worker类型(如从Standard切换到G.1X/G.2X)并增加Worker数量,解决CPU/内存瓶颈;
- 调整Spark核心参数:
# 减少不必要的Shuffle分区数 spark.conf.set("spark.sql.shuffle.partitions", "100") # 开启小表自动广播,避免Shuffle spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "104857600") # 100MB
3. Iceberg专属优化
- 启用Merge-On-Read模式,避免立即重写整个分区,降低IO开销:
ALTER TABLE glue_catalog.default.your_table SET TBLPROPERTIES ('write.merge.mode' = 'merge-on-read') - 为
user_id建立二级索引,加速匹配查询:CREATE INDEX user_id_idx ON glue_catalog.default.your_table (user_id)
二、并发合并耗时优化(解决锁等待问题)
1. 启用Iceberg分区级锁
Iceberg支持分区级乐观锁,而非默认表级锁,需手动开启:
ALTER TABLE glue_catalog.default.your_table SET TBLPROPERTIES ( 'lock.enabled' = 'true', 'lock.partition-level' = 'true' )
确保每个并发任务仅处理单个product_id分区,避免跨分区操作触发表级锁。
2. 规避全局资源竞争
- 为并发任务分配独立的Glue Worker队列,避免资源争抢;
- 控制并发任务数量(比如从8个降至4个),避免超出Glue或S3的并发限制。
3. 批量合并替代单分区并发
若锁问题无法彻底解决,可将多个product_id的合并逻辑整合到单个任务中批量处理,减少任务间的锁竞争。
4. Delta表补充优化(若继续使用)
启用Delta的分区级写入控制,限定锁的范围:
ALTER TABLE your_delta_table SET TBLPROPERTIES ('delta.enablePartitionWrites' = 'true')
内容的提问来源于stack exchange,提问作者P.Zaw
相关产品推荐
相关产品推荐

