Go + MongoDB 驱动实战:连接池、事务与 Change Streams

使用官方 mongo-driver 在 Go 中操作 MongoDB,覆盖 CRUD、聚合、事务、Change Streams

在 Go 生态中操作 MongoDB,官方提供的 mongo-driver(即 go.mongodb.org/mongo-driver)是唯一经得起生产环境考验的选择。它由 MongoDB 官方团队维护,完全遵循 MongoDB Wire Protocol,提供了从底层 BSON 编解码到高层 CRUD、聚合、事务、Change Streams 的完整能力栈。本文将以生产级视角,系统讲解如何在 Go 中使用 mongo-driver 构建健壮、高性能的数据访问层。

安装与环境准备

mongo-driver 要求 Go 1.18 或更高版本。在模块管理的项目中,只需一条命令即可引入:

go get go.mongodb.org/mongo-driver/v2/mongo

从 v2 版本开始,mongo-driver 将连接管理、BSON 编解码和操作 API 统一到了更清晰的包结构下。本文所有示例均基于 v2 系列编写。如果你的项目仍在使用 v1 版本,大部分 API 概念仍然适用,但建议迁移到 v2 以获得更好的性能和维护支持。

在开始编码前,确保你有一个可连接的 MongoDB 实例。本地开发推荐以副本集模式启动(事务和 Change Streams 需要此环境):

docker run -d -p 27017:27017 --name mongodb mongo:7 --replSet rs0
docker exec -it mongodb mongosh --eval "rs.initiate()"

连接与连接池配置

mongo-client 的客户端设计遵循对象池模式mongo.Client 内部管理 TCP 连接池,整个应用生命周期内应只创建一个实例并复用。

基础连接

package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "go.mongodb.org/mongo-driver/v2/bson"
    "go.mongodb.org/mongo-driver/v2/mongo"
    "go.mongodb.org/mongo-driver/v2/mongo/options"
)

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()

    client, err := mongo.Connect(options.Client().ApplyURI("mongodb://localhost:27017"))
    if err != nil {
        log.Fatal(err)
    }
    defer client.Disconnect(ctx)

    if err = client.Ping(ctx, nil); err != nil {
        log.Fatal(err)
    }
    fmt.Println("MongoDB 连接成功")
}

mongo.Connect 返回的 *mongo.Client 是连接管理的核心对象。defer client.Disconnect 确保程序退出时优雅关闭连接池。所有 mongo-driver 的 API 都要求传入 context.Context,用于控制操作超时和取消,这也是保证服务稳定性的关键。

连接字符串进阶

生产环境中,连接字符串通常包含更多参数。以下是常见的连接字符串示例:

// 单节点,带认证
uri := "mongodb://user:pass@localhost:27017/mydb?authSource=admin"

// 三节点副本集
uri := "mongodb://user:pass@node1:27017,node2:27017,node3:27017/mydb?replicaSet=rs0&w=majority"

// 连接 MongoDB Atlas
uri := "mongodb+srv://user:pass@cluster0.mongodb.net/mydb?retryWrites=true&w=majority"

关键参数说明:

  • authSource=admin:指定认证数据库,MongoDB 用户信息通常存储在 admin 库中
  • replicaSet=rs0:副本集名称,驱动会据此自动发现主从节点拓扑
  • w=majority:写操作等待大多数节点确认后才返回成功,保证数据持久性
  • retryWrites=true:自动重试可幂等的写操作,应对瞬时网络抖动
  • readPreference=secondaryPreferred:优先从 Secondary 节点读取,分担主节点压力

连接池精细化配置

options.Client() 返回的 *options.ClientOptions 允许你精细控制连接池行为。在生产环境中,默认配置往往不能满足需求:

clientOpts := options.Client().
    ApplyURI("mongodb://localhost:27017").
    SetMaxPoolSize(100).              // 连接池最大连接数,默认 100
    SetMinPoolSize(10).               // 连接池最小保留连接数,默认 0
    SetMaxConnIdleTime(30 * time.Second). // 连接最大空闲时间
    SetServerSelectionTimeout(5 * time.Second). //  server 选择超时
    SetConnectTimeout(10 * time.Second).        //  单次连接建立超时
    SetSocketTimeout(0).                        //  socket 读写超时,0 表示不限制
    SetHeartbeatInterval(10 * time.Second)      //  心跳检测间隔

