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_(每产品一表),记录设备码真实产品归属——比 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 }