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.
 
 
 
 
 
 

445 lines
20 KiB

package console
import (
"encoding/json"
"fmt"
"strconv"
"strings"
"time"
"yunyan/comm"
"yunyan/lego/sys/mysql"
"yunyan/lego/sys/postgres"
"yunyan/pb"
"github.com/gin-gonic/gin"
)
// 历史数据分析:扫描某注册应用的 MySQL 业务库原始表,按天聚合成 stats_global_day 行,
// 补齐实时管道(analyze 双写 NATS)上线前的历史空白。
//
// 仅能还原"带时间戳的"日维度指标:
// - 新增用户(user.createtime)
// - 订单:创建/支付/失败数、创建/支付金额、付费人数(payorder.create_time / pay_time),按最后绑定产品归属
// - 资源发放:VIP天/AI积分/翻译时长/会议时长(useruselog.ts,仅发放类日志),按最后绑定产品归属
// - 激活/绑定设备:按产品拆(userdevice.productid),JOIN user 取注册时间近似为激活时间
// - 资源使用:翻译/会议的使用次数与时长(useruselog logtype=UserConsume,add* 负值为消耗),按最后绑定产品归属
//
// 无法还原(保持 0,由实时管道往后采集):登录、活跃用户(当日去重 HLL)、AI对话 token、翻译字数。
//
// 落库维度:app_id(注册应用)+ region(应用注册区域)+ product_id + stat_day。
// product_id 归属:设备类指标按 userdevice.productid;订单/发放/使用按用户「最后绑定产品」归属——
// 直接从 userdevice 取该用户最新绑定设备(最大 id)的 productid 推导,不依赖 user.lastbindproductid 列
//(该列由实时埋点往后写,存量数据/未迁移的业务库里可能尚不存在);新增用户及无时间戳指标无产品维度,用 product_id=0。
// 无设备记录的用户其订单/使用落 product_id=0。分析前会先清理 (app_id, region[, 日期区间]) 的旧行,避免重叠。
// regionLabelMap 应用注册的区域中文标签 → pb.Region(与 apps.html 下拉、pb.Region 枚举一致)。
var regionLabelMap = map[string]pb.Region{
"中国": pb.Region_RegionChina, "美国": pb.Region_RegionUSA, "东亚": pb.Region_RegionEastAsia,
"东南亚": pb.Region_RegionSoutheastAsia, "南亚": pb.Region_RegionSouthAsia, "中东": pb.Region_RegionMiddleEast,
"欧洲": pb.Region_RegionEurope, "俄罗斯": pb.Region_RegionRussia, "非洲": pb.Region_RegionAfrica,
"南美": pb.Region_RegionSouthAmerica, "北美": pb.Region_RegionNorthAmerica, "大洋洲": pb.Region_RegionOceania,
"巴西": pb.Region_RegionBrazil, "印度": pb.Region_RegionIndia,
}
func regionFromLabel(label string) pb.Region {
if r, ok := regionLabelMap[strings.TrimSpace(label)]; ok {
return r
}
return pb.Region_RegionUnknown
}
// regionCodeFromLabel 应用注册的中文区域标签 → 英文区域代码(cn/us/...),与 comm.RegionCode 同源。
func regionCodeFromLabel(label string) string {
return comm.RegionCode(int32(regionFromLabel(label)))
}
// analyzeStep 分析报告里的单步结果。
type analyzeStep struct {
Name string `json:"name"` // 步骤名
Groups int64 `json:"groups"` // 命中的分组数(按天/按产品聚合后的行数)
}
// analyzeReport 分析报告,回前端展示过程与结果。
type analyzeReport struct {
AppName string `json:"app_name"`
Region string `json:"region"`
StartDay uint32 `json:"start_day"`
EndDay uint32 `json:"end_day"`
ClearedRows int64 `json:"cleared_rows"` // 分析前清理的旧行数
Steps []analyzeStep `json:"steps"` // 各聚合步骤命中情况
DaysWritten int `json:"days_written"` // 覆盖天数
RowsWritten int `json:"rows_written"` // 写入行数(app×region×product×day)
}
func (this *serverComp) appsAnalyze(c *gin.Context) {
var req struct {
AppId uint32 `json:"app_id"`
StartDay uint32 `json:"start_day"` // YYYYMMDD,0 不限
EndDay uint32 `json:"end_day"` // YYYYMMDD,0 不限
}
_ = c.ShouldBindJSON(&req)
// NDJSON 流式输出:每行一个 JSON 事件,逐步 flush,前端边读边显示进度条与分步状态。
// 事件类型:{"type":"progress",index,total,name} / {"type":"done",report} / {"type":"error",msg}
c.Header("Content-Type", "application/x-ndjson; charset=utf-8")
c.Header("Cache-Control", "no-cache")
c.Header("X-Accel-Buffering", "no") // 关掉反代缓冲,保证逐行下发
enc := json.NewEncoder(c.Writer)
emit := func(v interface{}) {
_ = enc.Encode(v) // Encode 自带换行 → NDJSON
c.Writer.Flush()
}
if req.AppId == 0 {
emit(gin.H{"type": "error", "msg": "app_id 必填"})
return
}
// 并发保护:同一应用同一时刻只允许一个分析在跑(含 DELETE+重写,并发会互相清表)。
// 重复点击 / 多会话同时触发 → 直接拒绝,不排队。完成(含出错 return)后释放。
if _, busy := this.analyzing.LoadOrStore(req.AppId, true); busy {
emit(gin.H{"type": "error", "msg": "该应用正在分析中,请等当前分析完成后再试"})
return
}
defer this.analyzing.Delete(req.AppId)
app, err := this.module.model.getApp(req.AppId)
if err != nil {
emit(gin.H{"type": "error", "msg": "应用未注册: " + err.Error()})
return
}
svc, err := this.module.registry.getServiceDB(req.AppId)
if err != nil {
emit(gin.H{"type": "error", "msg": "连接应用业务库失败: " + err.Error()})
return
}
rep, err := analyzeAppHistory(svc, app.Name, regionCodeFromLabel(app.Region), req.StartDay, req.EndDay,
func(idx, total int, name string) {
emit(gin.H{"type": "progress", "index": idx, "total": total, "name": name})
})
if err != nil {
emit(gin.H{"type": "error", "msg": err.Error()})
return
}
rep.AppName = app.Name
rep.Region = app.Region // 展示用中文标签;落库 region 列存的是英文代码
emit(gin.H{"type": "done", "report": rep})
}
// dayStrToNum "2026-06-12" -> 20260612
func dayStrToNum(s string) uint32 {
n, _ := strconv.ParseUint(strings.ReplaceAll(s, "-", ""), 10, 32)
return uint32(n)
}
// analyzeStepTotal 进度总步数(与下方 progress 调用一一对应,供前端进度条计算百分比)。
const analyzeStepTotal = 9
// analyzeAppHistory 跑各项按天聚合,合并成每日行,清理旧行后落库 stats_global_day,返回分析报告。
// appName 落 stats_global_day.app_id(应用名称),regionCode 落 region(英文代码)。
// onProgress 非 nil 时,在每个阶段开始前回调 (当前步, 总步数, 步骤名),供前端实时展示进度;可传 nil。
func analyzeAppHistory(svc mysql.ISys, appName, regionCode string, startDay, endDay uint32, onProgress func(idx, total int, name string)) (*analyzeReport, error) {
now := time.Now().Unix()
rep := &analyzeReport{StartDay: startDay, EndDay: endDay}
// progress 在每个阶段开始前推送一次进度(步号自增),onProgress 为 nil 时静默。
stepNo := 0
progress := func(name string) {
stepNo++
if onProgress != nil {
onProgress(stepNo, analyzeStepTotal, name)
}
}
// 行按 (product, day) 归集:设备类指标按真实 product 落,其余指标用 product=0。
type pk struct{ product, day uint32 }
rowMap := map[pk]*pb.StatsGlobalDay{}
ensure := func(product, day uint32) *pb.StatsGlobalDay {
k := pk{product, day}
r := rowMap[k]
if r == nil {
r = &pb.StatsGlobalDay{AppId: appName, ProductId: product, Region: regionCode, StatDay: day, UpdateTime: now}
rowMap[k] = r
}
return r
}
addStep := func(name string, groups int) {
rep.Steps = append(rep.Steps, analyzeStep{Name: name, Groups: int64(groups)})
}
type cntRow struct {
D string `gorm:"column:d"`
Cnt int64 `gorm:"column:cnt"`
}
// 「最后绑定产品」子查询:从 userdevice 取每个用户最新绑定设备(最大 id)的 productid,
// 据此把订单/发放/使用归属到产品。直接读设备列表推导,不依赖 user.lastbindproductid 列
//(该列由实时埋点往后写,存量数据/未迁移的业务库里可能尚不存在)。LEFT JOIN,无设备的用户落 0。
lastBindSub := "(SELECT ud.uid AS uid, ud.productid AS pid FROM " + comm.TableUserdevice + " ud " +
"JOIN (SELECT uid, MAX(id) AS mid FROM " + comm.TableUserdevice + " WHERE uid <> '' GROUP BY uid) m " +
"ON ud.uid = m.uid AND ud.id = m.mid)"
// 1) 创建订单:按 create_time 切日,count + sum(amount),按最后绑定产品归属
progress("创建订单")
type ocRow struct {
Pid uint32 `gorm:"column:pid"`
D string `gorm:"column:d"`
Cnt int64 `gorm:"column:cnt"`
Amt int64 `gorm:"column:amt"`
}
var oc []*ocRow
if err := svc.Table(comm.TablePayOrder + " AS po").
Select("COALESCE(lb.pid,0) AS pid, DATE_FORMAT(FROM_UNIXTIME(po.create_time), '%Y-%m-%d') AS d, COUNT(*) AS cnt, COALESCE(SUM(po.amount),0) AS amt").
Joins("LEFT JOIN " + lastBindSub + " AS lb ON po.uid = lb.uid").
Where("po.create_time > 0").Group("pid, d").Scan(&oc).Error; err != nil {
return nil, fmt.Errorf("聚合创建订单失败: %w", err)
}
for _, r := range oc {
row := ensure(r.Pid, dayStrToNum(r.D))
row.OrderCreateCount = r.Cnt
row.OrderCreateAmount = r.Amt
}
addStep("创建订单", len(oc))
// 2) 失败订单:按 create_time 切日,status=FAILED,LEFT JOIN user 按最后绑定产品归属
progress("失败订单")
type pidCntRow struct {
Pid uint32 `gorm:"column:pid"`
D string `gorm:"column:d"`
Cnt int64 `gorm:"column:cnt"`
}
var of []*pidCntRow
if err := svc.Table(comm.TablePayOrder+" AS po").
Select("COALESCE(lb.pid,0) AS pid, DATE_FORMAT(FROM_UNIXTIME(po.create_time), '%Y-%m-%d') AS d, COUNT(*) AS cnt").
Joins("LEFT JOIN "+lastBindSub+" AS lb ON po.uid = lb.uid").
Where("po.create_time > 0 AND po.status = ?", int32(pb.PayOrderStatus_PAY_ORDER_FAILED)).
Group("pid, d").Scan(&of).Error; err != nil {
return nil, fmt.Errorf("聚合失败订单失败: %w", err)
}
for _, r := range of {
ensure(r.Pid, dayStrToNum(r.D)).OrderFailedCount = r.Cnt
}
addStep("失败订单", len(of))
// 3) 支付订单:按 pay_time 切日,status=PAID,count + sum(amount) + distinct uid,LEFT JOIN user 按最后绑定产品归属
progress("支付订单")
type paidRow struct {
Pid uint32 `gorm:"column:pid"`
D string `gorm:"column:d"`
Cnt int64 `gorm:"column:cnt"`
Amt int64 `gorm:"column:amt"`
UCnt int64 `gorm:"column:ucnt"`
}
// 注意:微信/支付宝/IAP 的 notify 回调只置 status=PAID 未写 pay_time,故不能用 pay_time>0 过滤
//(会漏掉绝大多数已支付订单)。改为只按 status=PAID 计,切日优先用 pay_time,缺失则回退 create_time。
var op []*paidRow
if err := svc.Table(comm.TablePayOrder+" AS po").
Select("COALESCE(lb.pid,0) AS pid, DATE_FORMAT(FROM_UNIXTIME(IF(po.pay_time > 0, po.pay_time, po.create_time)), '%Y-%m-%d') AS d, COUNT(*) AS cnt, COALESCE(SUM(po.amount),0) AS amt, COUNT(DISTINCT po.uid) AS ucnt").
Joins("LEFT JOIN "+lastBindSub+" AS lb ON po.uid = lb.uid").
Where("po.status = ? AND (po.pay_time > 0 OR po.create_time > 0)", int32(pb.PayOrderStatus_PAY_ORDER_PAID)).
Group("pid, d").Scan(&op).Error; err != nil {
return nil, fmt.Errorf("聚合支付订单失败: %w", err)
}
for _, r := range op {
row := ensure(r.Pid, dayStrToNum(r.D))
row.OrderPaidCount = r.Cnt
row.OrderPaidAmount = r.Amt
row.PayUserCount = r.UCnt
}
addStep("支付订单", len(op))
// 4) 新增用户:按 createtime 切日
progress("新增用户")
var su []*cntRow
if err := svc.Table(comm.TableUser).
Select("DATE_FORMAT(FROM_UNIXTIME(createtime), '%Y-%m-%d') AS d, COUNT(*) AS cnt").
Where("createtime > 0").Group("d").Scan(&su).Error; err != nil {
return nil, fmt.Errorf("聚合新增用户失败: %w", err)
}
for _, r := range su {
ensure(0, dayStrToNum(r.D)).NewUserCount = r.Cnt
}
addStep("新增用户", len(su))
// 5) 资源发放:按 ts 切日,SUM 四项。仅发放类日志(排除 UserConsume),避免负的消耗冲抵发放量。LEFT JOIN user 按最后绑定产品归属。
progress("资源发放")
type grantRow struct {
Pid uint32 `gorm:"column:pid"`
D string `gorm:"column:d"`
Vipday int64 `gorm:"column:vipday"`
Aichat int64 `gorm:"column:aichat"`
Tradesec int64 `gorm:"column:tradesec"`
Meetsec int64 `gorm:"column:meetsec"`
}
var gr []*grantRow
if err := svc.Table(comm.TableUserUseLog+" AS l").
Select("COALESCE(lb.pid,0) AS pid, DATE_FORMAT(FROM_UNIXTIME(l.ts), '%Y-%m-%d') AS d, "+
"COALESCE(SUM(l.addvipday),0) AS vipday, COALESCE(SUM(l.addagentintegral),0) AS aichat, "+
"COALESCE(SUM(l.addtradesecond),0) AS tradesec, COALESCE(SUM(l.addmeetsecond),0) AS meetsec").
Joins("LEFT JOIN "+lastBindSub+" AS lb ON l.uid = lb.uid").
Where("l.ts > 0 AND l.logtype <> ?", int32(pb.UserLogType_UserConsume)).Group("pid, d").Scan(&gr).Error; err != nil {
return nil, fmt.Errorf("聚合资源发放失败: %w", err)
}
for _, r := range gr {
row := ensure(r.Pid, dayStrToNum(r.D))
row.GrantVipDay = r.Vipday
row.GrantAiIntegral = r.Aichat
row.GrantTradeSecond = r.Tradesec
row.GrantMeetSecond = r.Meetsec
}
addStep("资源发放", len(gr))
// 6) 激活/绑定设备:从 console 主库 license_<产品> 码表统计(uid非空=绑定,status>0=激活)。
// userdevice.productid 是客户端声称值(校验失败还会回退到 45058),不可靠;license 码表按产品分表、
// 是设备码真实产品归属。license 表无应用/区域维度,故设备数落全局行(app_id='' region=''),切日用 usedtime。
// 本步与所分析的具体 app 无关,每次分析全量刷新一次全局设备数(见末尾落库)。
progress("激活/绑定设备")
deviceRows, derr := analyzeDeviceFromLicense(now)
if derr != nil {
return nil, fmt.Errorf("聚合激活设备失败: %w", derr)
}
addStep("激活/绑定设备(按产品·license)", len(deviceRows))
// 7) 资源使用(消耗):logtype=UserConsume,add* 负值为消耗。按 ts 切日。LEFT JOIN user 按最后绑定产品归属。
// addagentintegral<0 → 翻译/AI 使用;addmeetsecond<0 → 会议使用。
progress("资源使用")
type useRow struct {
Pid uint32 `gorm:"column:pid"`
D string `gorm:"column:d"`
TradeCnt int64 `gorm:"column:tcnt"`
TradeSec int64 `gorm:"column:tsec"`
MeetCnt int64 `gorm:"column:mcnt"`
MeetSec int64 `gorm:"column:msec"`
}
var us []*useRow
if err := svc.Table(comm.TableUserUseLog+" AS l").
Select("COALESCE(lb.pid,0) AS pid, DATE_FORMAT(FROM_UNIXTIME(l.ts), '%Y-%m-%d') AS d, "+
"SUM(CASE WHEN l.addagentintegral < 0 THEN 1 ELSE 0 END) AS tcnt, "+
"COALESCE(SUM(CASE WHEN l.addagentintegral < 0 THEN -l.addagentintegral ELSE 0 END),0) AS tsec, "+
"SUM(CASE WHEN l.addmeetsecond < 0 THEN 1 ELSE 0 END) AS mcnt, "+
"COALESCE(SUM(CASE WHEN l.addmeetsecond < 0 THEN -l.addmeetsecond ELSE 0 END),0) AS msec").
Joins("LEFT JOIN "+lastBindSub+" AS lb ON l.uid = lb.uid").
Where("l.ts > 0 AND l.logtype = ?", int32(pb.UserLogType_UserConsume)).
Group("pid, d").Scan(&us).Error; err != nil {
return nil, fmt.Errorf("聚合资源使用失败: %w", err)
}
for _, r := range us {
row := ensure(r.Pid, dayStrToNum(r.D))
row.TradeCount = r.TradeCnt
row.TradeSecond = r.TradeSec
row.MeetCount = r.MeetCnt
row.MeetSecond = r.MeetSec
}
addStep("资源使用(翻译/会议)", len(us))
// 分析前清理:删除该 (app, region[, 区间]) 的旧统计行,避免与上次分析/残留数据重叠。
progress("清理旧统计行")
clearSQL := "DELETE FROM " + comm.TableStatsGlobalDay + " WHERE app_id = ? AND region = ?"
clearArgs := []interface{}{appName, regionCode}
if startDay > 0 {
clearSQL += " AND stat_day >= ?"
clearArgs = append(clearArgs, startDay)
}
if endDay > 0 {
clearSQL += " AND stat_day <= ?"
clearArgs = append(clearArgs, endDay)
}
if res := postgres.Exec(clearSQL, clearArgs...); res.Error != nil {
return nil, fmt.Errorf("清理旧统计失败: %w", res.Error)
} else {
rep.ClearedRows = res.RowsAffected
}
// 落库:仅写入区间内的天(start/end 为 0 表示不限)。
progress("写入统计库")
daysSet := map[uint32]bool{}
batch := make([]*pb.StatsGlobalDay, 0, len(rowMap))
for k, row := range rowMap {
if k.day == 0 {
continue
}
if startDay > 0 && k.day < startDay {
continue
}
if endDay > 0 && k.day > endDay {
continue
}
batch = append(batch, row)
daysSet[k.day] = true
}
// 批量 upsert:逐行写远程库会被 RTT 放大(实测单行 ~700ms),批量后整体压成几次往返。
// 分块写 + 每块回调一次心跳进度:既保持流式连接有数据流动(避免浏览器掐掉空闲长连接),
// 又能显示写入进度(N/总数)。块大小兼顾单条 INSERT 效率与心跳频率。
const writeChunk = 200
total := len(batch)
for i := 0; i < total; i += writeChunk {
end := i + writeChunk
if end > total {
end = total
}
if err := upsertStatDayBatch(batch[i:end]); err != nil {
return rep, fmt.Errorf("批量写入 stats_global_day(%d 行)失败: %w", total, err)
}
if onProgress != nil {
onProgress(analyzeStepTotal, analyzeStepTotal, fmt.Sprintf("写入统计库 (%d/%d)", end, total))
}
}
// 设备行(全局 app_id=''):与 per-app 行解耦,单独全量刷新——清掉旧的全局设备行再批量写入。
// 每次分析都重算(license 表是全局的),幂等。
if res := postgres.Exec("DELETE FROM " + comm.TableStatsGlobalDay + " WHERE app_id = ''"); res.Error != nil {
return rep, fmt.Errorf("清理全局设备行失败: %w", res.Error)
}
for i := 0; i < len(deviceRows); i += writeChunk {
end := i + writeChunk
if end > len(deviceRows) {
end = len(deviceRows)
}
if err := upsertStatDayBatch(deviceRows[i:end]); err != nil {
return rep, fmt.Errorf("写入全局设备行(%d 行)失败: %w", len(deviceRows), err)
}
}
rep.RowsWritten = total + len(deviceRows)
rep.DaysWritten = len(daysSet)
return rep, nil
}
// analyzeDeviceFromLicense 从 console 主库的 license_<产品> 码表统计每产品的激活/绑定设备数。
// 表名 license_<pid 十六进制>(每产品一表),记录设备码真实产品归属——比 userdevice.productid
// (客户端声称、校验失败回退)准确。status>0=已激活、uid<>”=已绑定用户;切日用 usedtime(缺失兜底 2020-01-01)。
// license 表无应用/区域维度,故落全局行(app_id=” region=”)。返回待 upsert 的 StatsGlobalDay 行。
func analyzeDeviceFromLicense(now int64) ([]*pb.StatsGlobalDay, error) {
pg := postgres.GetSys()
var tnames []string
if err := pg.Raw("SELECT table_name FROM information_schema.tables WHERE table_name ~ '^license_[0-9a-f]+$'").Scan(&tnames).Error; err != nil {
return nil, fmt.Errorf("列出 license 码表失败: %w", err)
}
type licRow struct {
D string `gorm:"column:d"`
Active int64 `gorm:"column:active"`
Bind int64 `gorm:"column:bind"`
}
rows := make([]*pb.StatsGlobalDay, 0, 256)
for _, tn := range tnames {
pid64, perr := strconv.ParseUint(strings.TrimPrefix(tn, "license_"), 16, 32)
if perr != nil {
continue
}
var lr []*licRow
q := "SELECT to_char(to_timestamp(COALESCE(NULLIF(usedtime,0), 1577808000)), 'YYYY-MM-DD') AS d, " +
"COUNT(*) FILTER (WHERE status > 0) AS active, COUNT(*) FILTER (WHERE uid <> '') AS bind " +
"FROM " + tn + " WHERE status > 0 OR uid <> '' GROUP BY d"
if err := pg.Raw(q).Scan(&lr).Error; err != nil {
continue // 个别表结构异常不阻断整体
}
for _, r := range lr {
if r.Active == 0 && r.Bind == 0 {
continue
}
rows = append(rows, &pb.StatsGlobalDay{
AppId: "", ProductId: uint32(pid64), Region: "", StatDay: dayStrToNum(r.D),
ActiveDeviceCount: r.Active, BindDeviceCount: r.Bind, UpdateTime: now,
})
}
}
return rows, nil
}