client, err := mongo.Connect(clientOpts)

调优建议:

  • MaxPoolSize:一般设置为并发 goroutine 数量的 1-2 倍,过大的连接池会造成内存浪费和连接管理开销
  • MinPoolSize:设置为 MaxPoolSize 的 10%-20%,避免请求突发时发生冷启动
  • MaxConnIdleTime:30 秒到 5 分钟比较合理,过短导致频繁建连,过长占用资源
  • ServerSelectionTimeout:网络环境复杂时可适当调高,故障转移时决定等待可用服务器的最长时间
  • SocketTimeout:默认不限制,建议为大数据量查询设置上限(如 60 秒),避免 goroutine 长时间阻塞

BSON 类型系统与编解码

MongoDB 使用 BSON(Binary JSON)而非纯 JSON 存储数据。BSON 扩展了 JSON 的能力,支持 ObjectIdDateDecimal128BinaryRegex 等类型。在 Go 中操作 MongoDB 时,你需要熟练掌握 mongo-driver 提供的 BSON 工具包。

BSON 基础类型

最核心的几个类型:

import "go.mongodb.org/mongo-driver/v2/bson"

// ObjectId:MongoDB 文档默认主键
oid, err := bson.ObjectIDFromHex("64a1b2c3d4e5f6a7b8c9d0e1")

// 生成新的 ObjectId
newOid := bson.NewObjectID()

// DateTime:毫秒时间戳
dt := bson.NewDateTimeFromTime(time.Now())

// Decimal128:高精度小数
dec, err := bson.ParseDecimal128("123.456")

文档表示方式

mongo-driver 提供了多种表示 BSON 文档的方式:

// bson.D:有序的键值对切片,确保字段顺序(对复合索引查询很重要)
doc := bson.D{
    {Key: "name", Value: "Alice"},
    {Key: "age", Value: 30},
    {Key: "created_at", Value: bson.NewDateTimeFromTime(time.Now())},
}

// bson.M:无序的 map,写起来更简洁,但不保证字段顺序
filter := bson.M{"status": "active", "age": bson.M{"$gte": 18}}

// bson.A:BSON 数组
arr := bson.A{"item1", "item2", 42}

// bson.E:单个键值对,是 bson.D 的元素类型
elem := bson.E{Key: "score", Value: 99.5}

何时用 bson.D,何时用 bson.M

  • bson.D:需要确保字段顺序时使用,例如复合索引查询,字段顺序必须和索引定义一致才能命中索引。bson.M 遍历顺序随机,可能导致无法使用索引
  • bson.M:普通等值查询和更新操作,不关心字段顺序,追求代码简洁

Struct 标签映射

最常用、最类型安全的方式是定义 Go struct,通过 bson tag 映射 MongoDB 字段:

type User struct {
    ID        bson.ObjectID `bson:"_id,omitempty"`
    Name      string        `bson:"name"`
    Email     string        `bson:"email"`
    Age       int           `bson:"age"`
    Tags      []string      `bson:"tags,omitempty"`
    IsActive  bool          `bson:"is_active"`
    CreatedAt time.Time     `bson:"created_at"`
    Score     float64       `bson:"score,omitempty"`
    Profile   Profile       `bson:"profile,omitempty"`
}

type Profile struct {
    Bio      string `bson:"bio"`
    Avatar   string `bson:"avatar"`
    Location string `bson:"location,omitempty"`
}

tag 关键选项:

  • bson:"_id":映射到 MongoDB 的 _id 字段
  • omitempty:字段为零值时插入/更新不包含,避免覆盖已有数据
  • -(短横线):完全忽略该字段

Marshal 与 Unmarshal

user := User{
    ID:        bson.NewObjectID(),
    Name:      "Bob",
    Email:     "bob@example.com",
    Age:       28,
    Tags:      []string{"go", "mongodb"},
    IsActive:  true,
    CreatedAt: time.Now(),
}

