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

多实例批量作业更新PostgreSQL表触发死锁问题排查

批量作业更新时的PostgreSQL死锁问题分析与解决建议

问题描述

当用户在UI中更新数据行时,会触发批量作业。用户可同时更新多行,从而触发多个带有唯一run_id的批量作业实例。

该作业会生成CSV文件并将数据插入allocations_update表,随后使用该表的数据更新第二个表allocations_od。

更新allocations_od的SQL语句如下:

UPDATE db.allocations_od target
SET rec_alloc = src.rec_alloc 
FROM db.allocations_update src
WHERE  src.run_id = '{run_id}' 
  AND  src.col1 = target.col1
  AND  src.col2 = target.col2

但当用户同时触发多个作业实例(即同时更新多行)时,执行更新allocations_od的语句会出现死锁错误。完整错误信息如下:

psycopg2.errors.DeadlockDetected: deadlock detected
DETAIL:  Process 15455 waits for ShareLock on transaction 62597603; blocked by process 15538.
Process 15538 waits for ShareLock on transaction 62597592; blocked by process 15455.
HINT:  See server log for query details.
CONTEXT:  while updating tuple (479821,43) in relation ""allocations_od_20230514""

初步猜测是其他作业实例仍在执行插入语句,占用了allocations_update表的锁,导致两个进程互相阻塞,想了解引发死锁的具体原因及解决办法。

相关核心代码

批量更新主方法

def update_alloc_query(self, final_data, stage_location):
    """ Method to bulk update allocations od table"""

    # stage_location is the s3 path of csv file.

    last_created_date = self.get_last_created_date()
    last_created_date = last_created_date.strftime('%Y-%m-%d')
    final_data['created_date'] = last_created_date
    run_id = final_data['run_id'].unique()[0]
    s3.s3_upload_df(stage_location, final_data)
    UITableLoader.bulk_upload_from_csv(db_model=AllocationsUpdate,
                                       file_location=stage_location,
                                       data_types={"rsid": "str", "passenger_class": "str",
                                                   "journey_origin": "str",
                                                   "journey_destination": "str",
                                                   "bucket_code": "str",
                                                   "eff_departure_date": "str",
                                                   "recommended_allocation": "float",
                                                   "run_id": "str"},
                                       sep="|",
                                       created_date=last_created_date)
    self.logger.info("Added table into new data")
    allo_sql = f"UPDATE db.allocations_od target\
                set rec_alloc = src.rec_alloc FROM\
                db.allocations_update src\
                WHERE src.run_id = '{run_id}' AND  \
                src.col1 = target.col1 AND\
                src.col2 = target.col2'"
    execute_sql_statement(allo_sql)
    self.logger.info("executed update query")

CSV批量插入方法

# UITableLoader.bulk_upload_from_csv
@staticmethod
def bulk_upload_from_csv(db_model, file_location, data_types=None, sep=',',
               created_date=None, chunk_size=1000):
    """Function uploads data from local csv file to sql alchemy db."""
    LOGGER.info("Bulk loading data.",
                file_location=file_location, table=db_model.__table__)
    record_count = 0
    chunks = pd.read_csv(
        file_location,
        dtype=data_types,
        chunksize=chunk_size,
        sep=sep,
        on_bad_lines='skip'
    )

    for chunk in chunks:
        chunk = chunk.where((pd.notnull(chunk)), None)
        chunk = chunk.replace({np.nan: None})
        record_count += chunk.shape[0]
        if created_date is not None:
            chunk['created_date'] = created_date
        rows = chunk.to_dict(orient='records')
        sqa_save(db_model, rows, save_many=True)

    return record_count

SQL执行工具方法

def execute_sql_statement(sql_statement, conn_string=None):  # pragma: no cover
    """Executes the given sql_statement"""
    if not sql_statement:
        return
    if not conn_string:
        conn_string = get_db_connection_string()
    dbsession = get_db_session(conn_string)
    try:
        dbsession.execute(sql_statement)
        dbsession.commit()
    except SQLAlchemyError as ex:
        LOGGER.exception(f"Error executing sql statement '{sql_statement}'")
        dbsession.rollback()
        raise ex
    finally:
        dbsession.close()

死锁原因分析

你的初步猜测方向有误,死锁并非来自allocations_update表的插入锁,而是发生在allocations_od表的更新过程中:

  1. 每个作业实例的UPDATE语句会匹配allocations_od中的多条记录,PostgreSQL会按照数据在磁盘上的物理存储顺序逐个获取行级排他锁(Row Exclusive Lock)。
  2. 当两个作业实例需要更新的行集合存在重叠,但锁的获取顺序相反时,就会形成循环等待:
    • 进程A先锁定行X,然后尝试锁定行Y
    • 进程B先锁定行Y,然后尝试锁定行X
    • 双方互相等待对方释放锁,最终触发死锁检测

从错误信息的CONTEXT字段可以看到,死锁发生在allocations_od_20230514(看起来是分区表)的元组更新操作上,进一步验证了这一点。

另外,当前的批量插入逻辑是分chunk提交事务,插入操作的锁在提交后就会释放,不会长期持有,所以不是插入锁导致的死锁。

解决方案

1. 强制更新时的行锁获取顺序

修改UPDATE语句,通过ORDER BY强制PostgreSQL按照固定顺序获取行锁,避免循环等待:

UPDATE db.allocations_od target
SET rec_alloc = src.rec_alloc 
FROM db.allocations_update src
WHERE  src.run_id = '{run_id}' 
  AND  src.col1 = target.col1
  AND  src.col2 = target.col2
ORDER BY target.col1, target.col2;

PostgreSQL 9.5+支持在UPDATE语句中使用ORDER BY,确保所有进程都按相同顺序锁定行,从根源上避免死锁。

2. 增加死锁重试机制

在执行UPDATE语句时添加重试逻辑,因为死锁是偶发的,重试通常能解决问题:

import time
from sqlalchemy.exc import SQLAlchemyError

def execute_sql_statement(sql_statement, conn_string=None, retry_times=3):
    """Executes the given sql_statement with deadlock retry"""
    if not sql_statement:
        return
    if not conn_string:
        conn_string = get_db_connection_string()
    for attempt in range(retry_times):
        dbsession = get_db_session(conn_string)
        try:
            dbsession.execute(sql_statement)
            dbsession.commit()
            return
        except SQLAlchemyError as ex:
            if "deadlock detected" in str(ex) and attempt < retry_times -1:
                LOGGER.warning(f"Deadlock detected, retrying attempt {attempt+1}/{retry_times}")
                dbsession.rollback()
                time.sleep(0.5)
                continue
            LOGGER.exception(f"Error executing sql statement '{sql_statement}'")
            dbsession.rollback()
            raise ex
        finally:
            dbsession.close()

3. 优化事务范围

当前插入操作是分chunk提交,而更新是单独事务。可以考虑将单个作业实例的插入+更新操作放在同一个事务中,减少并发事务数量,但要注意如果数据量过大,事务可能过长导致其他问题。

4. 检查分区表的索引

如果allocations_od是分区表,确保col1和col2上有合适的索引,让UPDATE语句能快速定位目标行,减少锁持有时间,降低死锁概率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 12:20:03