11 changed files with 380 additions and 15 deletions
@ -0,0 +1,296 @@ |
|||
// 一次性迁移工具:把各应用**业务库**(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", ®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) |
|||
} |
|||
Loading…
Reference in new issue