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

使用mongo-go driver的Go应用MongoDB连接泄漏问题排查

问题描述

我正在开发一个集成MongoDB的Go应用,用于实现图片的创建、读取和删除操作。为测试目的将连接池最大连接数设为4,但发现连接数最高可达10,即使未调用任何MongoDB相关接口,连接数也维持在2-3之间。代码存在连接泄漏问题,我无法正确释放资源使其返回连接池。

根据文档说明,mongo-go driver是goroutine安全的,因此我编写了初始化MongoDB的函数,将初始化后的Mongo client设为全局变量:

package persistence

import (
    "context"
    "go.mongodb.org/mongo-driver/bson"
    "go.mongodb.org/mongo-driver/mongo"
    "go.mongodb.org/mongo-driver/mongo/options"
    "time"
)

var MongoClient *mongo.Client

func InitMongo(ctx context.Context, URI string) error {
    if MongoClient != nil {
        return nil
    }
    serverAPI := options.ServerAPI(options.ServerAPIVersion1)
    opts := options.Client().ApplyURI(URI).SetServerAPIOptions(serverAPI)

    opts.SetMinPoolSize(2)
    opts.SetMaxPoolSize(4)
    opts.SetMaxConnIdleTime(2 * time.Second)

    // Create a new client and connect to the server
    client, err := mongo.Connect(ctx, opts)
    if err != nil {
        return err
    }
    MongoClient = client

    // Send a ping to confirm a successful connection
    var result bson.M
    if err = client.Database("admin").RunCommand(context.TODO(), bson.D{{"ping", 1}}).Decode(&result); err != nil {
        return err
    }

    return nil

}

我在main函数中调用该初始化函数:

func main() {
    configuration := config.NewEnvConfigs()

    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()
    err = persistence.InitMongo(ctx, configuration.MongoDBURI)
    if err != nil {
        fmt.Println("Mongo not ready")
        panic(err)
    }
    fmt.Println("mongo connected")
    defer func() {
        if err = persistence.MongoClient.Disconnect(context.TODO()); err != nil {
            panic(err)
        }
    }()
    r := server.Routes()
    http.ListenAndServe("0.0.0.0:3333", r)
}

使用Mongo client的代码是图片接口的处理器:

package handlers

import (
    "context"
    "encoding/base64"
    "errors"
    "fmt"
    "github.com/bycultivaet/backend/internal/infrastructure/persistence"
    "github.com/go-chi/chi"
    "github.com/go-chi/render"
    "github.com/google/uuid"
    "go.mongodb.org/mongo-driver/bson"
    "go.mongodb.org/mongo-driver/bson/primitive"
    "io/ioutil"
    "mime/multipart"
    "net/http"
)

func GetImageByID() http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        ctx, cancel := context.WithCancel(r.Context())
        defer cancel()
        Uuid := chi.URLParam(r, "uuid")
        _, err := uuid.Parse(Uuid)
        if err != nil {
            render.Status(r, 400)
            render.Respond(w, r, errors.New("invalid uuid").Error())
            return
        }
        filter := bson.M{"uuid": Uuid}
        var result bson.M
        collection := persistence.MongoClient.Database("hassad-media").Collection("images")
        err = collection.FindOne(ctx, filter).Decode(&result)
        if err != nil {
            render.Status(r, 404)
            render.Respond(w, r, errors.New("image not found").Error())
            return
        }
        imageData, ok := result["image"].(primitive.Binary)
        if !ok {
            render.Status(r, 500)
            render.Respond(w, r, errors.New("error parsing image").Error())
            return
        }
        imageBase64 := base64.StdEncoding.EncodeToString(imageData.Data)
        render.Status(r, 200)
        render.JSON(w, r, map[string]string{"image": imageBase64})
    }
}

func UploadImage() http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        ctx, cancel := context.WithCancel(r.Context())
        defer cancel()
        cancelled := r.Context().Done()
        select {
        case <-cancelled:
            render.Status(r, 499)
            render.Respond(w, r, errors.New("request cancelled").Error())
            return
        default:
            err := r.ParseMultipartForm(10 << 20)
            if err != nil {
                render.Status(r, http.StatusBadRequest)
                render.Respond(w, r, err.Error())
                return
            }
            file, _, err := r.FormFile("image")
            if err != nil {
                render.Status(r, http.StatusBadRequest)
                render.Respond(w, r, err.Error())
                return
            }
            defer file.Close()
            Uuid, err := uuid.NewRandom()
            if err != nil {
                render.Status(r, http.StatusInternalServerError)
                render.Respond(w, r, errors.New("error generating uuid").Error())
                return
            }
            // Pass the contents of the file to GetMongoDB
            _, err = UploadPhoto(file, Uuid.String(), ctx)
            if err != nil {
                render.Status(r, http.StatusInternalServerError)
                render.Respond(w, r, err.Error())
                return
            }

            render.Status(r, http.StatusOK)
            render.JSON(w, r, map[string]string{"uuid": Uuid.String()})

        }
    }
}

