Files
rootandClaude Opus 5 8679200f41 Initial commit
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-15 13:57:10 +08:00

144 lines
5.0 KiB
Go

package statuser
import (
"fmt"
"time"
"91porn-server/common/db"
"91porn-server/common/log"
"91porn-server/common/timeutil/timerange"
"91porn-server/models/s/statrecordmod"
"91porn-server/models/s/statusermod"
"91porn-server/skd/skdg"
)
func RunStat(defaultStatDate time.Time) {
log.Info("user stat running...")
defer log.Info("user stat finished...")
defaultTimeMap := map[string]time.Time{
string(statusermod.VidIncome): defaultStatDate,
string(statusermod.UploadCount): defaultStatDate,
}
//从历史提交记录中获取上次记录时间
lastStatAtMap, err := statrecordmod.LastRecordTimeMap(statrecordmod.UserStatJob, defaultTimeMap) //第一次从默认记录时间开始统计
if err != nil {
log.Error("user stat statuser LastRecordTimeMap faild!", log.E(err))
return
}
now := time.Now()
dayHead := timerange.LocDayRange(now).Head
for item, lastStatAt := range lastStatAtMap {
_lastStatAt := lastStatAt
//过去几天的数据按一日为单位统计,帮助数据恢复
for dayHead.After(_lastStatAt) {
end := _lastStatAt.AddDate(0, 0, 1)
//防止时间超出范围
if end.After(dayHead) {
end = dayHead
}
subTimeRange := timerange.TimeRange{ //连续子切片
Head: _lastStatAt,
Tail: end,
}
if err = statByTimeRange(item, subTimeRange); err != nil {
log.Error("user stat statuser statByTimeRange 1 faild!", log.E(err))
return
}
_lastStatAt = end
}
lastStatAtMap[item] = _lastStatAt
}
recentMinute := timerange.RecentMinute(now, statrecordmod.FiveMinuteScale) //对齐本次统计时间 Minute % frequency == 0
//拆分时间区间并依次提交,控制数据库读写压力
for item, lastStatAt := range lastStatAtMap {
_lastStatAt := lastStatAt
for recentMinute.After(_lastStatAt) {
end := _lastStatAt.Add(statrecordmod.FiveMinuteScale * time.Minute)
//防止时间超出范围
if end.After(recentMinute) {
end = recentMinute
}
subTimeRange := timerange.TimeRange{ //连续子切片
Head: _lastStatAt,
Tail: end,
}
if err = statByTimeRange(item, subTimeRange); err != nil {
log.Error("user stat statuser statByTimeRange 2 faild!", log.E(err))
return
}
_lastStatAt = end
}
lastStatAtMap[item] = _lastStatAt
}
}
func statByTimeRange(item string, timeRange timerange.TimeRange) error {
method, ok := methodMap[Itemtype(item)]
if !ok {
return nil
}
if method.MethodType == IncType {
//装配INC基础统计数据
incDocList, err := method.GetIncDocListMethod(timeRange)
if err != nil {
log.Error("user stat statByTimeRange GetIncDocListMethod faild!", log.Any("item", item), log.E(err))
return err
}
//提交Inc统计数据
log.Info("user stat statByTimeRange IncCommitMethod committing...", log.Any("item", item))
if err = method.IncCommitMethod(incDocList, item, timeRange.Tail); err != nil {
log.Error("user stat statByTimeRange IncCommitMethod faild!", log.Any("item", item), log.E(err))
return err
}
log.Info("user stat statByTimeRange IncCommitMethod committed.", log.Any("item", item))
} else if method.MethodType == SetType {
//装配SET基础统计数据
setDocList, err := method.GetSetDocListMethod(timeRange)
if err != nil {
log.Error("user stat statByTimeRange GetSetDocListMethod faild!", log.Any("item", item), log.E(err))
return err
}
//提交SET统计数据
log.Info("user stat statByTimeRange SetCommitMethod committing...", log.Any("item", item))
if err = method.SetCommitMethod(setDocList, item, timeRange.Tail); err != nil {
log.Error("user stat statByTimeRange SetCommitMethod faild!", log.Any("item", item), log.E(err))
return err
}
log.Info("user stat statByTimeRange SetCommitMethod committed.", log.Any("item", item))
}
return nil
}
func commitIncDoc(incDocList []IncDoc, item string, recordAt time.Time) error {
opt := (&db.TransOpts{}).SetReEntry(10)
return skdg.StatDB.Trans(func(trans *db.MongoTool) error {
if len(incDocList) != 0 {
if err := statusermod.ChangeIncStatTrans(trans, incDocList); err != nil {
return fmt.Errorf("user stat commitCount ChangeIncStatTrans faild!, err:%+v\n", err)
}
}
//提交本次修改记录
if err := statrecordmod.UpsertOneItemTrans(trans, statrecordmod.UserStatJob, item, recordAt); err != nil {
return fmt.Errorf("user stat commitIncDoc UpsertOneItemTrans faild!, err:%+v\n", err)
}
return nil
}, opt)
}
func commitSetDoc(setDocList []SetDoc, item string, recordAt time.Time) error {
opt := (&db.TransOpts{}).SetReEntry(10)
return skdg.StatDB.Trans(func(trans *db.MongoTool) error {
if len(setDocList) != 0 {
err := statusermod.ChangeSetStatTrans(trans, setDocList)
if err != nil {
return fmt.Errorf("user stat commitCount ChangeSetStatTrans faild!, err:%+v\n", err)
}
}
//提交本次修改记录
if err := statrecordmod.UpsertOneItemTrans(trans, statrecordmod.UserStatJob, item, recordAt); err != nil {
return fmt.Errorf("user stat commitSetDoc UpsertOneItemTrans faild!, item: %s+v, err:%+v\n", item, err)
}
return nil
}, opt)
}