使用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
相关产品推荐
相关产品推荐

