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.
 
 
 
 
 
 

132 lines
4.8 KiB

// 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)})
}
}