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

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]
15k395
45k450, 420, 423, 452
85k630, 628, 672, 677, 695, 613, 631, 641
110k432
410k628, 508, 597, 619
810k840, 861, 809, 846, 882, 876, 887, 861

我也曾尝试Delta table格式,遇到相同问题:针对同一张表不同product_id并行执行时,合并耗时大幅增加,疑似AWS Glue任务在合并时存在锁等待情况。

疑问

  1. 是否可以优化单任务合并耗时(5k行耗时395s过长)?
  2. 如何优化不同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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 07:54:52