bsonBytes, err := bson.Marshal(user)
// ...
var decoded User
err = bson.Unmarshal(bsonBytes, &decoded)

自定义编解码器

有时候默认的编解码行为不满足需求,你可以注册自定义编解码器。例如将 Go 的 time.Time 统一以 Unix 时间戳形式存储:

import (
    "go.mongodb.org/mongo-driver/v2/bson"
    "reflect"
    "time"
)

type unixTimeCodec struct{}

func (utc unixTimeCodec) EncodeValue(_ bson.EncodeContext, vw bson.ValueWriter, val reflect.Value) error {
    t := val.Interface().(time.Time)
    return vw.WriteInt64(t.Unix())
}

func (utc unixTimeCodec) DecodeValue(_ bson.DecodeContext, vr bson.ValueReader, val reflect.Value) error {
    i, err := vr.ReadInt64()
    if err != nil {
        return err
    }
    val.Set(reflect.ValueOf(time.Unix(i, 0)))
    return nil
}

// 注册到 ClientOptions
reg := bson.NewRegistry()
reg.RegisterTypeDecoder(reflect.TypeOf(time.Time{}), bson.ValueDecoderFunc(unixTimeCodec{}.DecodeValue))
reg.RegisterTypeEncoder(reflect.TypeOf(time.Time{}), bson.ValueEncoderFunc(unixTimeCodec{}.EncodeValue))

clientOpts := options.Client().ApplyURI("mongodb://localhost:27017").SetRegistry(reg)

自定义编解码器通常用于处理遗留数据格式或对时间/货币的精确控制。

CRUD 操作实战

获取数据库和集合对象是所有数据操作的第一步:

db := client.Database("myapp")
coll := db.Collection("users")

Collection 对象是线程安全的,你可以在多个 goroutine 中并发使用同一个 *mongo.Collection 实例。

创建:InsertOne 与 InsertMany

// 插入单条文档
func insertOne(ctx context.Context, coll *mongo.Collection) (*mongo.InsertOneResult, error) {
    user := User{
        ID:        bson.NewObjectID(),
        Name:      "Alice",
        Email:     "alice@example.com",
        Age:       25,
        IsActive:  true,
        CreatedAt: time.Now(),
    }

    result, err := coll.InsertOne(ctx, user)
    if err != nil {
        return nil, err
    }
    return result, nil
}

// 批量插入
func insertMany(ctx context.Context, coll *mongo.Collection) (*mongo.InsertManyResult, error) {
    users := []interface{}{
        User{Name: "Bob", Email: "bob@example.com", Age: 30, CreatedAt: time.Now()},
        User{Name: "Charlie", Email: "charlie@example.com", Age: 35, CreatedAt: time.Now()},
        User{Name: "Diana", Email: "diana@example.com", Age: 28, CreatedAt: time.Now()},
    }

    result, err := coll.InsertMany(ctx, users)
    if err != nil {
        return nil, err
    }
    return result, nil
}

注意 InsertMany 的第一个参数是 []interface{},这意味着你可以混合不同类型插入到同一个集合。但在实际项目中,建议保持集合内文档结构的一致性。

查询:FindOne 与 Find

// 根据 ID 查询单条文档
func findByID(ctx context.Context, coll *mongo.Collection, id bson.ObjectID) (*User, error) {
    var user User
    err := coll.FindOne(ctx, bson.M{"_id": id}).Decode(&user)
    if err == mongo.ErrNoDocuments {
        return nil, fmt.Errorf("文档不存在")
    }
    if err != nil {
        return nil, err
    }
    return &user, nil
}

// 条件查询多条记录,带排序和分页
func findUsers(ctx context.Context, coll *mongo.Collection, status string, page, pageSize int) ([]User, error) {
    // 构造过滤条件
    filter := bson.M{"is_active": true}
    if status != "" {
        filter["status"] = status
    }

    // 选项:排序 + 分页 + 投影
    opts := options.Find().
        SetSort(bson.D{{Key: "created_at", Value: -1}}).
        SetSkip(int64((page - 1) * pageSize)).
        SetLimit(int64(pageSize)).
        SetProjection(bson.M{"password": 0}) // 排除敏感字段

    cursor, err := coll.Find(ctx, filter, opts)
    if err != nil {
        return nil, err
    }
    defer cursor.Close(ctx)

    var users []User
    if err = cursor.All(ctx, &users); err != nil {
        return nil, err
    }
    return users, nil
}

