package analyze import ( "context" "encoding/json" "strconv" "strings" "time" "yunyan/comm" "yunyan/lego/core" "yunyan/lego/core/cbase" "yunyan/lego/sys/log" redissys "yunyan/lego/sys/redis" "yunyan/pb" natssys "yunyan/sys/nats" "yunyan/utils" ) // syncComp 定时把 Redis 中的当日统计【按分区打包成绝对快照】推送到 NATS 统计管道, // 由 console 消费端覆盖落库 stats_global_day(console 为唯一写库方); // 已结束的历史日推送后顺手清理其 Redis 数据,开始新一天的统计。 type syncComp struct { cbase.ModuleCompBase module *Analyze options *Options cancel context.CancelFunc done chan struct{} } func (this *syncComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, opt core.IModuleOptions) (err error) { this.ModuleCompBase.Init(service, module, comp, opt) this.module = module.(*Analyze) this.options = opt.(*Options) this.done = make(chan struct{}) return } func (this *syncComp) Start() (err error) { if err = this.ModuleCompBase.Start(); err != nil { return } ctx, cancel := context.WithCancel(context.Background()) this.cancel = cancel go this.loop(ctx) // console 已是唯一写库方:未配 Subject(或缺 AppName/Region) 则统计无处可去——数据仍在本地 Redis // 累加不丢,但不会被推送/落库。显著告警提醒修正 modules.analyze 配置。 if this.options.Subject == "" || this.options.AppName == "" || this.options.Region == 0 { this.module.Errorf("analyze: 未配置转发(Subject=%q, AppName=%q, Region=%d),统计将只留在本地 Redis、不会推送到 console 落库", this.options.Subject, this.options.AppName, this.options.Region) } log.Info("analyze 统计同步器已启动", log.Field{Key: "prefix", Value: prefix()}, log.Field{Key: "subject", Value: this.options.Subject}, log.Field{Key: "interval", Value: this.options.SyncInterval}) return } // Destroy 取消 loop 并等待其结束,确保进行中的一轮同步整轮完成后再关停。 func (this *syncComp) Destroy() (err error) { if this.cancel != nil { this.cancel() select { case <-this.done: case <-time.After(10 * time.Second): this.module.Warn("analyze: 关停时等待同步结束超时") } } return this.ModuleCompBase.Destroy() } // loop 启动即推送一次,之后按 SyncInterval 周期推送 Redis 当日快照到 NATS,并在统计时区 // 0 点对齐触发一次跨日收尾。周期 tick 与 0 点 timer 同在一个 select 里,天然串行、互不并发; // 一轮 sync 进行中不会被打断;ctx 取消只在轮次之间生效——统计累计仍在 Redis,重启后由启动快照对账,不丢数据。 func (this *syncComp) loop(ctx context.Context) { defer close(this.done) this.sync() // 启动即同步一次 ticker := time.NewTicker(time.Duration(this.options.SyncInterval) * time.Second) defer ticker.Stop() // 0 点收尾 timer:对齐到统计时区下一个 0 点(+5s),跨日即把上一日落库并清缓存, // 不必等周期 tick(SyncInterval 较大时尤为重要)。每次触发后重置到下一个 0 点。 midnight := time.NewTimer(utils.DurationToNextMidnight(statLocation())) defer midnight.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: this.sync() case <-midnight.C: this.sync() midnight.Reset(utils.DurationToNextMidnight(statLocation())) } } } // sync 遍历 Redis 中所有统计分区,把每个分区的当日绝对聚合打包成快照推送到 NATS; // 历史日分区推送后清理 Redis。未配置转发(Subject 为空)则跳过推送,仅做历史分区清理。 func (this *syncComp) sync() { ctx := context.Background() members, err := redissys.Conn().SMembers(ctx, indexKey()).Result() if err != nil { this.module.Warn("analyze: 读取分区索引失败", log.Field{Key: "err", Value: err.Error()}) return } now := time.Now() today := dayNumber(now) // 昨日仍每轮按全量推送,但延后到「比昨天更早」才清缓存:多留一天吸收迟到事件, // 避免分区被清后又被一条迟到事件以单条增量重建、再被绝对值覆盖把当日全量冲小。 yesterday := dayNumber(now.AddDate(0, 0, -1)) ts := now.Unix() for _, m := range members { productId, statDay, ok := parsePartition(m) if !ok { this.module.Warn("analyze: 非法分区标识", log.Field{Key: "member", Value: m}) continue } row, ok := this.readRow(ctx, productId, statDay, ts) if !ok { // 分区已无累加数据(被清理过的索引残留),历史日则顺手摘除索引。 if statDay < today { redissys.Conn().SRem(ctx, indexKey(), m) } continue } this.publishSnapshot(row) if statDay < yesterday { this.cleanup(ctx, productId, statDay) } } } // publishSnapshot 把一行当日绝对聚合包成 comm.StatSnapshot 推送到 NATS 统计管道(异步、绝对值幂等)。 // 未配 Subject/AppName/Region 则不推送(数据仍在本地 Redis);失败即丢弃,由下一轮快照自愈。 func (this *syncComp) publishSnapshot(row *pb.StatsGlobalDay) { if this.options.Subject == "" || this.options.AppName == "" || this.options.Region == 0 { return } snap := &comm.StatSnapshot{Kind: comm.StatSnapshotKind, Ts: time.Now().Unix(), Row: row} data, err := json.Marshal(snap) if err != nil { this.module.Warn("analyze: 快照序列化失败,本轮丢弃", log.Field{Key: "err", Value: err.Error()}) return } if err := natssys.PublishAsync(this.options.Subject, data); err != nil { this.module.Warn("analyze: 推送统计快照失败,下轮重试", log.Field{Key: "err", Value: err.Error()}) return } // 推送日志:与 console 端「收到 NATS 消息」对照——home 打了这条但 console 收不到,即两端 NATS 不互通。 log.Infof("analyze: 已推送统计快照 subject=%s app=%s region=%s product=%d day=%d", this.options.Subject, row.AppId, row.Region, row.ProductId, row.StatDay) } // readRow 从 Redis 读出某分区的统计并组装成 DB 行;分区无累加数据时返回 ok=false。 func (this *syncComp) readRow(ctx context.Context, productId, statDay uint32, ts int64) (*pb.StatsGlobalDay, bool) { h, err := redissys.Conn().HGetAll(ctx, hashKey(productId, statDay)).Result() if err != nil { this.module.Warn("analyze: 读取统计哈希失败", log.Field{Key: "err", Value: err.Error()}) return nil, false } if len(h) == 0 { return nil, false } return &pb.StatsGlobalDay{ AppId: this.options.AppKey(), ProductId: productId, Region: comm.RegionCode(this.options.Region), StatDay: statDay, UpdateTime: ts, LoginUserCount: this.pfcount(ctx, productId, statDay, uLogin), LoginCount: hashInt(h, fLoginCount), NewUserCount: hashInt(h, fNewUserCount), ActiveUserCount: this.pfcount(ctx, productId, statDay, uActive), ActiveDeviceCount: hashInt(h, fActiveDeviceCount), BindDeviceCount: hashInt(h, fBindDeviceCount), QrActiveCount: hashInt(h, fQRActiveCount), OrderCreateCount: hashInt(h, fOrderCreateCount), OrderPaidCount: hashInt(h, fOrderPaidCount), OrderFailedCount: hashInt(h, fOrderFailedCount), OrderCreateAmount: hashInt(h, fOrderCreateAmount), OrderPaidAmount: hashInt(h, fOrderPaidAmount), PayUserCount: this.pfcount(ctx, productId, statDay, uPay), GrantVipDay: hashInt(h, fGrantVipDay), GrantAiIntegral: hashInt(h, fGrantAiIntegral), GrantTradeSecond: hashInt(h, fGrantTradeSecond), GrantMeetSecond: hashInt(h, fGrantMeetSecond), AiChatCount: hashInt(h, fAiChatCount), AiUpToken: hashInt(h, fAiUpToken), AiDownToken: hashInt(h, fAiDownToken), TradeCount: hashInt(h, fTradeCount), TradeSecond: hashInt(h, fTradeSecond), TradeWords: hashInt(h, fTradeWords), MeetCount: hashInt(h, fMeetCount), MeetSecond: hashInt(h, fMeetSecond), }, true } // pfcount 取某去重指标 HyperLogLog 的基数估计(误差约 0.81%)。 func (this *syncComp) pfcount(ctx context.Context, productId, statDay uint32, metric string) int64 { n, err := redissys.Conn().PFCount(ctx, hllKey(productId, statDay, metric)).Result() if err != nil { this.module.Warn("analyze: PFCOUNT 失败", log.Field{Key: "err", Value: err.Error()}) return 0 } return n } // cleanup 删除某历史分区在 Redis 中的全部数据。 func (this *syncComp) cleanup(ctx context.Context, productId, statDay uint32) { redissys.Conn().Del(ctx, hashKey(productId, statDay), hllKey(productId, statDay, uLogin), hllKey(productId, statDay, uActive), hllKey(productId, statDay, uPay)) redissys.Conn().SRem(ctx, indexKey(), partition(productId, statDay)) } // hashInt 从 HGETALL 结果里取一个整型字段,缺失或非法按 0 计。 func hashInt(h map[string]string, field string) int64 { n, _ := strconv.ParseInt(h[field], 10, 64) return n } // parsePartition 解析索引成员 "产品ID:统计日"。 func parsePartition(m string) (productId, statDay uint32, ok bool) { i := strings.IndexByte(m, ':') if i < 0 { return 0, 0, false } p, e1 := strconv.ParseUint(m[:i], 10, 32) d, e2 := strconv.ParseUint(m[i+1:], 10, 32) if e1 != nil || e2 != nil { return 0, 0, false } return uint32(p), uint32(d), true }