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.
 
 
 
 
 
 

296 lines
10 KiB

// 一次性迁移工具:把各应用**业务库**(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 '<console主库DSN>' \
// -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", &registered).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)
}