Go中SQL连接池适配RabbitMQ高负载写入的优化问题
解决方案
1. 用带缓冲通道限制并发写入数
核心思路是通过缓冲通道作为「并发令牌桶」,令牌数量与数据库连接池的最大打开连接数(MaxOpenConns)保持一致。每次启动写入goroutine前先获取令牌,写入完成后释放令牌,确保同时执行的写入请求不超过连接池容量,从根源避免连接耗尽。
修改后的消费逻辑示例:
// 初始化令牌通道,容量等于连接池最大连接数 tokenChan := make(chan struct{}, db.Stats().MaxOpenConns) for m := range msgs { // 获取令牌,无可用令牌时阻塞,避免无限制创建goroutine tokenChan <- struct{}{} se := &sqlEntity{ body: string(m.Body), cnt: m.MessageCount, timeStamp: m.Timestamp.Format("2006-01-02 15:04:05"), // 无需额外fmt.Sprintf uuid: u, } go func(se *sqlEntity, msg amqp.Delivery) { defer func() { // 写入完成后释放令牌 <-tokenChan }() err := writeSQLWithRetry(se) if err != nil { // 写入失败时将消息重新入队(需避免死循环,可设置重试次数阈值) msg.Nack(false, true) log.Printf("write failed, requeue msg: %v", err) return } // 写入成功后确认消息 msg.Ack(false) }(se, m) }
2. 给数据库操作添加重试逻辑
遇到连接池耗尽、临时网络波动等可恢复错误时,直接丢弃数据是不合理的。通过有限次数的重试,配合指数退避策略,能有效降低数据丢失概率,同时避免短时间内重试压垮数据库。
func writeSQLWithRetry(se *sqlEntity) error { const maxRetries = 3 var err error for i := 0; i < maxRetries; i++ { err = writeSQL(se) if err == nil { return nil } // 判断是否为可重试错误(根据MS SQL驱动返回的错误特征调整) if isRetryable(err) { // 指数退避等待,避免频繁重试 time.Sleep(time.Duration(math.Pow(2, float64(i))) * 100 * time.Millisecond) continue } // 不可重试错误直接返回 return err } return fmt.Errorf("write failed after %d retries: %w", maxRetries, err) } func isRetryable(err error) bool { errStr := err.Error() // 匹配连接相关的错误关键词,可根据实际驱动返回信息调整 return strings.Contains(errStr, "connection") || strings.Contains(errStr, "timeout") || strings.Contains(errStr, "reset by peer") } func writeSQL(se *sqlEntity) error { _, err := db.Exec("INSERT INTO target_table (body, cnt, time_stamp, uuid) VALUES (?, ?, ?, ?)", se.body, se.cnt, se.timeStamp, se.uuid) return err }
3. 批量写入优化
百万级消息单条写入效率极低,且会频繁占用连接。将多条消息攒成一批后批量插入,能大幅减少数据库连接的使用次数,提升整体吞吐量,同时降低数据库负载。
示例批量处理逻辑:
const batchSize = 100 // 批量大小根据数据库性能调整 batchChan := make(chan *sqlEntity, batchSize*2) // 启动独立的批量写入goroutine go func() { var batch []*sqlEntity for se := range batchChan { batch = append(batch, se) if len(batch) >= batchSize { if err := writeSQLBatch(batch); err != nil { log.Printf("batch write failed: %v", err) // 批量失败时,需将这批消息对应的RabbitMQ消息重新入队 for _, item := range batch { item.delivery.Nack(false, true) } } else { // 批量成功后确认所有消息 for _, item := range batch { item.delivery.Ack(false) } } batch = batch[:0] // 重置批次 } } // 处理剩余不足一批的消息 if len(batch) > 0 { if err := writeSQLBatch(batch); err != nil { log.Printf("final batch write failed: %v", err) for _, item := range batch { item.delivery.Nack(false, true) } } else { for _, item := range batch { item.delivery.Ack(false) } } } }() // 消费逻辑改为发送消息到批量通道 for m := range msgs { se := &sqlEntity{ body: string(m.Body), cnt: m.MessageCount, timeStamp: m.Timestamp.Format("2006-01-02 15:04:05"), uuid: u, delivery: m, // 保存Delivery用于后续消息确认 } batchChan <- se }
批量写入SQL示例:
func writeSQLBatch(batch []*sqlEntity) error { if len(batch) == 0 { return nil } // 构建批量插入的占位符和参数 placeholders := make([]string, len(batch)) args := make([]interface{}, 0, len(batch)*4) for i, se := range batch { placeholders[i] = "(?, ?, ?, ?)" args = append(args, se.body, se.cnt, se.timeStamp, se.uuid) } sqlStmt := fmt.Sprintf("INSERT INTO target_table (body, cnt, time_stamp, uuid) VALUES %s", strings.Join(placeholders, ",")) _, err := db.Exec(sqlStmt, args...) return err }
4. RabbitMQ消费端流量控制
通过设置RabbitMQ的Prefetch Count,限制消费者同时持有的未确认消息数量,避免消息堆积过多导致内存压力过大。同时可结合批量通道的缓冲状态,动态调整消费速度。
// 设置Prefetch Count,建议值与批量大小或并发令牌数匹配 err := ch.Qos( 200, // 预取消息数量 0, // 预取大小(0表示无限制) false, // 是否应用于整个信道 ) if err != nil { log.Fatalf("failed to set QoS: %v", err) }
内容的提问来源于stack exchange,提问作者Oleg Sydorov
相关产品推荐
相关产品推荐

