如何使用Golang实现MongoDB集合的月度自动轮转归档
MongoDB集合每月自动轮转的Golang实现
要实现resultData集合的每月自动轮转,核心是通过定时任务触发集合重命名,并重建原集合以继续接收新数据。以下是具体实现方案:
核心逻辑
每月固定时间(比如每月1日凌晨)执行以下操作:
- 将
resultData重命名为resultData_YYYYMM格式(YYYY是年份,MM是两位月份) - 检查目标集合是否已存在,避免重复操作报错
- 创建新的空
resultData集合,并同步原集合的索引、约束等配置
依赖准备
需要两个核心依赖:
- 官方MongoDB Go驱动:
go.mongodb.org/mongo-driver/mongo - 定时任务库:
github.com/robfig/cron/v3(支持灵活的 cron 表达式配置)
完整代码实现
主程序与定时任务配置
package main import ( "context" "fmt" "time" "github.com/robfig/cron/v3" "go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo/options" ) func main() { // 初始化MongoDB客户端 clientOpts := options.Client().ApplyURI("mongodb://localhost:27017") client, err := mongo.Connect(context.TODO(), clientOpts) if err != nil { panic(fmt.Sprintf("MongoDB connect failed: %v", err)) } defer func() { if err := client.Disconnect(context.TODO()); err != nil { fmt.Printf("MongoDB disconnect failed: %v\n", err) } }() // 验证连接 if err := client.Ping(context.TODO(), nil); err != nil { panic(fmt.Sprintf("MongoDB ping failed: %v", err)) } // 初始化cron调度器,使用秒级精度 c := cron.New(cron.WithSeconds()) // 配置每月1日凌晨0点执行轮转任务(cron表达式:秒 分 时 日 月 周) _, err = c.AddFunc("0 0 0 1 * *", func() { rotateCollection(client, "your_db_name", "resultData") }) if err != nil { panic(fmt.Sprintf("Add cron task failed: %v", err)) } c.Start() fmt.Println("Collection rotation scheduler started.") select {} // 阻塞主进程,保持服务运行 }
集合轮转核心函数
func rotateCollection(client *mongo.Client, dbName, collName string) { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) defer cancel() db := client.Database(dbName) // 生成带年月后缀的目标集合名 now := time.Now() targetColl := fmt.Sprintf("%s_%d%02d", collName, now.Year(), now.Month()) // 检查目标集合是否已存在,避免重复操作 exists, err := checkCollectionExists(ctx, db, targetColl) if err != nil { fmt.Printf("Check collection existence failed: %v\n", err) return } if exists { fmt.Printf("Collection %s already exists, skip rotation.\n", targetColl) return } // 重命名原集合 err = db.Collection(collName).Rename(ctx, targetColl, nil) if err != nil { fmt.Printf("Rename collection failed: %v\n", err) return } fmt.Printf("Successfully renamed %s to %s\n", collName, targetColl) // 重建原集合并同步索引 if err := recreateCollectionWithIndexes(ctx, db, collName, targetColl); err != nil { fmt.Printf("Recreate collection with indexes failed: %v\n", err) return } fmt.Printf("Successfully created new %s collection with indexes synced.\n", collName) } // 检查集合是否存在 func checkCollectionExists(ctx context.Context, db *mongo.Database, collName string) (bool, error) { colls, err := db.ListCollectionNames(ctx, map[string]interface{}{"name": collName}) if err != nil { return false, err } return len(colls) > 0, nil } // 重建集合并同步原集合的索引 func recreateCollectionWithIndexes(ctx context.Context, db *mongo.Database, newCollName, oldCollName string) error { // 获取原集合的索引 oldColl := db.Collection(oldCollName) cursor, err := oldColl.Indexes().List(ctx) if err != nil { return err } var indexes []mongo.IndexModel if err := cursor.All(ctx, &indexes); err != nil { return err } // 创建新集合(如果需要配置capped、validator等,在这里设置CreateCollectionOptions) createOpts := options.CreateCollection() if err := db.CreateCollection(ctx, newCollName, createOpts); err != nil { return err } // 同步索引到新集合 newColl := db.Collection(newCollName) if _, err := newColl.Indexes().CreateMany(ctx, indexes); err != nil { return err } return nil }
生产环境注意事项
- 定时任务可靠性:如果服务可能重启,建议使用分布式定时任务或者在启动时检查是否错过当月轮转任务,手动补执行。
- 写入一致性:轮转过程中,原集合被重命名后,新写入会失败直到新集合创建完成。可以在轮转前短暂暂停写入流量,或者使用MongoDB的事务(仅副本集/分片集群支持)保证操作原子性。
- 集合配置同步:如果原集合有特殊配置(如
capped集合、文档校验器validator),需要在CreateCollectionOptions中同步这些配置。 - 告警机制:在错误处理逻辑中加入告警(如邮件、企业微信通知),确保轮转失败时能及时发现。
内容的提问来源于stack exchange,提问作者Ayush
相关产品推荐
相关产品推荐

