package kwstatmod import ( "fmt" "time" "91porn-server/common/constant" "91porn-server/common/db" "91porn-server/common/log" "91porn-server/common/timeutil/timerange" "91porn-server/models" "github.com/jinzhu/now" "go.mongodb.org/mongo-driver/bson" "go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo/options" ) var mdb *db.MongoDB const table = models.KeywordStat func coll(t *db.MongoTool) *db.MongoTool { if t == nil { return mdb.Coll(table) } return t.Coll(table) } // InitIndex 设置index func initIndex() { many := []mongo.IndexModel{ { Keys: bson.D{{Key: "word", Value: 1}}, }, { Keys: bson.D{{Key: "realm", Value: 1}, {Key: "word", Value: 1}}, }, { Keys: bson.D{{Key: "sumDate", Value: -1}, {Key: "realm", Value: 1}}, }, { Keys: bson.D{{Key: "sumDate", Value: -1}, {Key: "realm", Value: 1}, {Key: "word", Value: 1}}, Options: options.Index().SetUnique(true), }, { Keys: bson.D{{Key: "count", Value: -1}}, }, { Keys: bson.D{{Key: "createdAt", Value: -1}}, }, { Keys: bson.D{{Key: "updatedAt", Value: -1}}, }, } if _, err := coll(nil).CreateIndex(many); err != nil { panic(fmt.Sprintf("%s model set index err ==>[%+v]", table, err)) } } func ChangeStatTrans(trans *db.MongoTool, position time.Time, recordTime time.Time, docList []KeywordIncDoc) error { sumDate := timerange.LocDayRange(position).Head // create the slice of write models writes := make([]mongo.WriteModel, len(docList)) for i, doc := range docList { filter := M{ "sumDate": sumDate, "realm": doc.Realm, "word": doc.Word, } update := M{ "$setOnInsert": M{ "sumDate": sumDate, "realm": doc.Realm, "word": doc.Word, "createdAt": time.Now(), }, "$set": M{ "recordAt": recordTime, "updatedAt": time.Now(), }, "$inc": M{ "count": doc.Count, }, } writes[i] = mongo.NewUpdateOneModel(). SetFilter(filter). SetUpdate(update). SetUpsert(true) log.Debug(fmt.Sprintf("sumDate:%s table:%s [ realm:%+v word:%+v filter:%+v update:%+v ]\n", sumDate, table, doc.Realm, doc.Word, filter, update)) } if len(writes) == 0 { return nil } //bulkWrite 不是原子操作 不具备事务性 opt := (&options.BulkWriteOptions{}).SetOrdered(false) //设为无序,触发并行写,提升写效率 _, err := coll(trans).Bulk(writes, opt) return err } func list(sumDate time.Time, realms []constant.RealmType, enable bool, sort bson.D, skip, limit int64) ([]Keyword, error) { filter := M{ "sumDate": sumDate, "realm": M{"$in": realms}, } realmCount := len(realms) if realmCount == 0 { return []Keyword{}, nil } if realmCount == 1 { //只有一个元素 不使用$in 提升效率 filter["realm"] = realms[0] } opt := (&options.FindOptions{}). SetSort(sort). SetSkip(skip). SetLimit(limit) keywords := make([]Keyword, 0, limit) if err := coll(nil).Find(&keywords, filter, opt); err != nil { log.ZapLog.Error(fmt.Sprintf("coll:%s list fail error:%+v:", table, err)) return nil, err } return keywords, nil } // RealmRanking RealmRanking func RealmRanking(sumDate time.Time, realm constant.RealmType, skip int64, limit int64) ([]Keyword, error) { sort := bson.D{ {Key: "sortKey", Value: -1}, {Key: "count", Value: -1}, } keywords, err := list(sumDate, []constant.RealmType{realm}, true, sort, skip, limit) if err != nil { log.ZapLog.Error(fmt.Sprintf("coll:%s RealmRanking fail error:%+v:", table, err)) return nil, err } return keywords, nil } // KeywordCountMapBySumDate func KeywordCountMapBySumDate(sumDate time.Time) (map[string]int64, error) { pipeline := []M{ { "$match": M{ "sumDate": sumDate, }, }, { "$group": M{ "_id": "$word", "count": M{"$sum": "$count"}, }, }, { "$sort": M{ "count": -1, }, }, } var list []struct { Word string `bson:"_id"` Count int64 `bson:"count"` } if err := coll(nil).Aggregate(&list, pipeline); err != nil { return nil, err } m := make(map[string]int64, len(list)) for _, v := range list { m[v.Word] = v.Count } return m, nil } // HotSearchWordsToday 获取今日搜索较高的关键词 视频关键字 func HotSearchWordsToday(top int) ([]string, error) { var kw []Keyword var words []string opts := options.FindOptions{} sort := bson.D{{Key: "count", Value: -1}} opts.SetSort(sort).SetLimit(int64(top)) if err := coll(nil).Find(&kw, M{"realm": "video", "sumDate": now.BeginningOfDay()}, &opts); err != nil { log.Error("keywordStat model HotSearchWordsToday error", log.E(err), log.Any("top", top)) return words, err } words = make([]string, len(kw)) for i, k := range kw { words[i] = k.Word } return words, nil } func GetKeywordsExcludeWords(sumDate time.Time, realm string, words []string, skip int64, limit int64) (data []Keyword, err error) { if words == nil { words = []string{} } query := bson.M{"word": bson.M{"$nin": words}} opts := options.Find() sort := bson.D{ {Key: "sortKey", Value: -1}, {Key: "count", Value: -1}, } opts.SetSort(sort).SetSkip(skip).SetLimit(limit) if err = coll(nil).Find(&data, query, opts); err != nil { log.Error("keywordStat model GetKeywordsExcludeWords error", log.E(err), log.Any("sumDate", sumDate), log.Any("words", words)) return } return } // 根据时间获取今日加昨日的关键字 func GetKeywordsBySumDate(sumDate time.Time, realm string, skip int64, limit int64) (data []Keyword, err error) { var query = bson.M{"sumDate": bson.M{"$gte": sumDate}, "realm": realm} var opts = options.Find() opts.SetSort(bson.D{{Key: "sortKey", Value: -1}, {Key: "count", Value: -1}}).SetSkip(skip).SetLimit(limit) if err = coll(nil).Find(&data, query, opts); err != nil { log.Error("keywordStat model GetKeywordsBySumDate error", log.E(err), log.Any("sumDate", sumDate), log.Any("realm", realm)) return } return } // 根据时间段获取关键字 func GetKeywordsByTimeRange(start time.Time, end time.Time) (data []Keyword, err error) { var query = bson.M{"updatedAt": bson.M{"$gte": start, "$lt": end}} if err = coll(nil).Find(&data, query); err != nil { log.Error("keywordStat model GetKeywordsByTimeRange error", log.E(err)) return } return } func DeleteBeforeSumDate(t *db.MongoTool, tm time.Time) error { _, err := coll(t).DeleteMany(bson.M{"sumDate": bson.M{"$lt": tm}}) return err }