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.
185 lines
7.6 KiB
185 lines
7.6 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)
|
|
}
|
|
// 唯一写入自检:复合主键缺列会让 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 day=%d | %s",
|
|
row.AppId, row.Region, row.ProductId, row.StatDay, snapshotSummary(row))
|
|
this.recordSnapshot(recvTs, &snap, "info", true, "")
|
|
_ = msg.Ack()
|
|
}
|
|
|