多实例批量作业更新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表的更新过程中:
- 每个作业实例的UPDATE语句会匹配
allocations_od中的多条记录,PostgreSQL会按照数据在磁盘上的物理存储顺序逐个获取行级排他锁(Row Exclusive Lock)。 - 当两个作业实例需要更新的行集合存在重叠,但锁的获取顺序相反时,就会形成循环等待:
- 进程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

