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

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

原因分析
  1. Upsert的锁机制冲突:InnoDB执行INSERT ... ON DUPLICATE KEY UPDATE时,先尝试获取插入意向锁,遇唯一键冲突会转为排他锁。并行场景下,多个事务可能交叉等待对方释放锁,触发死锁。
  2. 主键索引的Supremum锁竞争:死锁日志显示两个事务都在等待主键索引末尾的虚拟Supremum记录的插入意向锁,同时其中一个事务持有该记录的排他锁。批量插入时InnoDB会锁定主键索引范围,并行批量插入同一区域数据时,锁冲突概率陡增。
  3. 数据拆分无顺序保证:当前任务拆分数据未按主键或唯一键排序,不同任务的批量插入数据在索引上的位置随机,加剧了锁竞争的可能性。

解决方案
  1. 给批量插入数据排序:调用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
    })
    
  2. 增加死锁重试逻辑:捕获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
    }
    
  3. 降低并发度:减少协程池的并发数,比如将ants.NewPoolWithFunc(10, taskFunc)调整为ants.NewPoolWithFunc(5, taskFunc),降低同时执行的Upsert事务数量。
  4. 对齐唯一键定义:当前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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 15:37:10