在业务系统中,当需要对多张关联表执行大规模数据更新时,同步执行的方式很容易引发数据库锁竞争、事务超时等问题,尤其是更新数据量达到万级甚至十万级时,主业务的响应速度会受到严重影响。通过触发器捕获更新事件,再将更新任务推入异步队列处理,是兼顾性能与数据一致性的有效方案。

核心实现思路
整个方案分为三个核心部分:触发器捕获变更事件、中间表存储待处理任务、异步消费者执行关联更新。触发器在源表发生更新时自动触发,将需要更新的关联信息写入任务中间表,而不是直接执行关联更新操作,主事务可以快速提交。之后由独立的异步进程从中间表读取任务,批量执行关联更新,避免阻塞主业务。
触发器设计实现
首先需要创建任务中间表,用于存储待处理的关联更新任务,表结构可以参考以下设计:
-- 创建异步更新任务表
CREATE TABLE async_update_task (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
source_table VARCHAR(64) NOT NULL COMMENT '源表名称',
source_id BIGINT NOT NULL COMMENT '源表记录ID',
update_field VARCHAR(64) NOT NULL COMMENT '更新字段名',
new_value VARCHAR(256) NOT NULL COMMENT '新值',
status TINYINT DEFAULT 0 COMMENT '任务状态 0待处理 1处理中 2已完成 3失败',
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_status_create_time (status, create_time)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='关联更新异步任务表';
接下来在源表上创建更新触发器,当源表的指定字段发生更新时,自动向任务表插入任务记录:
-- 假设源表为 user_info,当 user_name 字段更新时触发任务写入
DELIMITER //
CREATE TRIGGER trigger_user_info_update
AFTER UPDATE ON user_info
FOR EACH ROW
BEGIN
-- 仅当 user_name 字段发生变更时插入任务
IF OLD.user_name != NEW.user_name THEN
INSERT INTO async_update_task (
source_table,
source_id,
update_field,
new_value
) VALUES (
'user_info',
NEW.id,
'user_name',
NEW.user_name
);
END IF;
END //
DELIMITER ;
异步队列消费者实现
异步消费者需要定时从任务表中读取待处理任务,批量执行关联更新操作。以下是基于Python的消费者示例,使用pymysql连接数据库:
import pymysql
import time
def get_db_conn():
# 连接数据库,实际使用中替换为自己的数据库地址
return pymysql.connect(
host='127.0.0.1',
port=3306,
user='root',
password='123456',
database='test_db',
charset='utf8mb4'
)
def process_async_tasks():
conn = get_db_conn()
cursor = conn.cursor()
try:
# 每次批量读取100条待处理任务,加行锁避免重复处理
cursor.execute(
"""
SELECT id, source_id, new_value
FROM async_update_task
WHERE status = 0
ORDER BY create_time ASC
LIMIT 100 FOR UPDATE
"""
)
tasks = cursor.fetchall()
if not tasks:
return
task_ids = [task[0] for task in tasks]
# 先标记任务为处理中
cursor.execute(
"""
UPDATE async_update_task
SET status = 1
WHERE id IN %s
""", (tuple(task_ids),)
)
conn.commit()
# 执行关联更新,这里示例是更新 user_log 表中对应用户的 user_name
for task in tasks:
task_id, source_id, new_value = task
cursor.execute(
"""
UPDATE user_log
SET user_name = %s
WHERE user_id = %s
""", (new_value, source_id)
)
# 标记任务为已完成
cursor.execute(
"""
UPDATE async_update_task
SET status = 2
WHERE id IN %s
""", (tuple(task_ids),)
)
conn.commit()
except Exception as e:
conn.rollback()
# 标记失败任务
if task_ids:
cursor.execute(
"""
UPDATE async_update_task
SET status = 3
WHERE id IN %s
""", (tuple(task_ids),)
)
conn.commit()
print(f"处理任务失败: {e}")
finally:
cursor.close()
conn.close()
if __name__ == '__main__':
# 每5秒执行一次任务处理
while True:
process_async_tasks()
time.sleep(5)
注意事项
- 触发器逻辑尽量简单,仅做任务写入操作,避免触发器执行耗时过长影响主事务性能。
- 任务中间表需要定期清理已完成的历史任务,避免表数据量过大影响查询性能。
- 异步消费者的批量处理数量需要根据实际数据库性能调整,避免单次更新数据量过大再次引发性能问题。
- 如果出现消费者处理失败的情况,需要增加重试机制,保证最终数据一致性。
适用场景说明
这种方案适合对数据一致性要求为最终一致、更新时效性要求不高的场景,比如用户资料变更后同步更新关联的历史记录表、订单状态变更后同步更新统计表等。如果业务要求更新必须实时同步,则不适合使用该方案。