You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

190 lines
8.0 KiB

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)
}
// 唯一写入自检。⚠️ 主键缺 channel_id 不是「可能产生重复行」而是**快照一条都写不进去**:
// upsert 的 ON CONFLICT 指定五列,缺一列 Postgres 直接 42P10 报错,
// 表现为仪表盘恒为 0、只在日志里刷 WARN。修复见
// docs/migrations/2026-08-08-channel-dimension.sql 的 A4 段。
// 这里仍只告警不自动改表:有重复行时 ADD PRIMARY KEY 会失败,得先人工清理。
if ok, cols := verifyStatTablePK(); !ok {
log.Warn("console.stat: stats_global_day 复合主键缺列,统计快照将全部写入失败(42P10),"+
"请执行 docs/migrations/2026-08-08-channel-dimension.sql 的 A4 段",
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()
}