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

如何在多个Goroutine中复用MongoDB Go驱动的同一游标

多Goroutine安全读取MongoDB游标方案

MongoDB Go驱动的游标不是并发安全的,直接在多个Goroutine中调用cursor.Next()会引发竞争条件,导致丢失文档或出现异常。要实现多Goroutine处理文档,正确的做法是用单生产者多消费者模式:由一个Goroutine负责遍历游标并将文档发送到通道,多个Goroutine从通道接收数据并处理。

完整实现代码

package main

import (
	"context"
	"fmt"
	"sync"

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

// 定义你的文档结构体
type YourDocument struct {
	ID bson.ObjectID `bson:"_id,omitempty"`
	// 其他业务字段...
}

func main() {
	ctx := context.TODO()
	clientOpts := options.Client().ApplyURI("mongodb://localhost:27017")
	client, err := mongo.Connect(ctx, clientOpts)
	if err != nil {
		panic(err)
	}
	defer client.Disconnect(ctx)

	coll := client.Database("your_db").Collection("your_coll")
	filter := bson.D{} // 替换为你的查询条件
	opts := options.Find() // 替换为你的查询选项

	cursor, err := coll.Find(ctx, filter, opts)
	if err != nil {
		fmt.Println("Finding all documents ERROR:", err)
		return
	}
	defer cursor.Close(ctx) // 确保游标最终关闭

	// 定义通道传递文档,缓冲大小按需调整
	docChan := make(chan YourDocument, 100)
	var wg sync.WaitGroup

	// 启动多个消费者协程处理文档
	consumerCount := 5 // 消费者数量可根据需求调整
	wg.Add(consumerCount)
	for i := 0; i < consumerCount; i++ {
		go func(workerID int) {
			defer wg.Done()
			for doc := range docChan {
				// 这里编写你的文档处理逻辑
				fmt.Printf("Worker %d processed document: %v\n", workerID, doc.ID)
			}
		}(i)
	}

	// 生产者协程:遍历游标并发送文档到通道
	for cursor.Next(ctx) {
		var doc YourDocument
		if err := cursor.Decode(&doc); err != nil {
			fmt.Println("Decode document ERROR:", err)
			continue
		}
		docChan <- doc
	}

	// 检查游标遍历过程中是否出现错误
	if err := cursor.Err(); err != nil {
		fmt.Println("Cursor iteration ERROR:", err)
	}

	// 关闭通道,通知消费者无更多数据
	close(docChan)
	// 等待所有消费者处理完成
	wg.Wait()

	fmt.Println("All documents processed")
}

关键要点说明

  • 游标隔离:仅在单个生产者Goroutine中调用cursor.Next()和cursor.Decode(),彻底避免并发操作游标引发的竞争问题。
  • 通道分发:通过Go通道实现文档的安全传递,天然保证并发场景下的数据同步。
  • 资源管理:用defer cursor.Close(ctx)确保游标无论是否出错都能正常关闭;通过sync.WaitGroup等待所有消费者处理完毕,防止程序提前退出。
  • 错误检查:必须调用cursor.Err(),因为Next()返回false可能是正常遍历结束,也可能是过程中出现了错误。

注意事项

  • 通道缓冲大小可根据文档总量和处理速度调整,过小可能导致生产者阻塞,过大则会占用更多内存。
  • 消费者数量不宜超过CPU核心数过多,否则会增加上下文切换开销,反而降低处理效率。
  • 禁止在多个Goroutine中直接共享或操作游标,这违反MongoDB Go驱动的使用规范。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 18:18:21