关键要点:

  • FindOne 返回 *mongo.SingleResult,通过 .Decode(&target) 反序列化。如果未找到文档,err 会是 mongo.ErrNoDocuments,这是正常业务情况,不要当作致命错误处理
  • Find 返回 *mongo.Cursor,必须 defer cursor.Close(ctx) 释放资源
  • cursor.All(ctx, &users) 会一次性将所有结果加载到内存,适合中小数据量。大数据集应该逐条迭代:for cursor.Next(ctx) { cursor.Decode(&doc) }
  • SetProjection 可以控制返回字段,1 表示包含,0 表示排除。不能在同一文档中混用包含和排除(_id 除外)

更新:UpdateOne、UpdateMany、ReplaceOne

// 更新单条文档的指定字段
func updateUserAge(ctx context.Context, coll *mongo.Collection, id bson.ObjectID, newAge int) error {
    filter := bson.M{"_id": id}
    update := bson.M{
        "$set": bson.M{
            "age":        newAge,
            "updated_at": time.Now(),
        },
    }

    result, err := coll.UpdateOne(ctx, filter, update)
    if err != nil {
        return err
    }
    if result.MatchedCount == 0 {
        return fmt.Errorf("未找到匹配的文档")
    }
    return nil
}

// 批量更新
func activateUsers(ctx context.Context, coll *mongo.Collection, minAge int) (int64, error) {
    filter := bson.M{"age": bson.M{"$gte": minAge}, "is_active": false}
    update := bson.M{"$set": bson.M{"is_active": true, "activated_at": time.Now()}}

    result, err := coll.UpdateMany(ctx, filter, update)
    if err != nil {
        return 0, err
    }
    return result.ModifiedCount, nil
}

// 完全替换文档(保留 _id)
func replaceUser(ctx context.Context, coll *mongo.Collection, id bson.ObjectID, newUser User) error {
    filter := bson.M{"_id": id}
    _, err := coll.ReplaceOne(ctx, filter, newUser)
    return err
}

