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
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", ®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)
|
|
}
|
|
|