// Package analyze 是平台统计模块。 // // 采集走「Redis 实时累加」,进程崩溃不丢数据: // - 业务模块在启动阶段经 service.GetModule(comm.ModuleAnalyze) 取到本模块、 // 断言为 comm.IAnalyze,调用 Report 投递埋点事件;事件被立即翻译成 // HINCRBY / PFADD 原子累加到 Redis 当日统计上——Redis 始终持有当日实时 // 统计,不存事件队列; // - 同步器每隔 SyncInterval 秒把 Redis 统计转存到 stats_global_day 表, // 跨天后顺手清理上一日数据,开始新一天的统计。 // // 埋点接口定义在 comm.IAnalyze,业务模块只依赖 comm,无需 import 本包。 package analyze import ( "context" "fmt" "time" "yunyan/comm" "yunyan/lego/sys/log" redissys "yunyan/lego/sys/redis" "github.com/redis/go-redis/v9" ) // partitionTTL 分区键的存活时间,每次 Report 刷新。保证活跃分区不过期; // 兜底回收同步器未能清理的历史分区,防止 Redis 无限增长。 const partitionTTL = 72 * time.Hour // setupKeys 确定本部署的 Redis 键前缀 {集群标签}:analyze:{应用名}:{区域}, // 由模块 Init 调用一次。应用名不含冒号(配置可控的应用标识)。 func setupKeys(clusterTag string, appName string, region int32) { keyPrefix.Store(fmt.Sprintf("%s:analyze:%s:%d", clusterTag, appName, region)) } // Report 实现 comm.IAnalyze:投递一个埋点事件,直接在 Redis 上做原子累加。 // // 事件不入队列——而是 HINCRBY 累加进当日统计哈希、PFADD 进去重 HLL,一次 // pipeline 完成。调用返回即代表统计已落到 Redis,进程崩溃不丢。Redis 不可用时 // 事件被丢弃并告警(降级),不阻断业务主流程。 func (this *Analyze) Report(e *comm.StatEvent) { if e == nil { return } if prefix() == "" { log.Warn("analyze: 键前缀未初始化,埋点事件丢弃") return } ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() statDay := dayNumber(time.Now()) if _, err := redissys.Conn().Pipelined(ctx, func(p redis.Pipeliner) error { writeEvent(ctx, p, e, statDay) return nil }); err != nil { log.Warn("analyze: 写入 Redis 失败,埋点事件丢弃", log.Field{Key: "err", Value: err.Error()}) } // 注意:不在埋点热路径上转发。汇报给 console 由同步器(syncComp)按分区【绝对快照】定时推送, // 见 sync.go——绝对值幂等、丢消息自愈、重启自动对账,避免逐事件转发的丢失与重复计数问题。 } // writeEvent 把一个事件翻译成 Redis 累加命令,写入 pipeline。 func writeEvent(ctx context.Context, p redis.Pipeliner, e *comm.StatEvent, statDay uint32) { // 分区 = 产品 × 渠道商 × 日:渠道商是设备生产时定下的归属,随事件带上来, // 这样同一产品下不同渠道商的数据天然分行,后台可按渠道商下钻与结算。 part := partition(e.ProductId, e.ChannelId, statDay) hkey := hashKey(part) cnt := e.Count if cnt == 0 { cnt = 1 // 0 视为未指定,按 1 计;负数保留(用于冲正) } // 维护分区索引并刷新分区 TTL。 p.SAdd(ctx, indexKey(), part) p.Expire(ctx, indexKey(), partitionTTL) p.Expire(ctx, hkey, partitionTTL) hincr := func(field string, n int64) { if n != 0 { p.HIncrBy(ctx, hkey, field, n) } } pfadd := func(metric string) { if e.Uid != "" { k := hllKey(part, metric) p.PFAdd(ctx, k, e.Uid) p.Expire(ctx, k, partitionTTL) } } pfadd(uActive) // 任意带 uid 的事件都计入活跃用户 switch e.Type { case comm.StatEventLogin: hincr(fLoginCount, cnt) pfadd(uLogin) case comm.StatEventRegister: hincr(fNewUserCount, cnt) case comm.StatEventDeviceActive: hincr(fActiveDeviceCount, cnt) case comm.StatEventDeviceBind: hincr(fBindDeviceCount, cnt) case comm.StatEventQRCodeActive: hincr(fQRActiveCount, cnt) case comm.StatEventOrderCreate: hincr(fOrderCreateCount, cnt) hincr(fOrderCreateAmount, e.Amount) case comm.StatEventOrderPaid: hincr(fOrderPaidCount, cnt) hincr(fOrderPaidAmount, e.Amount) pfadd(uPay) case comm.StatEventOrderFailed: hincr(fOrderFailedCount, cnt) case comm.StatEventTranslate: hincr(fTradeCount, cnt) hincr(fTradeSecond, e.Second) hincr(fTradeWords, e.Words) case comm.StatEventMeeting: hincr(fMeetCount, cnt) hincr(fMeetSecond, e.Second) case comm.StatEventAiChat: hincr(fAiChatCount, cnt) hincr(fAiUpToken, e.UpToken) hincr(fAiDownToken, e.DnToken) case comm.StatEventResourceGrant: hincr(fGrantVipDay, e.GrantVipDay) hincr(fGrantAiIntegral, e.GrantAiIntegral) hincr(fGrantTradeSecond, e.GrantTradeSecond) hincr(fGrantMeetSecond, e.GrantMeetSecond) default: log.Warn("analyze: 未知埋点事件类型", log.Field{Key: "type", Value: int32(e.Type)}) } }