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.
269 lines
11 KiB
269 lines
11 KiB
package console
|
|
|
|
import (
|
|
"fmt"
|
|
"strings"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"yunyan/comm"
|
|
"yunyan/lego/sys/postgres"
|
|
"yunyan/pb"
|
|
"yunyan/utils"
|
|
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
// 统计落库与 dashboard 聚合查询的数据模型。console 为 stats_global_day 的唯一写库方:
|
|
// 消费各业务部署 analyze 推送的【绝对快照】(comm.StatSnapshot),直接以绝对值覆盖 upsert
|
|
// (见 stat.go.onMessage)。console 自身不再维护 Redis 累加/HLL——去重人数由业务侧算好随快照带来。
|
|
|
|
// statLoc 统计日界时区,由 statComp Init 经 setStatLoc 写入(默认 Asia/Shanghai)。
|
|
// 与各 analyze 端共用同一口径,保证 0 点切日一致;未初始化时回退固定 UTC+8。
|
|
var statLoc atomic.Value // *time.Location
|
|
|
|
func setStatLoc(loc *time.Location) { statLoc.Store(loc) }
|
|
|
|
func statLocation() *time.Location {
|
|
if v := statLoc.Load(); v != nil {
|
|
return v.(*time.Location)
|
|
}
|
|
return time.FixedZone("UTC+8", 8*3600)
|
|
}
|
|
|
|
// statDayNumber 把时间转成 YYYYMMDD 形式的统计日(按固定统计时区,与 analyze 口径一致)。
|
|
func statDayNumber(t time.Time) uint32 {
|
|
return utils.StatDayNumber(t, statLocation())
|
|
}
|
|
|
|
// ensureStatTable 在 console 公共库(AdminDB/supabase)建好 stats_global_day。
|
|
func ensureStatTable() error {
|
|
if err := postgres.CreateTable(comm.TableStatsGlobalDay, &pb.StatsGlobalDay{}); err != nil {
|
|
return err
|
|
}
|
|
// postgres CreateTable 对已存在表跳过 AutoMigrate,补可能漂移缺失的列(幂等)。
|
|
// qr_active_count 为后加字段,老表常缺,缺失会让所有 SUM 查询(statSumSelect)整体报错。
|
|
if res := postgres.Exec("ALTER TABLE " + comm.TableStatsGlobalDay + " ADD COLUMN IF NOT EXISTS qr_active_count bigint DEFAULT 0"); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
// channel_id 为后加的维度列。这里只补列(NOT NULL DEFAULT ''),【不动主键】——
|
|
// 把它并进复合主键需要 DROP/ADD CONSTRAINT,在有重复行时会失败,必须走
|
|
// docs/migrations 下的迁移脚本在维护窗口执行;未迁移主键前 verifyStatTablePK 会告警。
|
|
if res := postgres.Exec("ALTER TABLE " + comm.TableStatsGlobalDay + " ADD COLUMN IF NOT EXISTS channel_id varchar(16) NOT NULL DEFAULT ''"); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// verifyStatTablePK 只读自检 stats_global_day 的复合主键是否含五个维度列
|
|
// (app_id, product_id, channel_id, region, stat_day)。OnConflict 的「唯一写入」依赖此主键:
|
|
// 同一复合键的重复写会走 UPDATE 而非 INSERT,故主键缺一不可。
|
|
// 仅返回结果供调用方告警,不自动改主键(避免在已有重复行时 ALTER 失败)。
|
|
func verifyStatTablePK() (ok bool, cols string) {
|
|
postgres.Table(comm.TableStatsGlobalDay).Raw(
|
|
`SELECT COALESCE(string_agg(a.attname, ','), '')
|
|
FROM pg_index i
|
|
JOIN pg_attribute a ON a.attrelid = i.indrelid AND a.attnum = ANY(i.indkey)
|
|
WHERE i.indrelid = ?::regclass AND i.indisprimary`,
|
|
comm.TableStatsGlobalDay).Scan(&cols)
|
|
ok = strings.Contains(cols, "app_id") && strings.Contains(cols, "product_id") &&
|
|
strings.Contains(cols, "channel_id") && strings.Contains(cols, "region") &&
|
|
strings.Contains(cols, "stat_day")
|
|
return
|
|
}
|
|
|
|
// upsertStatDay 以「绝对值」覆盖写入一行(OnConflict UpdateAll)。Redis 侧持有的就是当日累计
|
|
// 绝对值,故覆盖式写入幂等、可重复执行。
|
|
func upsertStatDay(row *pb.StatsGlobalDay) error {
|
|
return postgres.Table(comm.TableStatsGlobalDay).
|
|
Clauses(clause.OnConflict{UpdateAll: true}).
|
|
Create(row).Error
|
|
}
|
|
|
|
// upsertStatDayBatch 批量覆盖写入多行(一条多行 INSERT ... ON CONFLICT DO UPDATE)。
|
|
// 历史分析一次产出成百上千行,逐行 upsert 会被远程库(Supabase) RTT 放大成 N×往返;
|
|
// 批量后压成 ⌈N/batchSize⌉ 次往返。语义与 upsertStatDay 一致(绝对值覆盖、幂等)。
|
|
func upsertStatDayBatch(rows []*pb.StatsGlobalDay) error {
|
|
if len(rows) == 0 {
|
|
return nil
|
|
}
|
|
return postgres.Table(comm.TableStatsGlobalDay).
|
|
Clauses(clause.OnConflict{UpdateAll: true}).
|
|
CreateInBatches(rows, 200).Error
|
|
}
|
|
|
|
// ============================ dashboard 聚合查询 ============================
|
|
//
|
|
// 直接查 console 公共库的 stats_global_day(console 消费各部署快照后按 app×product×region×day 落行),
|
|
// 按 stat_day 区间 SUM 聚合,支持按 app/region/product 下钻。
|
|
//
|
|
// 注意:*_user_count 三列是「当日去重」值,跨日/跨区域 SUM 得到的是「人日」之和而非
|
|
// 唯一人数(HLL 不可跨行精确合并);次数/金额/时长类 SUM 则为精确累计。
|
|
|
|
// statsQueryReq 聚合查询请求。指针字段为 nil 表示该维度不过滤(全量);
|
|
// 用指针而非零值,是因为 product_id=0、空区域都是合法取值,不能拿零值当“全部”。
|
|
// app_id 为应用名称、region 为区域代码(与 stats_global_day 存储一致)。
|
|
type statsQueryReq struct {
|
|
StartDay uint32 `json:"start_day"` // 起始统计日 YYYYMMDD(含),0 不限
|
|
EndDay uint32 `json:"end_day"` // 结束统计日 YYYYMMDD(含),0 不限
|
|
AppId *string `json:"app_id"` // 应用名称过滤(单值),nil 全部
|
|
Region *string `json:"region"` // 区域代码过滤(单值),nil 全部
|
|
ProductId *uint32 `json:"product_id"` // 产品过滤(单值),nil 全部
|
|
ChannelId *string `json:"channel_id"` // 渠道商过滤(单值,空串=只看未分配渠道的数据),nil 全部
|
|
// 多值 IN 过滤:代理"全部(我的)"用,传绑定列表,落 app_id IN (...) 等。与单值并存(同维度二选一即可)。
|
|
Apps []string `json:"apps"` // 应用名称列表
|
|
Regions []string `json:"regions"` // 区域代码列表
|
|
ProductIds []uint32 `json:"product_ids"` // 产品 id 列表
|
|
Channels []string `json:"channels"` // 渠道商 id 列表
|
|
}
|
|
|
|
// StatMetrics stats_global_day 的全部可累加指标(聚合结果,列名与表一致)。
|
|
type StatMetrics struct {
|
|
LoginUserCount int64 `json:"login_user_count"`
|
|
LoginCount int64 `json:"login_count"`
|
|
NewUserCount int64 `json:"new_user_count"`
|
|
ActiveUserCount int64 `json:"active_user_count"`
|
|
|
|
ActiveDeviceCount int64 `json:"active_device_count"`
|
|
BindDeviceCount int64 `json:"bind_device_count"`
|
|
|
|
QrActiveCount int64 `json:"qr_active_count"`
|
|
|
|
OrderCreateCount int64 `json:"order_create_count"`
|
|
OrderPaidCount int64 `json:"order_paid_count"`
|
|
OrderFailedCount int64 `json:"order_failed_count"`
|
|
OrderCreateAmount int64 `json:"order_create_amount"`
|
|
OrderPaidAmount int64 `json:"order_paid_amount"`
|
|
PayUserCount int64 `json:"pay_user_count"`
|
|
|
|
GrantVipDay int64 `json:"grant_vip_day"`
|
|
GrantAiIntegral int64 `json:"grant_ai_integral"`
|
|
GrantTradeSecond int64 `json:"grant_trade_second"`
|
|
GrantMeetSecond int64 `json:"grant_meet_second"`
|
|
|
|
AiChatCount int64 `json:"ai_chat_count"`
|
|
AiUpToken int64 `json:"ai_up_token"`
|
|
AiDownToken int64 `json:"ai_down_token"`
|
|
TradeCount int64 `json:"trade_count"`
|
|
TradeSecond int64 `json:"trade_second"`
|
|
TradeWords int64 `json:"trade_words"`
|
|
MeetCount int64 `json:"meet_count"`
|
|
MeetSecond int64 `json:"meet_second"`
|
|
}
|
|
|
|
// StatTrendRow 按日聚合的一行(趋势图用)。
|
|
type StatTrendRow struct {
|
|
StatDay uint32 `json:"stat_day"`
|
|
StatMetrics
|
|
}
|
|
|
|
// statMetricCols 与 StatMetrics 字段一一对应的列名(SUM 用)。
|
|
var statMetricCols = []string{
|
|
"login_user_count", "login_count", "new_user_count", "active_user_count",
|
|
"active_device_count", "bind_device_count",
|
|
"qr_active_count",
|
|
"order_create_count", "order_paid_count", "order_failed_count", "order_create_amount", "order_paid_amount", "pay_user_count",
|
|
"grant_vip_day", "grant_ai_integral", "grant_trade_second", "grant_meet_second",
|
|
"ai_chat_count", "ai_up_token", "ai_down_token", "trade_count", "trade_second", "trade_words", "meet_count", "meet_second",
|
|
}
|
|
|
|
// statSumSelect 拼出 "extra..., COALESCE(SUM(col),0) AS col, ..." 的 select 子句。
|
|
// COALESCE 兜底空结果集的 SUM(NULL),避免扫描 NULL 进 int64 报错。
|
|
func statSumSelect(extra ...string) string {
|
|
parts := make([]string, 0, len(extra)+len(statMetricCols))
|
|
parts = append(parts, extra...)
|
|
for _, c := range statMetricCols {
|
|
parts = append(parts, fmt.Sprintf("COALESCE(SUM(%s),0) AS %s", c, c))
|
|
}
|
|
return strings.Join(parts, ", ")
|
|
}
|
|
|
|
// statWhere 按请求拼过滤条件。
|
|
func statWhere(req *statsQueryReq) (string, []interface{}) {
|
|
conds := make([]string, 0, 5)
|
|
args := make([]interface{}, 0, 5)
|
|
if req.StartDay > 0 {
|
|
conds = append(conds, "stat_day >= ?")
|
|
args = append(args, req.StartDay)
|
|
}
|
|
if req.EndDay > 0 {
|
|
conds = append(conds, "stat_day <= ?")
|
|
args = append(args, req.EndDay)
|
|
}
|
|
if req.AppId != nil {
|
|
conds = append(conds, "app_id = ?")
|
|
args = append(args, *req.AppId)
|
|
}
|
|
if req.Region != nil {
|
|
conds = append(conds, "region = ?")
|
|
args = append(args, *req.Region)
|
|
}
|
|
if req.ProductId != nil {
|
|
conds = append(conds, "product_id = ?")
|
|
args = append(args, *req.ProductId)
|
|
}
|
|
if req.ChannelId != nil {
|
|
conds = append(conds, "channel_id = ?")
|
|
args = append(args, *req.ChannelId)
|
|
}
|
|
// 多值 IN 过滤(代理"全部(我的)"):gorm 对 "col IN ?" + 切片会展开为 IN (v1,v2,...)。
|
|
if len(req.Apps) > 0 {
|
|
conds = append(conds, "app_id IN ?")
|
|
args = append(args, req.Apps)
|
|
}
|
|
if len(req.Regions) > 0 {
|
|
conds = append(conds, "region IN ?")
|
|
args = append(args, req.Regions)
|
|
}
|
|
if len(req.ProductIds) > 0 {
|
|
conds = append(conds, "product_id IN ?")
|
|
args = append(args, req.ProductIds)
|
|
}
|
|
if len(req.Channels) > 0 {
|
|
conds = append(conds, "channel_id IN ?")
|
|
args = append(args, req.Channels)
|
|
}
|
|
return strings.Join(conds, " AND "), args
|
|
}
|
|
|
|
// andWhere 把两个 WHERE 片段用 AND 合并(任一为空则取另一个),并按顺序拼接参数。
|
|
func andWhere(w1 string, a1 []interface{}, w2 string, a2 []interface{}) (string, []interface{}) {
|
|
switch {
|
|
case w1 == "" && w2 == "":
|
|
return "", nil
|
|
case w1 == "":
|
|
return w2, a2
|
|
case w2 == "":
|
|
return w1, a1
|
|
default:
|
|
return "(" + w1 + ") AND (" + w2 + ")", append(append([]interface{}{}, a1...), a2...)
|
|
}
|
|
}
|
|
|
|
// queryStatSummary 区间总量(全维度 SUM 成一行)。scopeWhere/scopeArgs 为账号作用域附加条件(与下钻 AND)。
|
|
func queryStatSummary(req *statsQueryReq, scopeWhere string, scopeArgs []interface{}) (*StatMetrics, error) {
|
|
out := &StatMetrics{}
|
|
where, args := statWhere(req)
|
|
where, args = andWhere(where, args, scopeWhere, scopeArgs)
|
|
tx := postgres.Table(comm.TableStatsGlobalDay).Select(statSumSelect())
|
|
if where != "" {
|
|
tx = tx.Where(where, args...)
|
|
}
|
|
err := tx.Scan(out).Error
|
|
return out, err
|
|
}
|
|
|
|
// queryStatTrend 按 stat_day 分组的逐日序列(升序)。
|
|
func queryStatTrend(req *statsQueryReq, scopeWhere string, scopeArgs []interface{}) ([]*StatTrendRow, error) {
|
|
rows := make([]*StatTrendRow, 0)
|
|
where, args := statWhere(req)
|
|
where, args = andWhere(where, args, scopeWhere, scopeArgs)
|
|
tx := postgres.Table(comm.TableStatsGlobalDay).
|
|
Select(statSumSelect("stat_day")).
|
|
Group("stat_day").Order("stat_day ASC")
|
|
if where != "" {
|
|
tx = tx.Where(where, args...)
|
|
}
|
|
err := tx.Scan(&rows).Error
|
|
return rows, err
|
|
}
|
|
|