package console import ( "encoding/json" "sync" "time" "yunyan/comm" "yunyan/lego/core" "yunyan/lego/core/cbase" "yunyan/lego/sys/log" natssys "yunyan/sys/nats" "yunyan/utils" gonats "github.com/nats-io/nats.go" ) // statComp 统计汇总消费端,跑在 console 服务内(console 为 stats_global_day 的唯一写库方): // - 订阅各业务部署 analyze 推送到 NATS(JetStream) 的【绝对快照】(comm.StatSnapshot)—— // 每条即某 app×region×product×day 分区的当日绝对聚合; // - 校验通过(已注册应用、核心维度齐全)后,直接以绝对值覆盖 upsert 到 stats_global_day。 // // 绝对值幂等:重复/乱序投递、重启重放都收敛到最新值,不会重复计数;偶发丢消息由下一次快照补回, // 业务侧重启时的启动快照即完成对账。console 自身不再维护 Redis 累加/HLL(去重人数由业务侧算好随快照带来)。 // // 后台 dashboard 直接查 stats_global_day(见 api_stats.go),按 stat_day 区间 SUM 聚合。 type statComp struct { cbase.ModuleCompBase module *Console options *Options sub *gonats.Subscription nameMu sync.RWMutex regCache map[string]bool // 已注册应用名 → true;落库 app_id 直接用应用名 } // isRegisteredApp 校验快照里的应用名是否为「console 已注册应用」(落 stats_global_day.app_id)。 // 仅当注册表存在同名应用时返回 true 并缓存;未注册返回 false 且【不缓存】——以便其后在 console // 补注册即可被识别,不会永久误判。未注册应用的快照不落库,避免脏行污染统计。 func (this *statComp) isRegisteredApp(name string) bool { if name == "" { return false } this.nameMu.RLock() ok := this.regCache[name] this.nameMu.RUnlock() if ok { return true } app, err := this.module.model.getAppByName(name) if err != nil || app == nil || app.Id == 0 { return false } this.nameMu.Lock() if this.regCache == nil { this.regCache = map[string]bool{} } this.regCache[name] = true this.nameMu.Unlock() return true } func (this *statComp) 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.(*Console) this.options = opt.(*Options) // 统计日界固定时区(默认 Asia/Shanghai),须与各 analyze 端一致,决定 0 点切日口径。 setStatLoc(utils.LoadStatLocation(this.options.Stat.Timezone)) if err = ensureStatTable(); err != nil { this.module.Errorln("console.stat: 建 stats_global_day 表失败", err) } // 唯一写入自检:复合主键缺列会让 OnConflict 退化成插入、产生重复行。仅告警,不自动改表。 if ok, cols := verifyStatTablePK(); !ok { log.Warn("console.stat: stats_global_day 复合主键不完整,唯一写入无法保证,请检查表结构", log.Field{Key: "pk", Value: cols}) } return } func (this *statComp) Start() (err error) { if err = this.ModuleCompBase.Start(); err != nil { return } // 未配置 subject → 不启用消费(仅保留查询接口对历史数据可用)。 if this.options.Stat.Subject == "" { log.Warn("console.stat: 未配置 Stat.Subject,统计消费未启用") return } if natssys.JetStream() == nil { log.Warn("console.stat: NATS 未就绪,统计消费未启用(NATS 就绪后重启 console 即可)") return } if err = this.subscribe(); err != nil { // 订阅失败不阻断 console 启动(后台其余功能照常),记录告警。 log.Errorf("console.stat: 订阅统计管道失败,消费未启用: %v", err) err = nil return } log.Infof("console.stat: 统计快照消费已启动 subject=%s stream=%s durable=%s", this.options.Stat.Subject, this.options.Stat.Stream, this.options.Stat.Durable) return } func (this *statComp) Destroy() (err error) { if this.sub != nil { _ = this.sub.Drain() } return this.ModuleCompBase.Destroy() } // subscribe 幂等建好 JetStream 流并起一个持久(push)消费者订阅 subject。 func (this *statComp) subscribe() error { js := natssys.JetStream() cfg := this.options.Stat // 流不存在则创建,捕获约定 subject。 si, err := js.StreamInfo(cfg.Stream) if err != nil { if si, err = js.AddStream(&gonats.StreamConfig{ Name: cfg.Stream, Subjects: []string{cfg.Subject}, Storage: gonats.FileStorage, Retention: gonats.LimitsPolicy, }); err != nil { return err } log.Infof("console.stat: 已创建 JetStream 流 %s(subject=%s)", cfg.Stream, cfg.Subject) } // 打印流当前堆积消息数:诊断用——若长期为 0 说明 home 根本没把快照推到本 NATS(上游/地址问题)。 log.Infof("console.stat: 订阅前 JetStream 流 %s 当前消息数=%d subjects=%v", cfg.Stream, si.State.Msgs, si.Config.Subjects) sub, err := js.Subscribe(cfg.Subject, this.onMessage, gonats.Durable(cfg.Durable), gonats.ManualAck(), gonats.AckExplicit(), gonats.DeliverAll(), ) if err != nil { return err } this.sub = sub return nil } // onMessage 处理一条统计快照:校验通过后以绝对值覆盖 upsert 到 stats_global_day。 // - 非快照消息(管道里残留的历史逐事件等)或解析失败:记 error 流水并 ack 丢弃(下次快照自愈); // - 缺核心维度(app_id/region/stat_day):记 warn 流水并 ack 丢弃; // - 应用未在 console 注册:记 warn 流水并 ack 丢弃(允许其后补注册,下次快照自动入库); // - 落库失败:nak 等待重投。 func (this *statComp) onMessage(msg *gonats.Msg) { recvTs := time.Now() // 接收日志:每收到一条都打印,用于确认 console 是否真的从 NATS 消费到消息(subject + 字节数)。 log.Infof("console.stat: 收到 NATS 消息 subject=%s bytes=%d", msg.Subject, len(msg.Data)) var snap comm.StatSnapshot if err := json.Unmarshal(msg.Data, &snap); err != nil || snap.Row == nil { log.Warn("console.stat: 快照反序列化失败/非快照消息,丢弃", log.Field{Key: "kind", Value: snap.Kind}) this.recordSnapshot(recvTs, &snap, "error", false, "快照解析失败或非快照消息") _ = msg.Ack() return } row := snap.Row // 核心维度兜底校验:app_id(应用名)、region、stat_day 是复合主键的必备维度,缺失则拼不出有效行。 if row.AppId == "" || row.Region == "" || row.StatDay == 0 { log.Warn("console.stat: 丢弃缺失核心维度的快照", log.Field{Key: "app", Value: row.AppId}, log.Field{Key: "region", Value: row.Region}, log.Field{Key: "stat_day", Value: row.StatDay}) // registered 如实反映应用名注册情况(缺维度与是否注册是两回事),避免面板「未注册」徽标误标。 this.recordSnapshot(recvTs, &snap, "warn", this.isRegisteredApp(row.AppId), "丢弃:缺核心维度(app/region/day)") _ = msg.Ack() return } // 仅已注册应用落库;未注册记 warn、不落库(允许其后在 console 补注册,下次快照自动入库)。 if !this.isRegisteredApp(row.AppId) { log.Warnf("console.stat: 快照应用未在 console 注册,未落库 app=%q region=%s day=%d", row.AppId, row.Region, row.StatDay) this.recordSnapshot(recvTs, &snap, "warn", false, "应用未在 console 注册,未落库") _ = msg.Ack() return } if err := upsertStatDay(row); err != nil { log.Warn("console.stat: 快照落库失败,等待重投", log.Field{Key: "err", Value: err.Error()}) _ = msg.Nak() return } log.Infof("console.stat: 快照已落库 app=%s region=%s product=%d channel=%s day=%d | %s", row.AppId, row.Region, row.ProductId, row.ChannelId, row.StatDay, snapshotSummary(row)) this.recordSnapshot(recvTs, &snap, "info", true, "") _ = msg.Ack() }