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

Golang如何使用ants库实现循环中两张表的数据库插入操作

基于ants协程池的实现方案

ants是Go生态的高性能协程池,可控制并发量级避免循环中大量创建协程抢占系统资源,针对你的插入场景,可根据业务依赖选择两种实现策略:

策略1:按行并发(推荐,兼容事务依赖场景)

适合同一行数据的User和Address插入存在依赖(比如需要获取User自增ID回填到Address关联字段)的场景,每行数据作为一个独立任务提交到协程池,同一个任务内串行执行两次插入操作。

import (
    "sync"
    "github.com/panjf2000/ants/v2"
)

// 你的rows数据源定义在上方
var rows []YourRowStruct

func main() {
    // 初始化协程池,最大并发数根据数据库负载调整,推荐不超过数据库最大连接数的1/2
    pool, err := ants.NewPool(10)
    if err != nil {
        panic(err)
    }
    // 程序退出前释放协程池资源
    defer pool.Release()

    var wg sync.WaitGroup
    // 错误收集通道,缓存长度等于行数避免阻塞
    errChan := make(chan error, len(rows))

    for _, row := range rows {
        // 拷贝循环变量,避免闭包捕获问题
        row := row
        wg.Add(1)
        // 提交任务到协程池
        submitErr := pool.Submit(func() {
            defer wg.Done()
            // 插入User表
            user := User{
                Name:  row.Name,
                Email: row.Email,
            }
            if err := dm.Insert(&user); err != nil {
                errChan <- err
                return
            }
            // 插入Address表
            address := Address{
                Address1: row.Address1,
                Address2: row.Address2,
                PinCode:  row.PinCode,
                City:     row.City,
            }
            if err := dm.Insert(&address); err != nil {
                errChan <- err
                return
            }
        })
        if submitErr != nil {
            errChan <- submitErr
            wg.Done()
        }
    }

    // 等待所有任务执行完成
    wg.Wait()
    close(errChan)

    // 统一处理错误
    for err := range errChan {
        if err != nil {
            // 可根据业务需求调整错误处理逻辑,比如日志打印、回滚等
            panic(err)
        }
    }
}

策略2:按插入操作并发(性能更高,无依赖场景可用)

如果两次插入没有业务依赖和事务一致性要求,可以把每个插入操作拆成独立任务提交,进一步提升并发性能:

// 协程池初始化、错误收集逻辑和策略1一致
for _, row := range rows {
    row := row
    // 提交User插入任务
    wg.Add(1)
    if err := pool.Submit(func() {
        defer wg.Done()
        user := User{Name: row.Name, Email: row.Email}
        if err := dm.Insert(&user); err != nil {
            errChan <- err
        }
    }); err != nil {
        errChan <- err
        wg.Done()
    }

    // 提交Address插入任务
    wg.Add(1)
    if err := pool.Submit(func() {
        defer wg.Done()
        address := Address{
            Address1: row.Address1,
            Address2: row.Address2,
            PinCode:  row.PinCode,
            City:     row.City,
        }
        if err := dm.Insert(&address); err != nil {
            errChan <- err
        }
    }); err != nil {
        errChan <- err
        wg.Done()
    }
}
// 后续等待任务、错误处理逻辑和策略1一致

注意事项

  • 协程池最大并发数不要设置过高,避免压垮数据库,建议压测后调整到合适值
  • 如果需要保证同一行两次插入的事务一致性,要把两次插入放到同一个数据库事务中执行,不要拆分到两个协程
  • 循环内的row := row是必须操作,避免闭包捕获循环变量导致所有任务读取到最后一行的数据
  • 有超时控制要求的场景,可以在提交的任务函数中加入context超时逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 12:54:03