Golang+Gorm并发Goroutine中BatchUpsert死锁问题排查
问题描述
定期从other_table获取数据,解析后写入example_table。串行处理正常但速度慢,改为并行处理后出现MySQL死锁。
相关信息
1. 数据库表结构
CREATE TABLE `other_table` ( `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT COMMENT 'id', `column1` VARCHAR(100) NOT NULL DEFAULT '', `column7` VARCHAR(10) NULL DEFAULT '' , `column8` TEXT NULL DEFAULT NULL COMMENT '' , PRIMARY KEY (`id`), UNIQUE KEY `udx_column1` (`column1`), ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='other_table'; CREATE TABLE `example_table` ( `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT COMMENT 'id', `column1` VARCHAR(100) NOT NULL DEFAULT '', `column2` VARCHAR(50) NULL DEFAULT '' , `column3` VARCHAR(150) NULL DEFAULT '' , `column4` VARCHAR(10) NULL DEFAULT '' , `column5` TEXT NULL DEFAULT NULL COMMENT '' , `deleted_at` bigint(20) UNSIGNED NOT NULL DEFAULT '0' COMMENT '删除时间', PRIMARY KEY (`id`), UNIQUE KEY `udx_column1_column2_column3_deleted_at` (`column1`,`column2`,`column3`,`deleted_at`), KEY `idx_deleted_at` (`deleted_at`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='example_table';
2. Gorm批量Upsert实现
// Update if exist, create if not . func (e *example) BatchUpsert(ctx context.Context, es []*model.Example) error { return e.db.Clauses(clause.OnConflict{ Columns: []clause.Column{ {Name: "column1"}, {Name: "column2"}, {Name: "column3"}, }, UpdateAll: true, }).Create(&es).Error }
3. 并发处理代码
通过拆分数据到协程池并行处理,解析other_table.column8为example_table的多个字段,定时任务执行:
const ( DataSize = 100 DataPerTask = 10 ) type Task struct { index int datas []model.OtherData sum int wg *sync.WaitGroup } func handler(s string) []model.Example { // 解析column8为Example切片的逻辑 return nil } func (t *Task) Do() { for _, d := range t.datas { examples := handler(d.column8) if err := example.BatchUpsert(ctx, examples); err != nil { log.Errorf(err.Error()) } } t.wg.Done() } func main() { taskFunc := func(data interface{}) { task := data.(*Task) task.Do() } // 从other_table获取100条数据 nums := make([]model.OtherData, 100) p, _ := ants.NewPoolWithFunc(10, taskFunc) defer p.Release() var wg sync.WaitGroup wg.Add(DataSize / DataPerTask) tasks := make([]*Task, 0, DataSize/DataPerTask) for i := 0; i < DataSize/DataPerTask; i++ { task := &Task{ index: i + 1, datas: nums[i*DataPerTask : (i+1)*DataPerTask], wg: &wg, } tasks = append(tasks, task) p.Invoke(task) } wg.Wait() }
错误信息
Deadlock found when trying to get lock; try restarting transaction [11.305ms] [rows:0] INSERT INTO `example_table` (`column1`,`column2`,`column3`,`deleted_at`,`column4`,`column5`) VALUES (),(),(),(),()... ON DUPLICATE KEY UPDATE `column1`=VALUES(`column1`),`column2`=VALUES(`column2`),`column3`=VALUES(`column3`),`deleted_at`=VALUES(`deleted_at`),`column4`=VALUES(`column4`),`column5`=VALUES(`column5`) ... Similar deadlock errors ...
MySQL死锁状态
mysql> show engine innodb status; ... ------------------------ LATEST DETECTED DEADLOCK ------------------------ 2023-03-01 14:20:33 0x7f2be4264700 *** (1) TRANSACTION: TRANSACTION 188178, ACTIVE 0 sec inserting mysql tables in use 1, locked 1 LOCK WAIT 10 lock struct(s), heap size 1136, 5 row lock(s), undo log entries 1 MySQL thread id 46334, OS thread handle 139826491008768, query id 1037015 10.86.40.52 root update INSERT INTO `example_table` (`column1`,`column2`,`column3`,`deleted_at`,`column4`,`column5`) VALUES (?,?,?,?,?,?),(?,?,?,?,?,?),(?,?,?,?,?,?),(?,?,?,?,?,?),(?,?,?,?,,,,, *** (1) WAITING FOR THIS LOCK TO BE GRANTED: RECORD LOCKS space id 191 page no 17 n bits 88 index PRIMARY of table `examle_db`.`example_table` trx id 188178 lock_mode X insert intention waiting Record lock, heap no 1 PHYSICAL RECORD: n_fields 1; compact format; info bits 0 0: len 8; hex 73757072656d756d; asc supremum;; *** (2) TRANSACTION: TRANSACTION 188179, ACTIVE 0 sec inserting mysql tables in use 1, locked 1 8 lock struct(s), heap size 1136, 6 row lock(s), undo log entries 1 MySQL thread id 46338, OS thread handle 139826488035072, query id 1037016 10.86.40.52 root update INSERT INTO `example_table` (`column1`,`column2`,`column3`,`deleted_at`,`column4`,`column5`) VALUES (?,?,?,?,?,?),(?,?,?,?,?,?),(?,?,?,?,?,?),(?,?,?,?,?,?),(?,?,?,?,,,,, *** (2) HOLDS THE LOCK(S): RECORD LOCKS space id 191 page no 17 n bits 88 index PRIMARY of table `examle_db`.`example_table` trx id 188179 lock_mode X Record lock, heap no 1 PHYSICAL RECORD: n_fields 1; compact format; info bits 0 0: len 8; hex 73757072656d756d; asc supremum;; *** (2) WAITING FOR THIS LOCK TO BE GRANTED: RECORD LOCKS space id 191 page no 17 n bits 88 index PRIMARY of table `examle_db`.`example_table` trx id 188179 lock_mode X insert intention waiting Record lock, heap no 1 PHYSICAL RECORD: n_fields 1; compact format; info bits 0 0: len 8; hex 73757072656d756d; asc supremum;; *** WE ROLL BACK TRANSACTION (2) ------------ TRANSACTIONS ------------ ...
使用版本
- go version: 1.16
- gorm.io/driver/mysql v1.3.2
- gorm.io/gorm v1.23.2
- gorm.io/plugin/soft_delete v1.2.0
原因分析
- Upsert的锁机制冲突:InnoDB执行
INSERT ... ON DUPLICATE KEY UPDATE时,先尝试获取插入意向锁,遇唯一键冲突会转为排他锁。并行场景下,多个事务可能交叉等待对方释放锁,触发死锁。 - 主键索引的Supremum锁竞争:死锁日志显示两个事务都在等待主键索引末尾的虚拟Supremum记录的插入意向锁,同时其中一个事务持有该记录的排他锁。批量插入时InnoDB会锁定主键索引范围,并行批量插入同一区域数据时,锁冲突概率陡增。
- 数据拆分无顺序保证:当前任务拆分数据未按主键或唯一键排序,不同任务的批量插入数据在索引上的位置随机,加剧了锁竞争的可能性。
解决方案
- 给批量插入数据排序:调用
BatchUpsert前,将examples按唯一键(column1+column2+column3)或主键排序,保证InnoDB按统一顺序获取锁,避免交叉等待:// 给examples按唯一键排序示例 sort.Slice(examples, func(i, j int) bool { if examples[i].Column1 != examples[j].Column1 { return examples[i].Column1 < examples[j].Column1 } if examples[i].Column2 != examples[j].Column2 { return examples[i].Column2 < examples[j].Column2 } return examples[i].Column3 < examples[j].Column3 }) - 增加死锁重试逻辑:捕获MySQL死锁错误码(1213),自动重试Upsert操作,可搭配指数退避策略:
func (e *example) BatchUpsertWithRetry(ctx context.Context, es []*model.Example, maxRetry int) error { var err error for i := 0; i < maxRetry; i++ { err = e.db.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "column1"}, {Name: "column2"}, {Name: "column3"}}, UpdateAll: true, }).Create(&es).Error if err == nil { return nil } // 检测是否为死锁错误 if mysqlErr, ok := err.(*mysql.MySQLError); ok && mysqlErr.Number == 1213 { time.Sleep(time.Duration(i+1)*100 * time.Millisecond) continue } return err } return err } - 降低并发度:减少协程池的并发数,比如将
ants.NewPoolWithFunc(10, taskFunc)调整为ants.NewPoolWithFunc(5, taskFunc),降低同时执行的Upsert事务数量。 - 对齐唯一键定义:当前
OnConflict指定的列不包含deleted_at,但表的唯一键包含该字段。若业务中Upsert不会修改deleted_at,可将OnConflict的Columns调整为包含deleted_at,让唯一键匹配更精准,减少不必要的锁竞争:clause.OnConflict{ Columns: []clause.Column{ {Name: "column1"}, {Name: "column2"}, {Name: "column3"}, {Name: "deleted_at"}, }, UpdateAll: true, }
内容的提问来源于stack exchange,提问作者moluzhui
相关产品推荐
相关产品推荐

