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

使用Goroutine批量插入PostgreSQL数据异常:仅部分写入成功

问题描述

我尝试使用2000个Goroutine向PostgreSQL插入Item数据,运行了如下代码:

package main

import (
    "fmt"
    "github.com/jmoiron/sqlx"
    _ "github.com/lib/pq"
    "log"
    "sync"
)

type Item struct {
    Id          int    `db:"id"`
    Title       string `db:"title"`
    Description string `db:"description"`
}

// ConnectPostgresDB -> connect postgres db
func ConnectPostgresDB() *sqlx.DB {
    connstring := "user=postgres dbname=postgres sslmode=disable password=postgres host=localhost port=8080"
    db, err := sqlx.Open("postgres", connstring)
    if err != nil {
        fmt.Println(err)
        return db
    }
    return db
}

func InsertItem(item Item, wg *sync.WaitGroup) {
    defer wg.Done()
    db := ConnectPostgresDB()
    defer db.Close()
    tx, err := db.Beginx()
    if err != nil {
        fmt.Println(err)
        return
    }

    _, err = tx.Queryx("INSERT INTO items(id, title, description) VALUES($1, $2, $3)", item.Id, item.Title, item.Description)
    if err != nil {
        fmt.Println(err)
    }

    err = tx.Commit()
    if err != nil {
        fmt.Println(err)
        return
    }

    fmt.Println("Data is Successfully inserted!!")
}

func main() {
    var wg sync.WaitGroup
    //db, err := sqlx.Connect("postgres", "user=postgres dbname=postgres sslmode=disable password=postgres host=localhost port=8080")
    for i := 1; i <= 2000; i++ {
        item := Item{Id: i, Title: "TestBook", Description: "TestDescription"}
        //go GetItem(db, i, &wg)
        wg.Add(1)
        go InsertItem(item, &wg)

    }
    wg.Wait()
    fmt.Println("All DB Connection is Completed")
}

运行代码后,我预期items表中会有2000条数据,但实际仅存在150条。起初数据库报“too many clients”错误,我将max_connections从100调整到4000后再次尝试,但结果依旧。请问这是什么原因?


问题分析与解决方案

核心问题1:每个Goroutine重复创建连接池,资源耗尽

sqlx.Open创建的是连接池实例,不是单个数据库连接。你现在的代码里每个InsertItem Goroutine都新建一个连接池,会导致:

  • 瞬间创建大量数据库连接,即使调大max_connections,操作系统的文件句柄、网络套接字也会被耗尽,大量连接请求失败
  • 连接池的频繁创建/销毁带来巨大性能开销,很多插入请求还没执行就因为资源不足失败

核心问题2:错误处理不严谨,事务未回滚

  • 当tx.Queryx执行失败时,仅打印错误但未回滚事务,导致事务长期占用数据库资源,阻塞后续操作
  • 插入失败后没有标记、统计,程序直接继续执行,你无法感知到大量失败的插入请求

核心问题3:用错SQL执行方法

Queryx是为查询返回结果集设计的,INSERT属于写操作,应该用Execx,避免不必要的结果集处理开销。

修复后的代码示例

package main

import (
    "fmt"
    "github.com/jmoiron/sqlx"
    _ "github.com/lib/pq"
    "log"
    "sync"
)

type Item struct {
    Id          int    `db:"id"`
    Title       string `db:"title"`
    Description string `db:"description"`
}

// 全局复用一个连接池,初始化一次即可
var db *sqlx.DB

func init() {
    connstring := "user=postgres dbname=postgres sslmode=disable password=postgres host=localhost port=8080"
    var err error
    db, err = sqlx.Connect("postgres", connstring)
    if err != nil {
        log.Fatalf("无法连接数据库: %v", err)
    }
    // 根据服务器资源配置连接池参数
    db.SetMaxOpenConns(100)  // 最大打开连接数,建议和CPU核心数正相关
    db.SetMaxIdleConns(20)   // 最大空闲连接数
    db.SetConnMaxLifetime(0) // 连接生命周期,0表示永久有效
}

func InsertItem(item Item, wg *sync.WaitGroup, failed *int) {
    defer wg.Done()
    tx, err := db.Beginx()
    if err != nil {
        log.Printf("开启事务失败: %v", err)
        *failed++
        return
    }
    // 使用Execx执行INSERT操作,无需处理结果集
    _, err = tx.Execx("INSERT INTO items(id, title, description) VALUES($1, $2, $3)", item.Id, item.Title, item.Description)
    if err != nil {
        log.Printf("插入数据失败(id=%d): %v", item.Id, err)
        // 失败必须回滚事务,释放资源
        if rollbackErr := tx.Rollback(); rollbackErr != nil {
            log.Printf("回滚事务失败: %v", rollbackErr)
        }
        *failed++
        return
    }

    err = tx.Commit()
    if err != nil {
        log.Printf("提交事务失败(id=%d): %v", item.Id, err)
        *failed++
        return
    }

    fmt.Printf("数据插入成功(id=%d)!!\n", item.Id)
}

func main() {
    var wg sync.WaitGroup
    var failed int
    total := 2000
    for i := 1; i <= total; i++ {
        item := Item{Id: i, Title: "TestBook", Description: "TestDescription"}
        wg.Add(1)
        go InsertItem(item, &wg, &failed)
    }
    wg.Wait()
    fmt.Printf("所有操作完成,成功插入%d条,失败%d条\n", total-failed, failed)
}

额外优化建议

  • 连接池参数调优:SetMaxOpenConns不要设置过大,比如8核服务器设置100左右即可,过大反而会增加数据库上下文切换开销
  • 批量插入优化:如果要插入大量数据,建议分组批量插入(比如一次插入100条),大幅减少事务和连接的开销
  • 日志标准化:使用log包的格式化输出,方便定位具体失败的插入请求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 06:05:02