更新操作符速查:

  • $set:设置字段值,不存在则新增
  • $unset:删除字段
  • $inc:原子自增(如 $inc: { "views": 1 }
  • $push:向数组添加元素
  • $pull:从数组移除元素
  • $addToSet:向数组添加不重复元素
  • $rename:重命名字段

原子自增是高并发计数场景下的利器:

coll.UpdateOne(ctx, bson.M{"_id": id},
    bson.M{"$inc": bson.M{"view_count": 1}})

删除:DeleteOne 与 DeleteMany

// 删除单条
func deleteUser(ctx context.Context, coll *mongo.Collection, id bson.ObjectID) error {
    result, err := coll.DeleteOne(ctx, bson.M{"_id": id})
    if err != nil {
        return err
    }
    if result.DeletedCount == 0 {
        return fmt.Errorf("文档不存在")
    }
    return nil
}

// 批量删除(软删除通常更好,但业务确实需要物理删除时)
func deleteInactiveUsers(ctx context.Context, coll *mongo.Collection, before time.Time) (int64, error) {
    result, err := coll.DeleteMany(ctx, bson.M{
        "is_active":   false,
        "created_at":  bson.M{"$lt": before},
    })
    if err != nil {
        return 0, err
    }
    return result.DeletedCount, nil
}

Upsert:不存在则插入

opts := options.UpdateOne().SetUpsert(true)
result, err := coll.UpdateOne(ctx,
    bson.M{"email": "alice@example.com"},
    bson.M{"$set": bson.M{"name": "Alice Updated", "updated_at": time.Now()}},
    opts,
)
if result.UpsertedCount > 0 {
    fmt.Printf("新插入文档,ID: %v\n", result.UpsertedID)
}

Upsert 非常适合 “有则更新,无则创建” 的幂等操作。

Count 与 Distinct

// 计数
count, err := coll.CountDocuments(ctx, bson.M{"is_active": true})

// 估算计数(基于元数据,速度快但不精确,适合大致了解规模)
estCount, err := coll.EstimatedDocumentCount(ctx)

// 获取某个字段的所有不重复值
 vals, err := coll.Distinct(ctx, "tags", bson.M{"is_active": true})

聚合管道 Aggregate

当简单的 CRUD 无法满足需求时,聚合管道(Aggregation Pipeline)是 MongoDB 最强大的查询工具。在 Go 中,聚合管道通过 bson.A[]bson.D 构建,传递给 coll.Aggregate 方法。

基础聚合:分组统计

统计每个年龄段的用户数量:

pipeline := mongo.Pipeline{
    {{Key: "$match", Value: bson.M{"is_active": true}}},
    {{Key: "$group", Value: bson.M{
        "_id": "$age",
        "count": bson.M{"$sum": 1},
        "avgScore": bson.M{"$avg": "$score"},
    }}},
    {{Key: "$sort", Value: bson.M{"count": -1}}},
    {{Key: "$limit", Value: 10}},
}

cursor, err := coll.Aggregate(ctx, pipeline)
if err != nil {
    log.Fatal(err)
}
defer cursor.Close(ctx)

var results []bson.M
if err = cursor.All(ctx, &results); err != nil {
    log.Fatal(err)
}

for _, r := range results {
    fmt.Printf("年龄 %v: %d 人, 平均分数 %.2f\n", r["_id"], r["count"], r["avgScore"])
}

注意这里使用了 mongo.Pipeline,它是 []bson.D 的别名,能够确保每个 stage 内部的字段顺序。

关联查询:$lookup

MongoDB 3.2 引入的 $lookup 相当于 SQL 的 LEFT JOIN:

// 订单集合关联用户集合
pipeline := mongo.Pipeline{
    {{Key: "$match", Value: bson.M{"status": "paid"}}},
    {{Key: "$lookup", Value: bson.M{
        "from":         "users",
        "localField":   "user_id",
        "foreignField": "_id",
        "as":           "user_info",
    }}},
    {{Key: "$unwind", Value: "$user_info"}},         // 将数组展开为单对象
    {{Key: "$project", Value: bson.M{
        "order_no":   1,
        "amount":     1,
        "status":     1,
        "user_name":  "$user_info.name",
        "user_email": "$user_info.email",
    }}},
}

cursor, err := orderColl.Aggregate(ctx, pipeline)

$lookup 返回的关联结果默认是数组形式。使用 $unwind 将其展开后,可以直接在 $project 中引用子字段。

分页优化:$facet

当需要同时返回列表数据和总数时,使用 $facet 可以在一次查询中完成:

pipeline := mongo.Pipeline{
    {{Key: "$match", Value: bson.M{"is_active": true}}},
    {{Key: "$facet", Value: bson.M{
        "data": bson.A{
            bson.M{"$sort": bson.M{"created_at": -1}},
            bson.M{"$skip": 20},
            bson.M{"$limit": 10},
        },
        "total": bson.A{
            bson.M{"$count": "count"},
        },
    }}},
}

cursor, err := coll.Aggregate(ctx, pipeline)
// 结果结构: [{ data: [...], total: [{count: N}] }]

这比两次查询更高效,MongoDB 只需扫描一次数据。

时间窗口聚合

统计每小时的日志数量:

pipeline := mongo.Pipeline{
    {{Key: "$match", Value: bson.M{
        "created_at": bson.M{"$gte": time.Now().Add(-24 * time.Hour)},
    }}},
    {{Key: "$group", Value: bson.M{
        "_id": bson.M{
            "$dateToString": bson.M{"format": "%Y-%m-%d %H:00", "date": "$created_at"},
        },
        "count": bson.M{"$sum": 1},
    }}},
    {{Key: "$sort", Value: bson.M{"_id": 1}}},
}

事务:session.WithTransaction

MongoDB 4.0 开始支持多文档 ACID 事务(副本集环境),4.2 扩展到了分片集群。在 Go 中使用事务需要通过 Session 对象管理。

副本集事务基础

func transferBalance(ctx context.Context, client *mongo.Client, fromID, toID bson.ObjectID, amount float64) error {
    // 使用 WithTransaction 自动处理提交与回滚
    session, err := client.StartSession()
    if err != nil {
        return err
    }
    defer session.EndSession(ctx)

    callback := func(sessCtx mongo.SessionContext) (interface{}, error) {
        accounts := client.Database("bank").Collection("accounts")

        // 扣款
        _, err := accounts.UpdateOne(sessCtx,
            bson.M{"_id": fromID, "balance": bson.M{"$gte": amount}},
            bson.M{"$inc": bson.M{"balance": -amount}},
        )
        if err != nil {
            return nil, err
        }

        // 加款
        _, err = accounts.UpdateOne(sessCtx,
            bson.M{"_id": toID},
            bson.M{"$inc": bson.M{"balance": amount}},
        )
        if err != nil {
            return nil, err
        }

        return nil, nil
    }

    _, err = session.WithTransaction(ctx, callback, nil)
    return err
}

WithTransaction 会自动处理重试逻辑,回调中发生可重试错误(如主节点切换)时驱动会自动回滚并重试,最多重试一次。回调函数接收的 sessCtx 绑定了当前事务会话,所有操作必须使用它,否则不会被纳入事务。

事务选项配置

你可以通过 options.TransactionOptions 精细控制事务行为:

txnOpts := options.Transaction().
    SetReadConcern(readconcern.Snapshot()).
    SetWriteConcern(writeconcern.New(writeconcern.WMajority())).
    SetReadPreference(readpref.Primary())

_, err = session.WithTransaction(ctx, callback, txnOpts)
  • ReadConcern.Snapshot:读取事务开始前的一致快照,类似 MVCC,避免幻读
  • WriteConcern.WMajority:写入操作等待副本集多数节点确认
  • ReadPreference.Primary:事务中的读操作必须走主节点,Secondary 默认不参与事务读

事务注意事项

事务使用原则:

  • 保持事务简短,长时间运行会占用 WiredTiger 快照和 oplog 空间
  • 避免事务内做非数据库操作(HTTP 调用、文件读写等),拉长事务持有时间
  • 异常自动回滚,回调返回非 nil error 即触发回滚,无需手动 AbortTransaction
  • 重试机制有限WithTransaction 只在特定错误时重试一次
  • 单文档原子性不需要事务,只有跨文档一致性保证时才使用事务
  • 最大事务运行时间默认 60 秒,超过会被强制中止

Change Streams:实时数据监听

Change Streams 是 MongoDB 3.6 引入的特性,允许应用程序实时监听数据库的变更事件(insert/update/delete 等),是实现实时推送、数据同步、审计日志等功能的利器。

基础监听

func watchChanges(ctx context.Context, coll *mongo.Collection) error {
    // 监听集合级别的变更
    pipeline := mongo.Pipeline{
        {{Key: "$match", Value: bson.M{
            "operationType": bson.M{"$in": bson.A{"insert", "update", "delete"}},
        }}},
    }

    opts := options.ChangeStream().SetFullDocument(options.UpdateLookup)
    stream, err := coll.Watch(ctx, pipeline, opts)
    if err != nil {
        return err
    }
    defer stream.Close(ctx)

    fmt.Println("开始监听变更...")
    for stream.Next(ctx) {
        var changeDoc bson.M
        if err := stream.Decode(&changeDoc); err != nil {
            log.Printf("解码变更事件失败: %v", err)
            continue
        }

        opType := changeDoc["operationType"]
        docID := changeDoc["documentKey"]
        fullDoc := changeDoc["fullDocument"]

        fmt.Printf("操作类型: %v, 文档ID: %v\n", opType, docID)
        if fullDoc != nil {
            fmt.Printf("完整文档: %v\n", fullDoc)
        }
    }

    if err := stream.Err(); err != nil {
        return fmt.Errorf("Change Stream 异常: %w", err)
    }
    return nil
}

关键参数:

  • SetFullDocument(options.UpdateLookup):update 操作默认不返回完整文档,设置后驱动会自动再查一次
  • operationType:变更类型,包括 insertupdatereplacedelete
  • documentKey:被变更文档的 _id
  • ns:命名空间信息 { db: "myapp", coll: "users" }
  • updateDescription:update 事件特有,包含 updatedFieldsremovedFields

断点续传:Resume Token

网络中断或服务重启后,可以使用 Resume Token 从断点继续监听,避免丢失变更事件:

type ChangeStreamState struct {
    ResumeToken bson.Raw `bson:"resume_token"`
}

func watchWithResume(ctx context.Context, coll *mongo.Collection, stateColl *mongo.Collection) error {
    var state ChangeStreamState
    stateColl.FindOne(ctx, bson.M{}).Decode(&state)

    opts := options.ChangeStream()
    if state.ResumeToken != nil {
        opts.SetResumeAfter(state.ResumeToken)
        fmt.Println("从断点恢复监听...")
    }

    stream, err := coll.Watch(ctx, mongo.Pipeline{}, opts)
    if err != nil {
        return err
    }
    defer stream.Close(ctx)

    for stream.Next(ctx) {
        var event bson.M
        if err := stream.Decode(&event); err != nil {
            log.Printf("解码失败: %v", err)
            continue
        }

        // 处理业务逻辑...
        fmt.Printf("收到变更: %v\n", event["operationType"])

        // 保存 Resume Token
        token := stream.ResumeToken()
        _, err = stateColl.UpdateOne(ctx,
            bson.M{},
            bson.M{"$set": bson.M{"resume_token": token}},
            options.UpdateOne().SetUpsert(true),
        )
        if err != nil {
            log.Printf("保存 Resume Token 失败: %v", err)
        }
    }

    return stream.Err()
}

Resume Token 建议持久化到独立于监听目标的数据库中。

监听数据库或实例级别变更

// 整个数据库
dbStream, err := db.Watch(ctx, mongo.Pipeline{})

// 整个实例(需要 readAnyDatabase 权限)
clientStream, err := client.Watch(ctx, mongo.Pipeline{})

最佳实践

正确使用 context.Context

所有 mongo-driver 的操作都接受 context.Context,这是保证服务韧性的第一道防线:

// HTTP handler 中复用请求 context
func (h *Handler) GetUser(w http.ResponseWriter, r *http.Request) {
    ctx := r.Context() // 包含请求生命周期
    user, err := h.userRepo.FindByID(ctx, userID)
    if err != nil {
        // 如果客户端断开连接,context 会被取消,操作会自动终止
        http.Error(w, err.Error(), 500)
        return
    }
    json.NewEncoder(w).Encode(user)
}

// 为数据库操作设置独立超时
func (r *Repo) FindByID(ctx context.Context, id bson.ObjectID) (*User, error) {
    ctx, cancel := context.WithTimeout(ctx, 3*time.Second)
    defer cancel()

    var user User
    err := r.coll.FindOne(ctx, bson.M{"_id": id}).Decode(&user)
    return &user, err
}

HTTP/RPC 服务中应传递请求自带的 context,客户端超时或取消时数据库操作会及时释放资源。

错误处理

import (
    "errors"
    "go.mongodb.org/mongo-driver/v2/mongo"
    "go.mongodb.org/mongo-driver/v2/mongo/writeconcern"
)

func handleMongoError(err error) error {
    if err == nil {
        return nil
    }
    if errors.Is(err, mongo.ErrNoDocuments) {
        return ErrNotFound
    }
    if mongo.IsDuplicateKeyError(err) {
        return ErrDuplicate
    }
    if mongo.IsTimeout(err) {
        return ErrTimeout
    }
    var wcErr *writeconcern.WriteConcernError
    if errors.As(err, &wcErr) {
        return fmt.Errorf("写入确认失败: %s", wcErr.Message)
    }
    return fmt.Errorf("数据库错误: %w", err)
}

连接生命周期管理

在生产应用中,mongo.Client 的生命周期通常与应用程序本身绑定:

type MongoStore struct {
    client *mongo.Client
    db     *mongo.Database
}

func NewMongoStore(uri string) (*MongoStore, error) {
    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()

    client, err := mongo.Connect(options.Client().ApplyURI(uri))
    if err != nil {
        return nil, err
    }

    if err := client.Ping(ctx, nil); err != nil {
        client.Disconnect(ctx)
        return nil, err
    }

    return &MongoStore{
        client: client,
        db:     client.Database("myapp"),
    }, nil
}

func (s *MongoStore) Close(ctx context.Context) error {
    return s.client.Disconnect(ctx)
}

func (s *MongoStore) UserCollection() *mongo.Collection {
    return s.db.Collection("users")
}

类型安全与零值陷阱

Go struct 的零值和 omitempty 组合使用时,false0 会被忽略,导致无法将该字段重置为零值:

type Config struct {
    Debug       bool   `bson:"debug,omitempty"`
    MaxRetries  int    `bson:"max_retries,omitempty"`
}

解决方案:使用指针类型

type Config struct {
    Debug      *bool   `bson:"debug,omitempty"`
    MaxRetries *int    `bson:"max_retries,omitempty"`
}

指针的零值是 nil,BSON 编码器能区分 “未设置” 和 “设置为零值”。

索引使用提示

聚合管道中 $match 始终放最前面,让查询优化器有机会使用索引:

// 正确
pipeline := mongo.Pipeline{
    {{Key: "$match", Value: bson.M{"status": "active"}}},
    {{Key: "$group", Value: ...}},
}

// 错误:$match 在 $group 后面,已全量展开
pipeline := mongo.Pipeline{
    {{Key: "$group", Value: ...}},
    {{Key: "$match", Value: ...}},
}

查询返回大数据量时,用 SetBatchSize 控制每次拉取量:

opts := options.Find().SetBatchSize(100)
cursor, err := coll.Find(ctx, filter, opts)

批量写入优化

InsertMany 批量大小建议控制在 1000-5000 条之间:

func batchInsert(ctx context.Context, coll *mongo.Collection, users []User) error {
    const batchSize = 1000
    for i := 0; i < len(users); i += batchSize {
        end := i + batchSize
        if end > len(users) {
            end = len(users)
        }
        batch := make([]interface{}, end-i)
        for j := range batch {
            batch[j] = users[i+j]
        }
        _, err := coll.InsertMany(ctx, batch)
        if err != nil {
            return fmt.Errorf("批量插入失败 [%d:%d]: %w", i, end, err)
        }
    }
    return nil
}

优雅关闭与信号处理

捕获系统信号,退出前完成数据库连接优雅关闭:

func main() {
    store, err := NewMongoStore("mongodb://localhost:27017")
    if err != nil {
        log.Fatal(err)
    }

    ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
    defer stop()

    srv := &http.Server{Addr: ":8080"}
    go func() { srv.ListenAndServe() }()

    <-ctx.Done()

    shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()

    srv.Shutdown(shutdownCtx)
    store.Close(shutdownCtx)
}

总结

mongo-driver 是 Go 语言操作 MongoDB 的官方利器,它不仅提供了完整的 CRUD 能力,还通过连接池管理、事务支持、Change Streams 监听等高级特性,让开发者能在生产环境中构建可靠的数据访问层。

本文覆盖的核心要点:

  • 连接管理:单例 Client、连接池参数调优、连接字符串参数
  • BSON 编解码bson.Mbson.D 的差异、struct tag 映射、自定义编解码器
  • CRUD:单条/批量插入、条件查询与游标管理、更新操作符、Upsert、删除
  • 聚合管道$match$group$lookup$facet 等 stage 的组合运用
  • 事务session.WithTransaction 的正确用法、事务选项配置与性能注意事项
  • Change Streams:集合级监听、完整文档获取、Resume Token 断点续传
  • 最佳实践:context 超时控制、错误分类处理、零值陷阱防范、批量写入优化、优雅关闭

建议结合具体 QPS 要求和数据规模,对连接池参数和查询模式进行基准测试,找到最适合你的配置组合。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. 缓存架构演进之路:从单机 Redis 到亿级分布式多级缓存体系
  2. Redis 7.x 重大新特性与架构升级深度解析
  3. Redis 消息队列深度对比:Pub/Sub、Streams 与 Kafka/RabbitMQ 选型指南