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

将MongoDB去重聚合查询转换为Go语言实现的技术求助

MongoDB Go驱动实现删除templates集合重复数据的正确方案

核心问题修正

你之前的代码存在几个关键问题:

  • 聚合管道缺少原Shell逻辑中的$sort和$match阶段
  • 分组阶段字段名不匹配:原Shell用templateid你写成了templatesent,数组字段名dups写成了notification
  • 未开启allowDiskUse选项,大数据量场景下会执行失败

完整实现代码

package main

import (
	"context"
	"log"

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

// 定义聚合结果结构体,用于解析返回数据
type DuplicateTemplate struct {
	ID    bson.M                `bson:"_id"`
	Dups  []primitive.ObjectID `bson:"dups"`
	Count int                   `bson:"count"`
}

func main() {
	// 初始化MongoDB客户端(替换为你的连接配置)
	client, err := mongo.Connect(context.TODO(), options.Client().ApplyURI("mongodb://localhost:27017"))
	if err != nil {
		log.Fatal(err)
	}
	defer client.Disconnect(context.TODO())

	coll := client.Database("your_db_name").Collection("templates")
	ctx := context.TODO()

	// 1. 构建聚合管道,完全对齐原Shell逻辑
	sortStage := bson.D{
		{"$sort", bson.D{
			{"accountid", 1},
			{"sessionid", 1},
			{"templateid", 1},
		}},
	}

	groupStage := bson.D{
		{"$group", bson.D{
			{"_id", bson.D{
				{"accountid", "$accountid"},
				{"sessionid", "$sessionid"},
				{"templateid", "$templateid"},
			}},
			{"dups", bson.D{{"$push", "$_id"}}},
			{"count", bson.D{{"$sum", 1}}},
		}},
	}

	matchStage := bson.D{
		{"$match", bson.D{{"count", bson.D{{"$gt", 1}}}}},
	}

	pipeline := mongo.Pipeline{sortStage, groupStage, matchStage}

	// 2. 设置聚合选项:允许磁盘使用,处理大数据量
	aggOpts := options.Aggregate().SetAllowDiskUse(true)

	// 3. 执行聚合查询
	cursor, err := coll.Aggregate(ctx, pipeline, aggOpts)
	if err != nil {
		log.Fatalf("聚合查询失败: %v", err)
	}
	defer cursor.Close(ctx)

	// 4. 遍历结果并批量删除重复数据
	var duplicates []DuplicateTemplate
	if err = cursor.All(ctx, &duplicates); err != nil {
		log.Fatalf("解析聚合结果失败: %v", err)
	}

	for _, dup := range duplicates {
		// 保留第一条数据,删除其余重复项
		if len(dup.Dups) > 1 {
			idsToDelete := dup.Dups[1:]
			filter := bson.D{{"_id", bson.D{{"$in", idsToDelete}}}}
			result, err := coll.DeleteMany(ctx, filter)
			if err != nil {
				log.Printf("删除重复项失败: %v", err)
				continue
			}
			log.Printf("已删除 %d 条重复数据,重复键: %v", result.DeletedCount, dup.ID)
		}
	}
}

代码说明

  1. 聚合管道对齐原逻辑:
    • $sort:按三个字段升序排序,确保分组时重复项顺序一致
    • $group:按accountid/sessionid/templateid分组,收集重复文档的_id并统计数量
    • $match:仅筛选重复数量大于1的分组
  2. 性能优化:开启allowDiskUse避免内存溢出,用DeleteMany批量删除重复项,比逐条删除效率更高
  3. 类型安全:用结构体解析聚合结果,避免直接使用bson.M的类型不确定性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 12:24:53