在 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 的能力,支持 ObjectId、Date、Decimal128、Binary、Regex 等类型。在 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:变更类型,包括insert、update、replace、delete等documentKey:被变更文档的_idns:命名空间信息{ db: "myapp", coll: "users" }updateDescription:update 事件特有,包含updatedFields和removedFields
断点续传: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 组合使用时,false 和 0 会被忽略,导致无法将该字段重置为零值:
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.M与bson.D的差异、struct tag 映射、自定义编解码器 - CRUD:单条/批量插入、条件查询与游标管理、更新操作符、Upsert、删除
- 聚合管道:
$match、$group、$lookup、$facet等 stage 的组合运用 - 事务:
session.WithTransaction的正确用法、事务选项配置与性能注意事项 - Change Streams:集合级监听、完整文档获取、Resume Token 断点续传
- 最佳实践:context 超时控制、错误分类处理、零值陷阱防范、批量写入优化、优雅关闭
建议结合具体 QPS 要求和数据规模,对连接池参数和查询模式进行基准测试,找到最适合你的配置组合。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。