func UploadPhoto(file multipart.File, uuid string, ctx context.Context) (interface{}, error) {
    imageBytes, err := ioutil.ReadAll(file)
    if err != nil {
        return nil, err
    }
    imageDoc := bson.M{"image": imageBytes, "uuid": uuid}
    defer file.Close()
    collection := persistence.MongoClient.Database("hassad-media").Collection("images")
    data, err := collection.InsertOne(ctx, imageDoc)
    if err != nil {
        return 0, err
    }

    fmt.Println("Inserted image into MongoDB!")
    fmt.Println(data.InsertedID)
    return data.InsertedID, nil
}

func DeleteImageByID() http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        ctx, cancel := context.WithCancel(r.Context())
        defer cancel()

        uuid := chi.URLParam(r, "uuid")
        filter := bson.M{"uuid": uuid}
        collection := persistence.MongoClient.Database("hassad-media").Collection("images")
        res, err := collection.DeleteOne(ctx, filter)
        if err != nil {
            if res.DeletedCount == 0 {
                render.Status(r, 404)
                render.Respond(w, r, errors.New("image not found").Error())
                return
            }
            render.Status(r, 500)
            render.Respond(w, r, err.Error())
            return
        }
        render.JSON(w, r, map[string]string{"answer": "deleted"})
    }
}

我在连接释放或配置上哪里出现了问题?我通过Python脚本高频调用接口进行测试,也曾尝试将collection和database设为全局变量,但问题依旧。


分析与解决

1. 未调用接口时的2-3个连接是正常行为

你设置了SetMinPoolSize(2),这意味着mongo-go driver会维持至少2个空闲连接在池中,即使没有任何请求,driver也会保持这些连接以避免频繁创建/销毁连接的开销。偶尔出现3个连接可能是driver内部的临时操作(比如心跳检测),属于正常现象,并非连接泄漏。

2. 连接数超过maxPoolSize的核心原因

SetMaxPoolSize(4)的作用是限制单个MongoDB节点的连接数,而非整个集群的总连接数。如果你连接的是MongoDB副本集(包含多个节点),driver会为每个节点维护独立的连接池,总连接数会是节点数 × maxPoolSize。比如3节点副本集的话,总连接数最多可达12,和你看到的10个连接的情况吻合。

3. 代码中的潜在优化点

(1)替换无期限的context.TODO()

初始化时的Ping操作使用了context.TODO(),这是一个无期限的context,可能导致连接被长时间占用。建议替换为初始化时传入的ctx:

if err = client.Database("admin").RunCommand(ctx, bson.D{{"ping", 1}}).Decode(&result); err != nil {
    return err
}

同理,main函数中的Disconnect操作也应使用带超时的context,避免阻塞:

defer func() {
    disconnectCtx, disconnectCancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer disconnectCancel()
    if err = persistence.MongoClient.Disconnect(disconnectCtx); err != nil {
        panic(err)
    }
}()

(2)升级mongo-go driver版本

旧版本的mongo-go driver存在连接池管理的bug,可能导致连接无法正确归还到池。建议升级到最新的稳定版本(go get go.mongodb.org/mongo-driver/mongo@latest)。

(3)添加连接池监控排查问题

可以通过添加CommandMonitor来监控连接的使用情况,明确连接被哪些操作占用:

import "go.mongodb.org/mongo-driver/event"

// 在InitMongo函数中添加监控
opts.SetMonitor(&event.CommandMonitor{
    Started: func(ctx context.Context, evt *event.CommandStartedEvent) {
        fmt.Printf("Command %s started on connection %d\n", evt.CommandName, evt.ConnectionID)
    },
    Succeeded: func(ctx context.Context, evt *event.CommandSucceededEvent) {
        fmt.Printf("Command %s succeeded on connection %d\n", evt.CommandName, evt.ConnectionID)
    },
    Failed: func(ctx context.Context, evt *event.CommandFailedEvent) {
        fmt.Printf("Command %s failed on connection %d\n", evt.CommandName, evt.ConnectionID)
    },
})

4. 其他注意事项

  • 你的handler中对context的使用是正确的:每个请求都创建独立的context并defer cancel,确保操作完成后连接能及时归还到池。
  • 将collection/database设为全局变量不会影响连接池行为,因为它们只是client的轻量引用,不会创建新连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 03:47:06