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

如何使用Golang实现MongoDB集合的月度自动轮转归档

MongoDB集合每月自动轮转的Golang实现

要实现resultData集合的每月自动轮转,核心是通过定时任务触发集合重命名,并重建原集合以继续接收新数据。以下是具体实现方案:

核心逻辑

每月固定时间(比如每月1日凌晨)执行以下操作:

  1. 将resultData重命名为resultData_YYYYMM格式(YYYY是年份,MM是两位月份)
  2. 检查目标集合是否已存在,避免重复操作报错
  3. 创建新的空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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 08:00:55