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.
242 lines
9.1 KiB
242 lines
9.1 KiB
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
|
|
}
|
|
|