// 一次性迁移工具:把各应用**业务库**(MySQL)的 goods 表(pb.DBGoods)拷贝到 // console 主库(Postgres)的公共商品表 app_goods(comm.AppGoods),按 app_name 隔离。 // // · 业务库只读——绝不修改/删除应用本地 goods 表; // · 目标表按 (app_name, goods_id) 判重:默认跳过已存在的商品(幂等,可重复执行), // 加 -update 则用业务库数据覆盖已存在行(仅覆盖商品字段,不动 id/app_name/createtime); // · -dry 只打印将要做什么,不写库。 // // 用法示例(DSN 也可用环境变量 MYSQL_DSN / PG_DSN 传,避免写进命令历史): // // go run ./tools/importappgoods -mysql '<业务库DSN>' -pg '' \ // -apps 'earphone=Voitrans,lumi=Lumi,deapsound=Deapsound,deepglass=Deepglass,echomeet=Echomeet' -dry // // -apps 是「MySQL 库名=应用名」的映射,应用名须与 console 注册表 app_registry.app_name 一致, // 否则后台按应用筛选时看不到这批商品。 package main import ( "flag" "fmt" "os" "sort" "strings" "time" "yunyan/comm" "gorm.io/driver/mysql" "gorm.io/driver/postgres" "gorm.io/gorm" "gorm.io/gorm/logger" ) // srcGoods 业务库 goods 表的读取结构(字段取自 pb.DBGoods,只读用,故不引 pb 以免 gorm tag 干扰)。 type srcGoods struct { Id string `gorm:"column:id"` Enable bool `gorm:"column:enable"` Name string `gorm:"column:name"` Localname string `gorm:"column:localname"` Price int32 `gorm:"column:price"` Usagetype int32 `gorm:"column:usagetype"` Ainum int64 `gorm:"column:ainum"` Meetnum int64 `gorm:"column:meetnum"` Tradenum int64 `gorm:"column:tradenum"` Viptime int32 `gorm:"column:viptime"` } func main() { var ( mysqlDsn = flag.String("mysql", os.Getenv("MYSQL_DSN"), "业务库 MySQL DSN(库名部分会被 -apps 的库名逐个替换)") pgDsn = flag.String("pg", os.Getenv("PG_DSN"), "console 主库 Postgres DSN") appsArg = flag.String("apps", "", "MySQL库名=应用名,逗号分隔,如 earphone=Voitrans,lumi=Lumi") dry = flag.Bool("dry", false, "只打印计划,不写库") verify = flag.Bool("verify", false, "只读校验:逐条比对业务库与 app_goods 的字段是否一致") update = flag.Bool("update", false, "已存在的商品用业务库数据覆盖(默认跳过)") ) flag.Parse() if *mysqlDsn == "" || *pgDsn == "" || *appsArg == "" { fmt.Fprintln(os.Stderr, "缺少参数:-mysql / -pg / -apps 必填") flag.Usage() os.Exit(2) } pairs, err := parseApps(*appsArg) if err != nil { fatal(err) } pg, err := gorm.Open(postgres.Open(withExecQueryMode(*pgDsn)), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) if err != nil { fatal(fmt.Errorf("连接 console 主库失败: %w", err)) } if sqlDB, err := pg.DB(); err == nil { // 与 lego/sys/postgres 一样对 pooler 客气点 sqlDB.SetMaxOpenConns(2) sqlDB.SetMaxIdleConns(1) } // 表通常已由 console 启动时建好;只在确实缺表时建(不能直接 AutoMigrate—— // 走 Supabase pooler 时它对已存在的表会报 42P07 relation already exists)。 // 用 to_regclass 按 search_path 找表:gorm 的 HasTable 查的是 CURRENT_SCHEMA(), // 经 Supabase pooler 连接时它与实际 search_path 不一致,会把已存在的表判成不存在。 var tableExists bool if err := pg.Raw("SELECT to_regclass(?) IS NOT NULL", comm.TableAppGoods).Scan(&tableExists).Error; err != nil { fatal(fmt.Errorf("检查表 %s 是否存在失败: %w", comm.TableAppGoods, err)) } if !tableExists { if err := pg.Table(comm.TableAppGoods).AutoMigrate(&comm.AppGoods{}); err != nil { fatal(fmt.Errorf("建表 %s 失败: %w", comm.TableAppGoods, err)) } fmt.Printf("已创建表 %s\n", comm.TableAppGoods) } // 注册表里的应用名,用于提醒拼写不一致(只提醒,不阻断) var registered []string if err := pg.Table("app_registry").Order("app_name").Pluck("app_name", ®istered).Error; err == nil { fmt.Printf("console 注册表现有应用: %s\n", strings.Join(registered, ", ")) } now := time.Now().Unix() var totalIns, totalUpd, totalSkip int for _, p := range pairs { dsn := replaceMySQLDB(*mysqlDsn, p.db) src, err := gorm.Open(mysql.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) if err != nil { fmt.Printf("✗ %-12s 连接业务库 %s 失败: %v\n", p.app, p.db, err) continue } var rows []srcGoods if err := src.Table(comm.TableGoods).Order("id").Find(&rows).Error; err != nil { fmt.Printf("✗ %-12s 读取 %s.goods 失败: %v\n", p.app, p.db, err) closeDB(src) continue } closeDB(src) if !contains(registered, p.app) && len(registered) > 0 { fmt.Printf("⚠ %-12s 不在 console 注册表 app_registry 中,后台按应用筛选可能看不到\n", p.app) } var ins, upd, skip int for i, r := range rows { var exist comm.AppGoods err := pg.Table(comm.TableAppGoods).Where("app_name=? AND goods_id=?", p.app, r.Id).Take(&exist).Error if *verify { switch { case err == gorm.ErrRecordNotFound: fmt.Printf("✗ %-12s %s 目标表缺失\n", p.app, r.Id) ins++ case err != nil: fmt.Printf("✗ %-12s 查询 %s 失败: %v\n", p.app, r.Id, err) default: if d := diffGoods(&r, &exist); d != "" { fmt.Printf("✗ %-12s %s 字段不一致: %s\n", p.app, r.Id, d) upd++ } else { skip++ } } continue } switch { case err == gorm.ErrRecordNotFound: m := &comm.AppGoods{ AppName: p.app, GoodsId: r.Id, Enable: r.Enable, Name: r.Name, Price: r.Price, Usagetype: r.Usagetype, Ainum: r.Ainum, Meetnum: r.Meetnum, Tradenum: r.Tradenum, Viptime: r.Viptime, Localname: r.Localname, Sort: int32((i + 1) * 10), Remark: "由业务库 " + p.db + ".goods 导入", Createtime: now, Updatetime: now, } if !*dry { if err := pg.Table(comm.TableAppGoods).Create(m).Error; err != nil { fmt.Printf("✗ %-12s 写入 %s 失败: %v\n", p.app, r.Id, err) continue } } ins++ case err != nil: fmt.Printf("✗ %-12s 查询 %s 失败: %v\n", p.app, r.Id, err) case *update: exist.Enable, exist.Name, exist.Price, exist.Usagetype = r.Enable, r.Name, r.Price, r.Usagetype exist.Ainum, exist.Meetnum, exist.Tradenum, exist.Viptime = r.Ainum, r.Meetnum, r.Tradenum, r.Viptime exist.Localname, exist.Updatetime = r.Localname, now if !*dry { if err := pg.Table(comm.TableAppGoods).Save(&exist).Error; err != nil { fmt.Printf("✗ %-12s 更新 %s 失败: %v\n", p.app, r.Id, err) continue } } upd++ default: skip++ } } if *verify { fmt.Printf("· %-12s(%s) 源 %2d 条 → 一致 %2d 不一致 %2d 缺失 %2d\n", p.app, p.db, len(rows), skip, upd, ins) } else { fmt.Printf("· %-12s(%s) 源 %2d 条 → 新增 %2d 更新 %2d 跳过(已存在) %2d\n", p.app, p.db, len(rows), ins, upd, skip) } totalIns, totalUpd, totalSkip = totalIns+ins, totalUpd+upd, totalSkip+skip } switch { case *verify: fmt.Printf("\n校验完成:一致 %d,不一致 %d,缺失 %d\n", totalSkip, totalUpd, totalIns) case *dry: fmt.Printf("\nDRY-RUN(未写库):新增 %d,更新 %d,跳过 %d\n", totalIns, totalUpd, totalSkip) default: fmt.Printf("\n已写入:新增 %d,更新 %d,跳过 %d\n", totalIns, totalUpd, totalSkip) } // 目标表现状 type stat struct { AppName string N int64 } var stats []stat if err := pg.Table(comm.TableAppGoods).Select("app_name, count(*) as n").Group("app_name").Order("app_name").Scan(&stats).Error; err == nil { fmt.Println("app_goods 现有数据:") for _, s := range stats { fmt.Printf(" %-12s %d\n", s.AppName, s.N) } } } // diffGoods 比对源(业务库)与目标(app_goods)的商品字段,返回不一致项的说明(空串=完全一致)。 // 只比对「拷贝过来的」字段,app_goods 独有的 sort/remark/时间戳不参与。 func diffGoods(s *srcGoods, d *comm.AppGoods) string { var bad []string cmp := func(name string, a, b any) { if fmt.Sprint(a) != fmt.Sprint(b) { bad = append(bad, fmt.Sprintf("%s(%v→%v)", name, a, b)) } } cmp("enable", s.Enable, d.Enable) cmp("name", s.Name, d.Name) cmp("price", s.Price, d.Price) cmp("usagetype", s.Usagetype, d.Usagetype) cmp("ainum", s.Ainum, d.Ainum) cmp("meetnum", s.Meetnum, d.Meetnum) cmp("tradenum", s.Tradenum, d.Tradenum) cmp("viptime", s.Viptime, d.Viptime) if s.Localname != d.Localname { bad = append(bad, "localname") } return strings.Join(bad, " ") } // withExecQueryMode 与 lego/sys/postgres 保持一致:Supabase 的事务型 pooler(6543) // 不保留会话级 prepared statement,pgx 默认的 cache_statement 会随机报 // "prepared statement ... does not exist/already exists",必须改用 cache_describe。 func withExecQueryMode(dsn string) string { if strings.Contains(dsn, "default_query_exec_mode") { return dsn } sep := "?" if strings.Contains(dsn, "?") { sep = "&" } return dsn + sep + "default_query_exec_mode=cache_describe" } type appPair struct{ db, app string } func parseApps(s string) (out []appPair, err error) { for _, seg := range strings.Split(s, ",") { seg = strings.TrimSpace(seg) if seg == "" { continue } db, app, ok := strings.Cut(seg, "=") if !ok || strings.TrimSpace(db) == "" || strings.TrimSpace(app) == "" { return nil, fmt.Errorf("-apps 片段格式应为 库名=应用名,收到: %q", seg) } out = append(out, appPair{strings.TrimSpace(db), strings.TrimSpace(app)}) } if len(out) == 0 { return nil, fmt.Errorf("-apps 为空") } sort.Slice(out, func(i, j int) bool { return out[i].app < out[j].app }) return out, nil } // replaceMySQLDB 把 DSN 里 ")/" 之后、"?" 之前的库名换成目标库(同一实例多库共用一份账号)。 func replaceMySQLDB(dsn, db string) string { i := strings.LastIndex(dsn, ")/") if i < 0 { return dsn } head := dsn[:i+2] rest := dsn[i+2:] if j := strings.Index(rest, "?"); j >= 0 { return head + db + rest[j:] } return head + db } func contains(list []string, s string) bool { for _, v := range list { if v == s { return true } } return false } func closeDB(db *gorm.DB) { if sqlDB, err := db.DB(); err == nil { _ = sqlDB.Close() } } func fatal(err error) { fmt.Fprintln(os.Stderr, "✗", err) os.Exit(1) }