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.
 
 
 
 
 
 

255 lines
9.7 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, channelId, statDay, ok := parsePartition(m)
if !ok {
this.module.Warn("analyze: 非法分区标识", log.Field{Key: "member", Value: m})
continue
}
row, ok := this.readRow(ctx, m, productId, channelId, statDay, ts)
if !ok {
// 分区已无累加数据(被清理过的索引残留),历史日则顺手摘除索引。
if statDay < today {
redissys.Conn().SRem(ctx, indexKey(), m)
}
continue
}
this.publishSnapshot(row)
if statDay < yesterday {
this.cleanup(ctx, m)
}
}
}
// 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 channel=%s day=%d",
this.options.Subject, row.AppId, row.Region, row.ProductId, row.ChannelId, row.StatDay)
}
// readRow 从 Redis 读出某分区的统计并组装成 DB 行;分区无累加数据时返回 ok=false。
// part 是索引里的分区标识原文(老格式无渠道商段也照常工作,见 core.go 的 hashKey 说明)。
func (this *syncComp) readRow(ctx context.Context, part string, productId uint32, channelId string, statDay uint32, ts int64) (*pb.StatsGlobalDay, bool) {
h, err := redissys.Conn().HGetAll(ctx, hashKey(part)).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,
ChannelId: channelId,
Region: comm.RegionCode(this.options.Region),
StatDay: statDay,
UpdateTime: ts,
LoginUserCount: this.pfcount(ctx, part, uLogin),
LoginCount: hashInt(h, fLoginCount),
NewUserCount: hashInt(h, fNewUserCount),
ActiveUserCount: this.pfcount(ctx, part, 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, part, 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, part, metric string) int64 {
n, err := redissys.Conn().PFCount(ctx, hllKey(part, metric)).Result()
if err != nil {
this.module.Warn("analyze: PFCOUNT 失败", log.Field{Key: "err", Value: err.Error()})
return 0
}
return n
}
// cleanup 删除某历史分区在 Redis 中的全部数据。part 为索引里的分区标识原文。
func (this *syncComp) cleanup(ctx context.Context, part string) {
redissys.Conn().Del(ctx,
hashKey(part),
hllKey(part, uLogin),
hllKey(part, uActive),
hllKey(part, uPay))
redissys.Conn().SRem(ctx, indexKey(), part)
}
// hashInt 从 HGETALL 结果里取一个整型字段,缺失或非法按 0 计。
func hashInt(h map[string]string, field string) int64 {
n, _ := strconv.ParseInt(h[field], 10, 64)
return n
}
// parsePartition 解析索引成员。新格式 "产品ID:渠道商:统计日",
// 老格式 "产品ID:统计日" 同样接受(渠道商按空处理)——版本切换当天 Redis 里
// 还留着老格式分区,兼容解析才能把它们照常推送并清理掉。
func parsePartition(m string) (productId uint32, channelId string, statDay uint32, ok bool) {
parts := strings.Split(m, ":")
var pidStr, dayStr string
switch len(parts) {
case 2: // 老格式:产品:日
pidStr, dayStr = parts[0], parts[1]
case 3: // 新格式:产品:渠道商:日
pidStr, channelId, dayStr = parts[0], parts[1], parts[2]
if channelId == emptyChannelMark {
channelId = ""
}
default:
return 0, "", 0, false
}
p, e1 := strconv.ParseUint(pidStr, 10, 32)
d, e2 := strconv.ParseUint(dayStr, 10, 32)
if e1 != nil || e2 != nil {
return 0, "", 0, false
}
return uint32(p), channelId, uint32(d), true
}