updae
This commit is contained in:
11
Makefile
11
Makefile
@ -8,13 +8,12 @@ build:
|
||||
run: build
|
||||
./bin/server
|
||||
|
||||
# 表结构由 GORM AutoMigrate 在 server 启动时幂等创建,无需手动 SQL。
|
||||
# 跨环境搬运全库数据(结构 + 数据)用 dbtool:
|
||||
# go run ./cmd/dbtool dump -out db_dump.json
|
||||
# go run ./cmd/dbtool import -in db_dump.json # 整机重建
|
||||
# go run ./cmd/dbtool import -in db_dump.json -data-only # 仅导数据(结构已由 AutoMigrate 建好)
|
||||
# 表结构不由服务启动自动创建:全库「结构 + 索引 + 数据」统一由 dbtool 导出的 SQL 维护。
|
||||
# 导出:go run ./cmd/dbtool dump -out db/backups/db_dump.sql
|
||||
# 导入:psql -U fashion -d fashion -v ON_ERROR_STOP=1 -f db/backups/db_dump.sql
|
||||
# 目标库须为空且已启用 pgvector 扩展(导入时可加 -with-extension 让 dump 自带 CREATE EXTENSION)。
|
||||
migrate:
|
||||
@echo "Schema is auto-migrated by server on startup (GORM AutoMigrate). Use 'go run ./cmd/dbtool' to move data between environments."
|
||||
@echo "Schema is maintained by 'go run ./cmd/dbtool dump' -> psql -f (server no longer auto-migrates)."
|
||||
|
||||
# 跑单测
|
||||
test:
|
||||
|
||||
15
README.md
15
README.md
@ -145,13 +145,18 @@ go build -o bin/dbtool ./cmd/dbtool
|
||||
|
||||
### 数据库
|
||||
|
||||
表结构随服务启动由 GORM AutoMigrate 幂等创建(含 pgvector 扩展、去重索引与语义 embedding 表)。跨环境搬运全库数据用 dbtool:
|
||||
表结构**不由服务启动自动创建**:全库「结构 + 索引 + 约束 + 数据」统一由 dbtool 导出的**单个纯 SQL 文件**维护(不依赖 pg_dump)。
|
||||
|
||||
```bash
|
||||
go run ./cmd/dbtool dump -out db_dump.json # 导出当前库
|
||||
go run ./cmd/dbtool import -in db_dump.json # 在目标环境重建并回灌
|
||||
# 导出(源机器)
|
||||
go run ./cmd/dbtool dump -out db/backups/db_dump.sql
|
||||
|
||||
# 导入(目标机器;库须为空,且已启用 pgvector 扩展)
|
||||
psql -U fashion -d fashion -v ON_ERROR_STOP=1 -f db/backups/db_dump.sql
|
||||
```
|
||||
|
||||
目标库若尚未启用 pgvector,导出时加 `-with-extension`,让文件自带 `CREATE EXTENSION IF NOT EXISTS vector`。
|
||||
|
||||
### 优雅关闭
|
||||
|
||||
监听 `SIGINT` / `SIGTERM`,等待在途请求完成(最长 `shutdown_timeout` 秒)后退出,并通知入库 worker 停止领取新任务。
|
||||
@ -165,12 +170,12 @@ backend/
|
||||
├── configs/config.yml # 唯一配置源
|
||||
├── cmd/
|
||||
│ ├── server/ # 组合根:三端口 HTTP + 入库 worker
|
||||
│ ├── dbtool/ # 跨环境数据搬运
|
||||
│ ├── dbtool/ # 全库导出为纯 SQL(结构 + 索引 + 数据)
|
||||
│ └── dbdiag/ # 数据库诊断
|
||||
├── internal/
|
||||
│ ├── config/ # 配置加载
|
||||
│ ├── model/ # 实体(brand / runway / street_snap / user / 草稿 / ingest)
|
||||
│ ├── database/ # 连接池 + AutoMigrate + 去重 schema
|
||||
│ ├── database/ # 连接池(结构由 dbtool 导出的 SQL 维护,不做 DDL)
|
||||
│ ├── repository/ # 数据访问接口 + GORM 实现
|
||||
│ ├── service/ # 业务逻辑
|
||||
│ ├── dto/ # 请求 / 响应结构体
|
||||
|
||||
3008
cmd/dbtool/db_dump.sql
Normal file
3008
cmd/dbtool/db_dump.sql
Normal file
File diff suppressed because one or more lines are too long
@ -1,28 +1,38 @@
|
||||
// Command dbtool 是后台管理用的数据库搬运脚本。
|
||||
// Command dbtool 是后台管理用的数据库导出脚本。
|
||||
//
|
||||
// 用途:在不同开发电脑 / 环境之间搬运 PostgreSQL 库(fashion)的「数据」,
|
||||
// 省去手动 pg_dump / 重新 seed 的麻烦。完全复用后端既有的 pgx 驱动与
|
||||
// configs/config.yml,不依赖任何外部二进制(pg_dump 等)。
|
||||
// 用途:把一个 PostgreSQL 库(fashion)的「结构 + 索引 + 约束 + 数据」导出为
|
||||
// **单个纯 SQL 文件**,拷到另一台机器后直接用 psql 灌入即可 —— 无需 pg_dump。
|
||||
//
|
||||
// 子命令:
|
||||
// dbtool dump [-config <yml>] [-out <file.sql>] [-with-extension]
|
||||
//
|
||||
// dbtool dump -out db_dump.json 导出当前库全部表的数据 + 列元信息为单个 JSON 文件
|
||||
// dbtool import -in db_dump.json 把 JSON 数据灌入目标库(结构先由 AutoMigrate 收敛)
|
||||
// 导入方式(目标机器):
|
||||
//
|
||||
// 结构在哪里定义:**只在代码里** —— model 的 GORM tag + database.EnsureDedupSchema。
|
||||
// import 会先调 database.AutoMigrate + EnsureDedupSchema,把目标库收敛到与当前代码
|
||||
// 一致(建表 / 加列 / 加索引 / 建 pgvector 扩展 / 建 HNSW 索引),然后再灌数据。
|
||||
// psql -U fashion -d fashion -v ON_ERROR_STOP=1 -f db_dump.sql
|
||||
//
|
||||
// 历史实现曾让 import 自己读 information_schema 拼 DDL 来重建表,那条路必然丢信息
|
||||
// (类型修饰符 vector(64)、identity 自增、索引、扩展),已移除。
|
||||
// 结构从哪来:**唯一入口就是本工具的 dump 输出**。项目已不再在服务启动时自动迁移表结构
|
||||
// (原 database.AutoMigrate / EnsureDedupSchema 已移除),重建一个库就是「dump 出来再灌进去」。
|
||||
//
|
||||
// 为什么没有 import 子命令:数据段用的是 SQL 标准的 COPY ... FROM stdin,它依赖 PostgreSQL
|
||||
// 前端的 copy 子协议,Go 的 database/sql 无法执行这类脚本。导入统一交给 psql。
|
||||
//
|
||||
// 与 pg_dump 的关系:本工具不依赖任何外部二进制,等于把「读系统目录 → 拼 DDL → 导数据」
|
||||
// 自己实现一遍,因此**只覆盖本项目实际用到的对象**:
|
||||
//
|
||||
// 表、列(类型含修饰符 / 默认值 / identity / NOT NULL)、主键、索引(含 pgvector 的 HNSW)、
|
||||
// 数据(COPY)、序列当前值。
|
||||
//
|
||||
// 视图 / 触发器 / 外键 / 注释 / 权限不导出(当前库中也不存在)。
|
||||
//
|
||||
// 前置条件:目标库必须已启用 pgvector 扩展,否则 vector(64) 列建不出来。
|
||||
// 默认**不**导出扩展语句(可用 -with-extension 带上)。若目标容器由 scripts/pgvector 的
|
||||
// compose 启动,initdb/01-extensions.sql 会在数据卷首次初始化时自动创建扩展。
|
||||
//
|
||||
// 连接信息来自 configs/config.yml 的 database 段(同后端服务),可用 -config 指定其它配置。
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log"
|
||||
@ -30,70 +40,31 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
|
||||
"fashionapi/internal/config"
|
||||
"fashionapi/internal/database"
|
||||
)
|
||||
|
||||
// columnMeta 是单列的元信息,用于导入时的类型转换与 identity 判定。
|
||||
type columnMeta struct {
|
||||
Name string `json:"name"`
|
||||
// Type 为 udt_name,如 int8 / varchar / timestamptz / vector。
|
||||
// 注意:它**不含**类型修饰符(vector(64) 会退化成 vector),
|
||||
// 因此不可用于重建 DDL —— 结构一律交给 database.AutoMigrate。
|
||||
Type string `json:"type"`
|
||||
IsNullable bool `json:"is_nullable"`
|
||||
// IsIdentity 取 information_schema.columns.is_identity。
|
||||
// 必须读这一列:identity 列的 column_default 是 NULL,
|
||||
// 用 column_default LIKE 'nextval(%' 去猜会把它全部漏判为「非自增」,
|
||||
// 于是重建出的表 id 没有默认值、插入即违反 NOT NULL。
|
||||
IsIdentity bool `json:"is_identity"`
|
||||
Default string `json:"default,omitempty"`
|
||||
}
|
||||
|
||||
// tableDump 是单张表的导出结构。
|
||||
type tableDump struct {
|
||||
Name string `json:"name"`
|
||||
Columns []columnMeta `json:"columns"`
|
||||
Pk []string `json:"pk"`
|
||||
Rows [][]any `json:"rows"` // 每行为一组值;nil 表示 NULL,string 表示值
|
||||
}
|
||||
|
||||
// dumpFile 是 dump 输出的顶层结构。
|
||||
type dumpFile struct {
|
||||
Version int `json:"version"`
|
||||
GeneratedAt string `json:"generated_at"`
|
||||
Database string `json:"database"`
|
||||
Tables []tableDump `json:"tables"`
|
||||
}
|
||||
|
||||
// verbose 控制是否打印每条执行的 SQL(用于排查)。由 dump/import 子命令的 -v 开关设置。
|
||||
var verbose bool
|
||||
|
||||
// vlog 在开启 verbose 时打印调试信息。
|
||||
func vlog(format string, args ...any) {
|
||||
if verbose {
|
||||
log.Printf("[sql] "+format, args...)
|
||||
}
|
||||
}
|
||||
|
||||
// pgDetail 提取 PostgreSQL 报错的位置/消息,便于定位语法错误。
|
||||
func pgDetail(err error) string {
|
||||
var pgErr *pgconn.PgError
|
||||
if errors.As(err, &pgErr) {
|
||||
return fmt.Sprintf(" [position=%d message=%q where=%q]", pgErr.Position, pgErr.Message, pgErr.Where)
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// quoteLit 安全包裹 SQL 字符串字面量(转义单引号),用于内联可信的内部标识符/表名。
|
||||
func quoteLit(s string) string {
|
||||
return "'" + strings.ReplaceAll(s, "'", "''") + "'"
|
||||
// column 是单列的结构信息,用于生成 CREATE TABLE、COPY 列清单与序列重置语句。
|
||||
type column struct {
|
||||
Name string
|
||||
// Type 为 format_type(atttypid, atttypmod) 的结果,**含**类型修饰符:
|
||||
// vector(64) / character varying(255) / bigint。
|
||||
// 不能用 information_schema.udt_name —— 它会丢掉修饰符(vector(64) 变成 vector),
|
||||
// 建出来的列没有维度,随后 HNSW 索引会报 "column does not have dimensions"。
|
||||
Type string
|
||||
// NotNull 取自 pg_attribute.attnotnull。
|
||||
NotNull bool
|
||||
// Default 取自 pg_get_expr(adbin, adrelid),如 0 / ''::character varying / nextval('x'::regclass)。
|
||||
Default string
|
||||
// Identity 取自 pg_attribute.attidentity:'d'=BY DEFAULT、'a'=ALWAYS、''=非 identity。
|
||||
Identity string
|
||||
// AutoIncrement 表示该列由序列自动赋值(identity 列,或默认值是 nextval 的 serial 列),
|
||||
// 导出后会追加 setval 把序列推到 MAX(col)+1。
|
||||
//
|
||||
// 判定必须读 attidentity:identity 列的 column_default 是 NULL,
|
||||
// 用 column_default LIKE 'nextval(%' 去猜会把它们全部漏判,重建出的表 id 就没有自增。
|
||||
AutoIncrement bool
|
||||
}
|
||||
|
||||
func main() {
|
||||
@ -104,8 +75,6 @@ func main() {
|
||||
switch os.Args[1] {
|
||||
case "dump":
|
||||
runDump(os.Args[2:])
|
||||
case "import":
|
||||
runImport(os.Args[2:])
|
||||
default:
|
||||
usage()
|
||||
os.Exit(2)
|
||||
@ -113,16 +82,20 @@ func main() {
|
||||
}
|
||||
|
||||
func usage() {
|
||||
fmt.Println(`dbtool - 后台数据库搬运脚本(PostgreSQL)
|
||||
fmt.Println(`dbtool - 后台数据库导出脚本(PostgreSQL)
|
||||
|
||||
用法:
|
||||
dbtool dump [-config <yml>] [-out <file>] 导出全库数据到 JSON
|
||||
dbtool import [-config <yml>] [-in <file>] 把 JSON 数据灌入目标库
|
||||
dbtool dump [-config <yml>] [-out <file>] [-with-extension]
|
||||
导出「结构 + 主键 / 索引 + 数据」为单个纯 SQL 文件
|
||||
|
||||
导入(目标机器):
|
||||
psql -U fashion -d fashion -v ON_ERROR_STOP=1 -f <file>
|
||||
|
||||
说明:
|
||||
连接信息读取 configs/config.yml 的 database 段(可用 -config 覆盖)。
|
||||
import 会先把目标库结构收敛到与当前代码一致(AutoMigrate + EnsureDedupSchema),
|
||||
再清空同名表并按 JSON 重灌数据。结构定义只在代码里,dump 文件不含 DDL。`)
|
||||
-with-extension 会在文件开头加 CREATE EXTENSION IF NOT EXISTS vector;
|
||||
默认不加,此时目标库须已启用 pgvector,否则 vector 列建不出来。
|
||||
目标库必须是**空库**(脚本按 CREATE TABLE IF NOT EXISTS + COPY 写入,不会清表)。`)
|
||||
}
|
||||
|
||||
// openDB 按 config 加载 DSN 并探活。
|
||||
@ -142,14 +115,16 @@ func openDB(cfgPath string) (*sql.DB, *config.DatabaseConfig, error) {
|
||||
return db, &cfg.Database, nil
|
||||
}
|
||||
|
||||
// runDump 导出全库。
|
||||
// runDump 导出「结构 → 数据 → 主键/索引 → 序列」四段到单个 SQL 文件。
|
||||
//
|
||||
// 段落顺序刻意与 pg_dump 一致:先建表、再灌数据、最后建约束与索引。
|
||||
// 索引放在数据之后建,既快又不会在导入过程中被反复维护。
|
||||
func runDump(args []string) {
|
||||
fs := flag.NewFlagSet("dump", flag.ExitOnError)
|
||||
out := fs.String("out", "db_dump.json", "导出文件路径")
|
||||
out := fs.String("out", "db_dump.sql", "导出文件路径")
|
||||
cfgPath := fs.String("config", "", "配置文件路径(默认 configs/config.yml)")
|
||||
v := fs.Bool("v", false, "打印每条执行的 SQL,便于排查")
|
||||
withExt := fs.Bool("with-extension", false, "在文件开头加 CREATE EXTENSION IF NOT EXISTS vector")
|
||||
_ = fs.Parse(args)
|
||||
verbose = *v
|
||||
|
||||
db, dcfg, err := openDB(*cfgPath)
|
||||
if err != nil {
|
||||
@ -157,127 +132,103 @@ func runDump(args []string) {
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
f, err := os.Create(*out)
|
||||
if err != nil {
|
||||
log.Fatalf("✗ 创建文件 %s 失败: %v", *out, err)
|
||||
}
|
||||
defer f.Close()
|
||||
w := bufio.NewWriter(f)
|
||||
|
||||
writeHeader(w, dcfg)
|
||||
if *withExt {
|
||||
fmt.Fprintln(w, "CREATE EXTENSION IF NOT EXISTS vector;")
|
||||
fmt.Fprintln(w)
|
||||
}
|
||||
|
||||
tables, err := listTables(db)
|
||||
if err != nil {
|
||||
log.Fatalf("✗ 列举表失败: %v", err)
|
||||
}
|
||||
|
||||
df := dumpFile{
|
||||
Version: 1,
|
||||
GeneratedAt: time.Now().UTC().Format(time.RFC3339),
|
||||
Database: dcfg.Name,
|
||||
Tables: make([]tableDump, 0, len(tables)),
|
||||
if len(tables) == 0 {
|
||||
log.Fatalf("✗ 库中没有表(schema=public),请确认 -config 指向的库是否正确")
|
||||
}
|
||||
|
||||
cols := make(map[string][]column, len(tables))
|
||||
|
||||
// ① 结构
|
||||
fmt.Fprintln(w, "-- ==================== 结构 ====================")
|
||||
for _, t := range tables {
|
||||
log.Printf("→ 导出表 %s ...", t)
|
||||
td, err := dumpTable(db, t)
|
||||
cs, err := listColumns(db, t)
|
||||
if err != nil {
|
||||
log.Fatalf("✗ 导出表 %s 失败: %v", t, err)
|
||||
log.Fatalf("✗ 读取 %s 的列失败: %v", t, err)
|
||||
}
|
||||
df.Tables = append(df.Tables, td)
|
||||
cols[t] = cs
|
||||
writeCreateTable(w, t, cs)
|
||||
}
|
||||
log.Printf("→ 结构:%d 张表", len(tables))
|
||||
|
||||
// ② 数据
|
||||
fmt.Fprintln(w, "-- ==================== 数据 ====================")
|
||||
var total int64
|
||||
for _, t := range tables {
|
||||
n, err := writeCopyData(w, db, t, cols[t])
|
||||
if err != nil {
|
||||
log.Fatalf("✗ 导出 %s 数据失败: %v", t, err)
|
||||
}
|
||||
total += n
|
||||
log.Printf(" ✓ %s: %d 行", t, n)
|
||||
}
|
||||
|
||||
data, err := json.MarshalIndent(df, "", " ")
|
||||
if err != nil {
|
||||
log.Fatalf("✗ 序列化失败: %v", err)
|
||||
// ③ 主键 / 索引
|
||||
fmt.Fprintln(w, "-- ==================== 主键 / 索引 ====================")
|
||||
for _, t := range tables {
|
||||
if err := writeConstraints(w, db, t); err != nil {
|
||||
log.Fatalf("✗ 导出 %s 约束失败: %v", t, err)
|
||||
}
|
||||
if err := writeIndexes(w, db, t); err != nil {
|
||||
log.Fatalf("✗ 导出 %s 索引失败: %v", t, err)
|
||||
}
|
||||
}
|
||||
if err := os.WriteFile(*out, data, 0o644); err != nil {
|
||||
|
||||
// ④ 序列当前值:COPY 写了显式 id,不会推进序列,必须手工推到 MAX+1。
|
||||
fmt.Fprintln(w, "-- ==================== 序列当前值 ====================")
|
||||
for _, t := range tables {
|
||||
writeSequenceResets(w, t, cols[t])
|
||||
}
|
||||
|
||||
if err := w.Flush(); err != nil {
|
||||
log.Fatalf("✗ 写文件 %s 失败: %v", *out, err)
|
||||
}
|
||||
log.Printf("✓ 已导出 %d 张表 -> %s", len(df.Tables), *out)
|
||||
log.Printf("✓ 已导出 %d 张表、%d 行 -> %s", len(tables), total, *out)
|
||||
}
|
||||
|
||||
// runImport 把 JSON 中的数据灌入目标库。
|
||||
//
|
||||
// 结构不在这里重建:先调 ensureSchema 用后端的 AutoMigrate + EnsureDedupSchema
|
||||
// 把目标库收敛到与当前代码一致(建表 / 加列 / 加索引 / 建 pgvector 扩展 / 建 HNSW 索引),
|
||||
// 本函数只负责搬数据。这样结构只有一个来源(model + EnsureDedupSchema),
|
||||
// 不会再出现「自拼 DDL 丢类型修饰符 / 丢 identity 自增 / 丢索引 / 不建扩展」那一类问题。
|
||||
//
|
||||
// 注意:导入先清空目标同名表再写入,属「以 dump 为准的整体覆盖」。
|
||||
func runImport(args []string) {
|
||||
fs := flag.NewFlagSet("import", flag.ExitOnError)
|
||||
in := fs.String("in", "db_dump.json", "导入文件路径")
|
||||
cfgPath := fs.String("config", "", "配置文件路径(默认 configs/config.yml)")
|
||||
_ = fs.Parse(args)
|
||||
|
||||
raw, err := os.ReadFile(*in)
|
||||
if err != nil {
|
||||
log.Fatalf("✗ 读文件 %s 失败: %v", *in, err)
|
||||
}
|
||||
var df dumpFile
|
||||
if err := json.Unmarshal(raw, &df); err != nil {
|
||||
log.Fatalf("✗ 解析 %s 失败: %v", *in, err)
|
||||
}
|
||||
|
||||
db, dcfg, err := openDB(*cfgPath)
|
||||
if err != nil {
|
||||
log.Fatalf("✗ %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
// ① 结构:与后端同一套定义(含 CREATE EXTENSION vector、HNSW 索引、老库兼容迁移)。
|
||||
log.Printf("→ 收敛目标库结构(AutoMigrate + EnsureDedupSchema)...")
|
||||
if err := ensureSchema(dcfg); err != nil {
|
||||
log.Fatalf("✗ 结构初始化失败: %v", err)
|
||||
}
|
||||
|
||||
// ② 数据:关闭外键 / 触发器,避免插入顺序受约束(整库覆盖无需保序)。
|
||||
if _, err := db.Exec("SET session_replication_role = 'replica'"); err != nil {
|
||||
log.Fatalf("✗ 关闭约束检查失败: %v", err)
|
||||
}
|
||||
defer db.Exec("SET session_replication_role = 'origin'")
|
||||
|
||||
for _, t := range df.Tables {
|
||||
// 先清空目标表(约束已由 replica 角色关闭),再插入。
|
||||
if _, err := db.Exec(fmt.Sprintf(`DELETE FROM "%s"`, t.Name)); err != nil {
|
||||
log.Fatalf("✗ 清空表 %s 失败: %v", t.Name, err)
|
||||
}
|
||||
if err := importRows(db, t); err != nil {
|
||||
log.Fatalf("✗ 导数据到 %s 失败: %v", t.Name, err)
|
||||
}
|
||||
log.Printf(" ✓ %s: %d 行", t.Name, len(t.Rows))
|
||||
}
|
||||
log.Printf("✓ 导入完成(%d 张表)", len(df.Tables))
|
||||
// writeHeader 写文件头与几个影响字面量解析的会话设置。
|
||||
func writeHeader(w *bufio.Writer, dcfg *config.DatabaseConfig) {
|
||||
fmt.Fprintf(w, "-- dbtool 导出:%s\n", dcfg.Addr())
|
||||
fmt.Fprintf(w, "-- 生成时间:%s\n", time.Now().Format(time.RFC3339))
|
||||
fmt.Fprintf(w, "-- 导入:psql -U %s -d %s -v ON_ERROR_STOP=1 -f <本文件>\n", dcfg.User, dcfg.Name)
|
||||
fmt.Fprintln(w, "--")
|
||||
fmt.Fprintln(w, "-- 注意:目标库必须是空库;本脚本不清表,重复执行会主键冲突。")
|
||||
fmt.Fprintln(w, "SET client_encoding = 'UTF8';")
|
||||
fmt.Fprintln(w, "SET standard_conforming_strings = on;")
|
||||
fmt.Fprintln(w, "SET check_function_bodies = false;")
|
||||
fmt.Fprintln(w, "SET client_min_messages = warning;")
|
||||
fmt.Fprintln(w)
|
||||
}
|
||||
|
||||
// ensureSchema 用后端的 database.AutoMigrate + EnsureDedupSchema 把目标库结构
|
||||
// 收敛到与当前代码一致。这是「结构定义只有一个来源」的落点:
|
||||
// 表 / 列来自 model 的 GORM tag,扩展与 HNSW 索引来自 EnsureDedupSchema。
|
||||
func ensureSchema(dcfg *config.DatabaseConfig) error {
|
||||
gdb, err := gorm.Open(postgres.Open(dcfg.DSN()), &gorm.Config{
|
||||
Logger: logger.Default.LogMode(logger.Warn),
|
||||
SkipDefaultTransaction: true,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("连接数据库失败: %w", err)
|
||||
}
|
||||
sqlDB, err := gdb.DB()
|
||||
if err != nil {
|
||||
return fmt.Errorf("获取底层连接池失败: %w", err)
|
||||
}
|
||||
defer sqlDB.Close()
|
||||
|
||||
if err := database.AutoMigrate(gdb); err != nil {
|
||||
return fmt.Errorf("AutoMigrate: %w", err)
|
||||
}
|
||||
if err := database.EnsureDedupSchema(gdb); err != nil {
|
||||
return fmt.Errorf("EnsureDedupSchema: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// listTables 返回 public 模式下所有基表名。
|
||||
// listTables 返回 public 下全部基表名(按名字排序,保证导出可复现)。
|
||||
func listTables(db *sql.DB) ([]string, error) {
|
||||
rows, err := db.Query(`
|
||||
SELECT table_name FROM information_schema.tables
|
||||
WHERE table_schema = 'public' AND table_type = 'BASE TABLE'
|
||||
ORDER BY table_name`)
|
||||
SELECT c.relname
|
||||
FROM pg_class c
|
||||
JOIN pg_namespace n ON n.oid = c.relnamespace
|
||||
WHERE n.nspname = 'public' AND c.relkind = 'r'
|
||||
ORDER BY c.relname`)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []string
|
||||
for rows.Next() {
|
||||
var name string
|
||||
@ -289,217 +240,245 @@ func listTables(db *sql.DB) ([]string, error) {
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// dumpTable 导出单张表的列元信息 + 数据。
|
||||
func dumpTable(db *sql.DB, name string) (tableDump, error) {
|
||||
cols, pk, err := dumpColumns(db, name)
|
||||
if err != nil {
|
||||
return tableDump{}, err
|
||||
}
|
||||
// listColumns 读取单表全部列的完整结构信息。
|
||||
func listColumns(db *sql.DB, table string) ([]column, error) {
|
||||
q := fmt.Sprintf(`
|
||||
SELECT a.attname,
|
||||
format_type(a.atttypid, a.atttypmod),
|
||||
a.attnotnull,
|
||||
COALESCE(a.attidentity, ''),
|
||||
COALESCE(pg_get_expr(d.adbin, d.adrelid), '')
|
||||
FROM pg_attribute a
|
||||
LEFT JOIN pg_attrdef d ON d.adrelid = a.attrelid AND d.adnum = a.attnum
|
||||
WHERE a.attrelid = %s::regclass AND a.attnum > 0 AND NOT a.attisdropped
|
||||
ORDER BY a.attnum`, ql(qname(table)))
|
||||
|
||||
dataQuery := fmt.Sprintf(`SELECT * FROM "%s"`, name)
|
||||
vlog("dump data: %s", dataQuery)
|
||||
rows, err := db.Query(dataQuery)
|
||||
rows, err := db.Query(q)
|
||||
if err != nil {
|
||||
return tableDump{}, fmt.Errorf("读取数据失败(表 %s): %w\nSQL: %s", name, err, dataQuery)
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
colNames, err := rows.Columns()
|
||||
if err != nil {
|
||||
return tableDump{}, err
|
||||
}
|
||||
n := len(colNames)
|
||||
td := tableDump{Name: name, Columns: cols, Pk: pk, Rows: make([][]any, 0)}
|
||||
|
||||
scanPtrs := make([]any, n)
|
||||
raw := make([]sql.RawBytes, n)
|
||||
for i := range raw {
|
||||
scanPtrs[i] = &raw[i]
|
||||
}
|
||||
var out []column
|
||||
for rows.Next() {
|
||||
if err := rows.Scan(scanPtrs...); err != nil {
|
||||
return tableDump{}, err
|
||||
var c column
|
||||
if err := rows.Scan(&c.Name, &c.Type, &c.NotNull, &c.Identity, &c.Default); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
row := make([]any, n)
|
||||
for i := range raw {
|
||||
if raw[i] == nil {
|
||||
row[i] = nil // NULL
|
||||
} else {
|
||||
row[i] = string(raw[i]) // 全部按字符串搬运,导入时按列类型转换
|
||||
}
|
||||
}
|
||||
td.Rows = append(td.Rows, row)
|
||||
c.AutoIncrement = c.Identity != "" || isSerial(c)
|
||||
out = append(out, c)
|
||||
}
|
||||
return td, rows.Err()
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// dumpColumns 读取列元信息与主键。
|
||||
func dumpColumns(db *sql.DB, table string) ([]columnMeta, []string, error) {
|
||||
// is_identity 必须直接读 information_schema.columns.is_identity:
|
||||
// identity 列的 column_default 是 NULL,用 column_default LIKE 'nextval(%' 判断
|
||||
// 会把 GENERATED BY DEFAULT AS IDENTITY 的列全部漏判成「非自增」。
|
||||
colQuery := strings.Replace(`
|
||||
SELECT c.column_name,
|
||||
c.udt_name,
|
||||
(c.is_nullable = 'YES'),
|
||||
COALESCE(c.column_default, ''),
|
||||
(c.is_identity = 'YES')
|
||||
FROM information_schema.columns c
|
||||
WHERE c.table_schema = 'public' AND c.table_name = $1
|
||||
ORDER BY c.ordinal_position`, "$1", quoteLit(table), 1)
|
||||
vlog("dump columns: table=%q query=%s", table, colQuery)
|
||||
rows, err := db.Query(colQuery)
|
||||
// writeCreateTable 由列信息生成 CREATE TABLE。
|
||||
//
|
||||
// 两类自增列统一归一为 GENERATED BY DEFAULT AS IDENTITY:
|
||||
// - 原生 identity 列(attidentity 非空);
|
||||
// - serial 列(默认值是 nextval,如 image_embeddings.id)。
|
||||
//
|
||||
// 归一后不必再单独导出 CREATE SEQUENCE —— identity 子句会自动建序列,
|
||||
// 也避免了「默认值引用一个还不存在的序列」导致建表失败。
|
||||
func writeCreateTable(w *bufio.Writer, table string, cols []column) {
|
||||
fmt.Fprintf(w, "CREATE TABLE IF NOT EXISTS %s (\n", qname(table))
|
||||
for i, c := range cols {
|
||||
fmt.Fprintf(w, " %s %s", qi(c.Name), c.Type)
|
||||
switch {
|
||||
case c.AutoIncrement:
|
||||
fmt.Fprint(w, " GENERATED BY DEFAULT AS IDENTITY")
|
||||
case c.Default != "":
|
||||
fmt.Fprintf(w, " DEFAULT %s", c.Default)
|
||||
}
|
||||
// identity 列本身即 NOT NULL,不重复声明。
|
||||
if c.NotNull && !c.AutoIncrement {
|
||||
fmt.Fprint(w, " NOT NULL")
|
||||
}
|
||||
if i < len(cols)-1 {
|
||||
fmt.Fprint(w, ",")
|
||||
}
|
||||
fmt.Fprintln(w)
|
||||
}
|
||||
fmt.Fprintln(w, ");")
|
||||
fmt.Fprintln(w)
|
||||
}
|
||||
|
||||
// writeCopyData 用 COPY ... FROM stdin 导出单表数据,返回行数。
|
||||
//
|
||||
// 全部值按服务端文本表示搬运;NULL 写成 \N,其余转义反斜杠 / 制表符 / 换行 / 回车。
|
||||
func writeCopyData(w *bufio.Writer, db *sql.DB, table string, cols []column) (int64, error) {
|
||||
names := make([]string, len(cols))
|
||||
for i, c := range cols {
|
||||
names[i] = qi(c.Name)
|
||||
}
|
||||
colList := strings.Join(names, ", ")
|
||||
|
||||
fmt.Fprintf(w, "COPY %s (%s) FROM stdin;\n", qname(table), colList)
|
||||
|
||||
rows, err := db.Query(fmt.Sprintf("SELECT %s FROM %s", colList, qname(table)))
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("读取列元信息失败(表 %s): %w%s\nSQL: %s", table, err, pgDetail(err), strings.TrimSpace(colQuery))
|
||||
return 0, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var cols []columnMeta
|
||||
raw := make([]sql.RawBytes, len(cols))
|
||||
ptrs := make([]any, len(cols))
|
||||
for i := range raw {
|
||||
ptrs[i] = &raw[i]
|
||||
}
|
||||
|
||||
var count int64
|
||||
for rows.Next() {
|
||||
var name, udt, def string
|
||||
var nullable, isIdentity bool
|
||||
if err := rows.Scan(&name, &udt, &nullable, &def, &isIdentity); err != nil {
|
||||
return nil, nil, err
|
||||
if err := rows.Scan(ptrs...); err != nil {
|
||||
return count, err
|
||||
}
|
||||
cols = append(cols, columnMeta{
|
||||
Name: name, Type: udt, IsNullable: nullable, IsIdentity: isIdentity, Default: def,
|
||||
})
|
||||
for i := range raw {
|
||||
if i > 0 {
|
||||
w.WriteByte('\t')
|
||||
}
|
||||
if raw[i] == nil {
|
||||
w.WriteString(`\N`)
|
||||
continue
|
||||
}
|
||||
w.WriteString(escapeCopy(raw[i]))
|
||||
}
|
||||
w.WriteByte('\n')
|
||||
count++
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, nil, err
|
||||
return count, err
|
||||
}
|
||||
|
||||
pkQuery := strings.Replace(`
|
||||
SELECT kcu.column_name
|
||||
FROM information_schema.table_constraints tc
|
||||
JOIN information_schema.key_column_usage kcu
|
||||
ON kcu.constraint_name = tc.constraint_name AND kcu.table_schema = tc.table_schema
|
||||
WHERE tc.table_schema = 'public' AND tc.table_name = $1 AND tc.constraint_type = 'PRIMARY KEY'
|
||||
ORDER BY kcu.ordinal_position`, "$1", quoteLit(table), 1)
|
||||
vlog("dump pk: %s", strings.TrimSpace(pkQuery))
|
||||
pkRows, err := db.Query(pkQuery)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("读取主键失败(表 %s): %w\nSQL: %s", table, err, strings.TrimSpace(pkQuery))
|
||||
}
|
||||
defer pkRows.Close()
|
||||
var pk []string
|
||||
for pkRows.Next() {
|
||||
var c string
|
||||
if err := pkRows.Scan(&c); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
pk = append(pk, c)
|
||||
}
|
||||
return cols, pk, pkRows.Err()
|
||||
fmt.Fprintln(w, `\.`)
|
||||
fmt.Fprintln(w)
|
||||
return count, nil
|
||||
}
|
||||
|
||||
// castFor 返回列类型对应的 pg 类型转换后缀(用于导入时把字符串值转为正确类型)。
|
||||
func castFor(udt string) string {
|
||||
switch udt {
|
||||
case "int2", "int4", "int8":
|
||||
return "bigint"
|
||||
case "numeric":
|
||||
return "numeric"
|
||||
case "float4", "float8":
|
||||
return "double precision"
|
||||
case "bool":
|
||||
return "boolean"
|
||||
case "timestamp":
|
||||
return "timestamp"
|
||||
case "timestamptz":
|
||||
return "timestamptz"
|
||||
case "date":
|
||||
return "date"
|
||||
case "time":
|
||||
return "time"
|
||||
case "json", "jsonb":
|
||||
return "jsonb"
|
||||
case "uuid":
|
||||
return "uuid"
|
||||
default:
|
||||
return "" // text / varchar / char 等字符串类型无需转换
|
||||
}
|
||||
}
|
||||
// writeConstraints 导出主键等约束定义(pg_get_constraintdef 给出完整定义,无需自己拼)。
|
||||
func writeConstraints(w *bufio.Writer, db *sql.DB, table string) error {
|
||||
q := fmt.Sprintf(`
|
||||
SELECT con.conname, pg_get_constraintdef(con.oid)
|
||||
FROM pg_constraint con
|
||||
WHERE con.conrelid = %s::regclass
|
||||
ORDER BY CASE con.contype
|
||||
WHEN 'p' THEN 1 WHEN 'u' THEN 2 WHEN 'c' THEN 3 WHEN 'x' THEN 4 WHEN 'f' THEN 5 ELSE 9
|
||||
END,
|
||||
con.conname`, ql(qname(table)))
|
||||
|
||||
// importRows 把一张表的数据批量 INSERT 进目标库(单表一个事务,每批多行)。
|
||||
func importRows(db *sql.DB, t tableDump) error {
|
||||
if len(t.Rows) == 0 {
|
||||
return nil
|
||||
}
|
||||
hasIdentity := false
|
||||
for _, c := range t.Columns {
|
||||
if c.IsIdentity {
|
||||
hasIdentity = true
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
tx, err := db.Begin()
|
||||
rows, err := db.Query(q)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
defer rows.Close()
|
||||
|
||||
colList := make([]string, len(t.Columns))
|
||||
casts := make([]string, len(t.Columns))
|
||||
for i, c := range t.Columns {
|
||||
colList[i] = `"` + c.Name + `"`
|
||||
casts[i] = castFor(c.Type)
|
||||
}
|
||||
|
||||
const batchSize = 200
|
||||
for start := 0; start < len(t.Rows); start += batchSize {
|
||||
end := start + batchSize
|
||||
if end > len(t.Rows) {
|
||||
end = len(t.Rows)
|
||||
}
|
||||
chunk := t.Rows[start:end]
|
||||
|
||||
var sb strings.Builder
|
||||
ov := ""
|
||||
if hasIdentity {
|
||||
ov = " OVERRIDING SYSTEM VALUE"
|
||||
}
|
||||
sb.WriteString(fmt.Sprintf(`INSERT INTO "%s" (%s)%s VALUES `, t.Name, strings.Join(colList, ","), ov))
|
||||
args := make([]any, 0, len(chunk)*len(t.Columns))
|
||||
param := 1
|
||||
for ri, row := range chunk {
|
||||
if ri > 0 {
|
||||
sb.WriteString(",")
|
||||
}
|
||||
sb.WriteString("(")
|
||||
for ci, val := range row {
|
||||
if ci > 0 {
|
||||
sb.WriteString(",")
|
||||
}
|
||||
if val == nil {
|
||||
sb.WriteString("NULL")
|
||||
} else {
|
||||
if casts[ci] != "" {
|
||||
sb.WriteString(fmt.Sprintf("$%d::%s", param, casts[ci]))
|
||||
} else {
|
||||
sb.WriteString(fmt.Sprintf("$%d", param))
|
||||
}
|
||||
args = append(args, val)
|
||||
param++
|
||||
}
|
||||
}
|
||||
sb.WriteString(")")
|
||||
}
|
||||
if _, err := tx.Exec(sb.String(), args...); err != nil {
|
||||
for rows.Next() {
|
||||
var name, def string
|
||||
if err := rows.Scan(&name, &def); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Fprintf(w, "ALTER TABLE ONLY %s ADD CONSTRAINT %s %s;\n", qname(table), qi(name), def)
|
||||
}
|
||||
fmt.Fprintln(w)
|
||||
return rows.Err()
|
||||
}
|
||||
|
||||
// 重置 identity 序列,避免后续自增插入与已导入的最大 ID 冲突。
|
||||
for _, c := range t.Columns {
|
||||
if c.IsIdentity {
|
||||
seq := fmt.Sprintf("pg_get_serial_sequence('%s','%s')", t.Name, c.Name)
|
||||
if _, err := tx.Exec(fmt.Sprintf(
|
||||
`SELECT setval(%s, COALESCE((SELECT MAX("%s") FROM "%s"), 1))`, seq, c.Name, t.Name)); err != nil {
|
||||
log.Printf("! 重置序列 %s.%s 失败: %v", t.Name, c.Name, err)
|
||||
}
|
||||
// writeIndexes 导出不承载约束的索引(含 pgvector 的 HNSW)。
|
||||
//
|
||||
// 排除承载约束的索引:主键索引已由 writeConstraints 以 ADD CONSTRAINT 的形式重建,
|
||||
// 再 CREATE INDEX 一次会重复。
|
||||
func writeIndexes(w *bufio.Writer, db *sql.DB, table string) error {
|
||||
q := fmt.Sprintf(`
|
||||
SELECT i.relname, pg_get_indexdef(i.oid)
|
||||
FROM pg_index x
|
||||
JOIN pg_class i ON i.oid = x.indexrelid
|
||||
WHERE x.indrelid = %s::regclass
|
||||
AND NOT EXISTS (SELECT 1 FROM pg_constraint con WHERE con.conindid = x.indexrelid)
|
||||
ORDER BY i.relname`, ql(qname(table)))
|
||||
|
||||
rows, err := db.Query(q)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
for rows.Next() {
|
||||
var name, def string
|
||||
if err := rows.Scan(&name, &def); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Fprintf(w, "%s;\n", injectIfNotExists(def))
|
||||
}
|
||||
fmt.Fprintln(w)
|
||||
return rows.Err()
|
||||
}
|
||||
|
||||
// writeSequenceResets 为每个自增列把序列推到 MAX(col)+1。
|
||||
//
|
||||
// COPY 写入的是显式 id,不会推进序列;不重置的话,后续 INSERT 会从 1 开始并与存量主键冲突。
|
||||
// setval 用 (值, false) 形式:false 表示「这个值还没被取走」,故下一次 nextval 正好是 max+1;
|
||||
// 空表时为 1,即从 1 开始。
|
||||
func writeSequenceResets(w *bufio.Writer, table string, cols []column) {
|
||||
for _, c := range cols {
|
||||
if !c.AutoIncrement {
|
||||
continue
|
||||
}
|
||||
fmt.Fprintf(w,
|
||||
"SELECT pg_catalog.setval(pg_get_serial_sequence(%s, %s), COALESCE((SELECT MAX(%s) FROM %s), 0) + 1, false);\n",
|
||||
ql(qname(table)), ql(c.Name), qi(c.Name), qname(table))
|
||||
}
|
||||
}
|
||||
|
||||
// isSerial 判断是否为 serial 列(默认值 nextval,且类型为整型)。
|
||||
func isSerial(c column) bool {
|
||||
if !strings.HasPrefix(c.Default, "nextval(") {
|
||||
return false
|
||||
}
|
||||
switch c.Type {
|
||||
case "bigint", "integer", "smallint":
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// escapeCopy 转义 COPY 文本格式里的特殊字符。
|
||||
// 先处理反斜杠,保证字面量 `\N` 被写成 `\\N` 而不会变成 NULL。
|
||||
func escapeCopy(b []byte) string {
|
||||
var sb strings.Builder
|
||||
sb.Grow(len(b) + 8)
|
||||
for _, c := range b {
|
||||
switch c {
|
||||
case '\\':
|
||||
sb.WriteString(`\\`)
|
||||
case '\t':
|
||||
sb.WriteString(`\t`)
|
||||
case '\n':
|
||||
sb.WriteString(`\n`)
|
||||
case '\r':
|
||||
sb.WriteString(`\r`)
|
||||
default:
|
||||
sb.WriteByte(c)
|
||||
}
|
||||
}
|
||||
return tx.Commit()
|
||||
return sb.String()
|
||||
}
|
||||
|
||||
// injectIfNotExists 给 pg_get_indexdef 的输出补上 IF NOT EXISTS。
|
||||
// pg_get_indexdef 不会输出它(那是 pg_dump 的清理语义),补上可让索引段落可重复执行。
|
||||
func injectIfNotExists(def string) string {
|
||||
if rest, ok := strings.CutPrefix(def, "CREATE UNIQUE INDEX "); ok {
|
||||
return "CREATE UNIQUE INDEX IF NOT EXISTS " + rest
|
||||
}
|
||||
if rest, ok := strings.CutPrefix(def, "CREATE INDEX "); ok {
|
||||
return "CREATE INDEX IF NOT EXISTS " + rest
|
||||
}
|
||||
return def
|
||||
}
|
||||
|
||||
// qname 返回 schema 限定且带引号的表名:public."brands"。
|
||||
func qname(table string) string { return "public." + qi(table) }
|
||||
|
||||
// qi 安全包裹 SQL 标识符(表名 / 列名 / 索引名)。
|
||||
func qi(s string) string { return `"` + strings.ReplaceAll(s, `"`, `""`) + `"` }
|
||||
|
||||
// ql 安全包裹 SQL 字符串字面量。
|
||||
func ql(s string) string { return "'" + strings.ReplaceAll(s, "'", "''") + "'" }
|
||||
|
||||
@ -67,14 +67,8 @@ func main() {
|
||||
log.Printf("! 关闭数据库连接失败: %v", err)
|
||||
}
|
||||
}()
|
||||
// 2.5 幂等自动迁移表结构(替代原手写旧迁移脚本)
|
||||
if err := database.AutoMigrate(db); err != nil {
|
||||
log.Fatalf("✗ 自动迁移失败: %v", err)
|
||||
}
|
||||
// 2.6 确保去重所需的 HNSW 索引与语义 embedding 表(幂等,随启动重复执行)。
|
||||
if err := database.EnsureDedupSchema(db); err != nil {
|
||||
log.Fatalf("✗ 去重 schema 初始化失败: %v", err)
|
||||
}
|
||||
// 2.5 结构不再随启动自动迁移:全库「结构 + 索引 + 数据」统一由 cmd/dbtool 导出的
|
||||
// 纯 SQL 维护(dbtool dump → psql -f)。这里只探活,不执行任何 DDL。
|
||||
log.Println("✓ PostgreSQL 已连接:", cfg.Database.Addr())
|
||||
|
||||
// 3. 确保上传目录存在(静态文件服务的根目录)
|
||||
|
||||
@ -79,7 +79,7 @@ client_sign:
|
||||
# 爬虫入库管线:暴露 :8092 /admin/internal/ingest 接收爬虫 HMAC 签名上报。
|
||||
# 设计:服务端到服务端(非前端 JS),密钥绝不下发前端、只存爬虫配置与 INGEST_SECRET,安全高。
|
||||
# - secret 必须与爬虫端 config.yml 的 ingest.secret 完全一致;env INGEST_SECRET 注入随机长串(生产)。
|
||||
# - 建表由 GORM AutoMigrate 自动完成(ingest_jobs / ingest_nonces 等),无需手工 SQL 迁移。
|
||||
# - 建表不再由服务启动自动完成:表结构随全库由 cmd/dbtool 导出的 SQL 维护(dbtool dump → psql -f)。
|
||||
# - 启用后 worker 随主进程拉起(cmd/server 内 goroutine),异步下载图、补 season_code、写正式表。
|
||||
ingest:
|
||||
secret: "dev-ingest-secret-2026" # env: INGEST_SECRET(须与爬虫 config.yml 的 ingest.secret 一致)
|
||||
|
||||
@ -4,7 +4,11 @@
|
||||
// 依赖关系清晰可见,也让 repository 可以在测试中替换为独立的数据库实例。
|
||||
//
|
||||
// 数据库已迁移至 PostgreSQL(见 scripts/pgvector 的本地 Docker 环境)。
|
||||
// 表结构由 GORM AutoMigrate 在启动时幂等创建,不再依赖手写的旧迁移脚本。
|
||||
//
|
||||
// 表结构不再由服务启动时自动迁移:全库「结构 + 索引 + 数据」统一由 cmd/dbtool 导出的
|
||||
// 纯 SQL 维护(dbtool dump → psql -f)。原 database.AutoMigrate / EnsureDedupSchema
|
||||
// 已随此决策移除 —— 本项目结构定义现在只有一个来源,就是那份 dump 脚本。
|
||||
// 因此本包只负责「连接」这一件事,不执行任何 DDL。
|
||||
package database
|
||||
|
||||
import (
|
||||
@ -12,7 +16,6 @@ import (
|
||||
"time"
|
||||
|
||||
"fashionapi/internal/config"
|
||||
"fashionapi/internal/model"
|
||||
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
@ -50,183 +53,6 @@ func New(cfg config.DatabaseConfig) (*gorm.DB, error) {
|
||||
return db, nil
|
||||
}
|
||||
|
||||
// AutoMigrate 幂等创建 / 更新全部表结构。
|
||||
//
|
||||
// 表结构由 GORM 模型定义托管(替代原 scripts/sql、db/migrations 下的手写旧迁移脚本)。
|
||||
// 仅在缺少列/索引时增量变更,已存在的表不会被重建;对存量大表(如 brand_runway_images)
|
||||
// 的加列操作可能短暂加锁,属一次性开销。
|
||||
func AutoMigrate(db *gorm.DB) error {
|
||||
// 确保 pgvector 扩展存在(dHash 近重复 HNSW 索引与语义 embedding 向量列均依赖它)。
|
||||
if err := db.Exec("CREATE EXTENSION IF NOT EXISTS vector").Error; err != nil {
|
||||
return fmt.Errorf("create extension vector: %w", err)
|
||||
}
|
||||
// 兼容旧库:phash 早期以 bigint 存储 64-bit dHash 指纹(uint64),现模型改为 vector(64)。
|
||||
// pgvector 不支持 bigint→vector 直接转换,需逐位拆成 0/1 向量(位 i → 第 i 维,与 ToVectorBits 一致,
|
||||
// 保证旧数据与新写入的向量在同一语义空间可比)。已为 vector/text 类型的列跳过,幂等可重复执行。
|
||||
const migratePhashSQL = `
|
||||
DO $$
|
||||
DECLARE
|
||||
t text;
|
||||
dt text;
|
||||
BEGIN
|
||||
FOR t IN
|
||||
SELECT tablename FROM pg_tables
|
||||
WHERE schemaname = 'public'
|
||||
AND tablename IN ('brand_runway_images','brand_runway_draft_images','street_snap_images','street_snap_draft_images')
|
||||
LOOP
|
||||
SELECT data_type INTO dt
|
||||
FROM information_schema.columns
|
||||
WHERE table_schema = 'public' AND table_name = t AND column_name = 'phash';
|
||||
|
||||
IF dt = 'bigint' THEN
|
||||
-- bigint 即 64-bit dHash 指纹;pgvector 不支持 bigint->vector 直接转换,且 ALTER ... USING 的
|
||||
-- transform 表达式不允许子查询(0A000)。改用纯表达式:bigint::bit(64) 取 64 位模式(MSB 在左),
|
||||
-- reverse 成 LSB 在左(与 phash.ToVectorBits 维度顺序一致:第 i 维 = bit i),插逗号后 text::vector(64)。
|
||||
-- NULL 自然透传。
|
||||
EXECUTE format($e$
|
||||
ALTER TABLE %I ALTER COLUMN phash TYPE vector(64)
|
||||
USING (
|
||||
('[' ||
|
||||
rtrim(regexp_replace(reverse(phash::bit(64)::text), '(.)', '\1,', 'g'), ',') ||
|
||||
']')::vector(64)
|
||||
)
|
||||
$e$, t);
|
||||
RAISE NOTICE 'migrated phash (bigint -> vector) on %', t;
|
||||
ELSIF dt IN ('character varying','text') THEN
|
||||
EXECUTE format($e$
|
||||
ALTER TABLE %I ALTER COLUMN phash TYPE vector(64) USING phash::vector(64)
|
||||
$e$, t);
|
||||
RAISE NOTICE 'migrated phash (text -> vector) on %', t;
|
||||
END IF;
|
||||
END LOOP;
|
||||
END $$;`
|
||||
if err := db.Exec(migratePhashSQL).Error; err != nil {
|
||||
return fmt.Errorf("migrate phash column: %w", err)
|
||||
}
|
||||
// 回填历史 NULL:迁移 012 给 4 张 image 表的 is_duplicate/dup_of 声明了 not null,
|
||||
// 但存量行可能是 NULL;GORM 直接 SET NOT NULL 会失败(23502)。先按默认值回填再交给 GORM(幂等)。
|
||||
// 注意:各表 schema 演化路径不同,可能缺其中某列——故按列存在性动态拼 SQL,只回填真实存在的列;
|
||||
// 缺失的列交由随后的 db.AutoMigrate 新建(带默认值,不会触发 NULL 问题)。
|
||||
const backfillDedupSQL = `
|
||||
DO $$
|
||||
DECLARE
|
||||
t text;
|
||||
setc text := '';
|
||||
wherec text := '';
|
||||
has_dup boolean;
|
||||
has_dupof boolean;
|
||||
BEGIN
|
||||
FOR t IN
|
||||
SELECT unnest(ARRAY[
|
||||
'brand_runway_images','brand_runway_draft_images',
|
||||
'street_snap_images','street_snap_draft_images'])
|
||||
LOOP
|
||||
SELECT
|
||||
EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema='public' AND table_name=t AND column_name='is_duplicate'),
|
||||
EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema='public' AND table_name=t AND column_name='dup_of')
|
||||
INTO has_dup, has_dupof;
|
||||
|
||||
setc := '';
|
||||
wherec := '';
|
||||
IF has_dup THEN
|
||||
setc := setc || 'is_duplicate = COALESCE(is_duplicate, 0)';
|
||||
wherec := wherec || 'is_duplicate IS NULL';
|
||||
END IF;
|
||||
IF has_dupof THEN
|
||||
IF setc <> '' THEN setc := setc || ', '; wherec := wherec || ' OR '; END IF;
|
||||
setc := setc || 'dup_of = COALESCE(dup_of, '''')';
|
||||
wherec := wherec || 'dup_of IS NULL';
|
||||
END IF;
|
||||
|
||||
IF setc <> '' THEN
|
||||
EXECUTE format('UPDATE %I SET %s WHERE %s', t, setc, wherec);
|
||||
END IF;
|
||||
END LOOP;
|
||||
END $$;`
|
||||
if err := db.Exec(backfillDedupSQL).Error; err != nil {
|
||||
return fmt.Errorf("backfill dedup columns: %w", err)
|
||||
}
|
||||
return db.AutoMigrate(
|
||||
&model.Brand{},
|
||||
&model.User{},
|
||||
&model.RefreshToken{},
|
||||
&model.IngestNonce{},
|
||||
&model.IngestJob{},
|
||||
&model.StreetSnap{},
|
||||
&model.StreetSnapImage{},
|
||||
&model.StreetSnapDraft{},
|
||||
&model.StreetSnapDraftImage{},
|
||||
&model.Favorite{},
|
||||
&model.History{},
|
||||
&model.BrandRunway{},
|
||||
&model.BrandRunwayImage{},
|
||||
&model.BrandRunwayDraft{},
|
||||
&model.BrandRunwayDraftImage{},
|
||||
)
|
||||
}
|
||||
|
||||
// EnsureDedupSchema 在 GORM AutoMigrate 之外,补充去重所需的索引与语义 embedding 表。
|
||||
// 全部幂等(IF NOT EXISTS / DROP IF EXISTS),可随服务启动重复执行。
|
||||
//
|
||||
// - phash 的 HNSW 索引(vector_l2_ops,L2²==汉明距离,替代原先的全表扫描);
|
||||
// - 幂等清理 content_sha1 精确去重废弃后遗留的旧索引(uq_*_sha1 及其 HNSW);
|
||||
// - image_embeddings 表 + cosine HNSW 索引:语义 embedding(CLIP/DINOv2)落库位,
|
||||
// 待后续推理接入填充,当前不写入。
|
||||
func EnsureDedupSchema(db *gorm.DB) error {
|
||||
// 统一去重方案:废弃早期独立 image_dhash 表(与 phash 直接挂各 image 表冲突)。
|
||||
// 该表从未被 Go 写入,启动即幂等清理遗留表,避免与现行方案并存造成混淆。
|
||||
if err := db.Exec("DROP TABLE IF EXISTS image_dhash CASCADE").Error; err != nil {
|
||||
return fmt.Errorf("drop legacy image_dhash: %w", err)
|
||||
}
|
||||
// content_sha1 精确去重已废弃:幂等清理其遗留的唯一索引 uq_*_sha1 与旧 HNSW 索引 uq_*_sha1_phash_hnsw。
|
||||
// 新 HNSW 索引统一改名为 <table>_phash_hnsw。
|
||||
legacyIdx := []string{
|
||||
"uq_br_imgs_sha1", "uq_br_draft_imgs_sha1",
|
||||
"uq_ss_imgs_sha1", "uq_ss_draft_imgs_sha1",
|
||||
}
|
||||
for _, idx := range legacyIdx {
|
||||
for _, suffix := range []string{"", "_phash_hnsw"} {
|
||||
if err := db.Exec(fmt.Sprintf("DROP INDEX IF EXISTS %s%s", idx, suffix)).Error; err != nil {
|
||||
return fmt.Errorf("drop legacy index %s%s: %w", idx, suffix, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
tables := []string{
|
||||
"brand_runway_images",
|
||||
"brand_runway_draft_images",
|
||||
"street_snap_images",
|
||||
"street_snap_draft_images",
|
||||
}
|
||||
for _, name := range tables {
|
||||
// 近重复检索:dHash 向量(vector(64) 的 0/1)汉明 HNSW 索引。
|
||||
// 用 vector_l2_ops:两 {0,1}^64 向量的 L2² 恰等于汉明距离,故 L2 距离即汉明语义;
|
||||
// 且 vector_l2_ops 在所有 pgvector 版本均可建 HNSW(bit_hamming_ops 在旧版不可用于 HNSW)。
|
||||
if err := db.Exec(fmt.Sprintf(
|
||||
"CREATE INDEX IF NOT EXISTS %s_phash_hnsw ON %s USING hnsw (phash vector_l2_ops)",
|
||||
name, name,
|
||||
)).Error; err != nil {
|
||||
return fmt.Errorf("create hnsw index %s: %w", name, err)
|
||||
}
|
||||
}
|
||||
// 语义 embedding 表(CLIP/DINOv2 等 512 维向量):待后续推理接入填充;当前不写入。
|
||||
if err := db.Exec(`
|
||||
CREATE TABLE IF NOT EXISTS image_embeddings (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
image_id INTEGER NOT NULL,
|
||||
kind SMALLINT NOT NULL,
|
||||
embedding vector(512),
|
||||
created_at INTEGER NOT NULL DEFAULT 0
|
||||
)`).Error; err != nil {
|
||||
return fmt.Errorf("create image_embeddings: %w", err)
|
||||
}
|
||||
if err := db.Exec(
|
||||
"CREATE INDEX IF NOT EXISTS image_embeddings_vec_idx ON image_embeddings USING hnsw (embedding vector_cosine_ops)",
|
||||
).Error; err != nil {
|
||||
return fmt.Errorf("create embeddings hnsw: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close 关闭数据库连接池。
|
||||
func Close(db *gorm.DB) error {
|
||||
if db == nil {
|
||||
|
||||
@ -1,8 +1,8 @@
|
||||
// Package phash 计算图片的感知哈希(dHash,Difference Hash)。
|
||||
//
|
||||
// dHash 是 64-bit 轻量指纹,用于判断「同一张图被转码 / 加水印 / 改尺寸」这类近重复。
|
||||
// 去重检索不在这里做——交给 PostgreSQL + pgvector 的 HNSW 索引(见 database 包的
|
||||
// ensureDedupSchema),本包只负责产出指纹,复杂度极低、零 CGO、零 ML 依赖。
|
||||
// 去重检索不在这里做——交给 PostgreSQL + pgvector 的 HNSW 索引(索引随全库结构一起,
|
||||
// 由 cmd/dbtool 导出的 SQL 维护),本包只负责产出指纹,复杂度极低、零 CGO、零 ML 依赖。
|
||||
//
|
||||
// 与之前被移除的全表扫描实现不同:现在指纹只写库、检索交给索引,O(1) 桶内比对即可,
|
||||
// 图片量级再大也扛得住。
|
||||
|
||||
@ -23,7 +23,14 @@ import (
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// testDB 连真实库并完成 AutoMigrate + 去重 schema 初始化。
|
||||
// testDB 连真实库并校验去重所需的结构已就绪。
|
||||
//
|
||||
// 结构不再由 AutoMigrate 自动创建(已移除),须先自行准备好库,例如:
|
||||
//
|
||||
// ./bin/dbtool dump -out db_dump.sql
|
||||
// psql -U fashion -d fashion -v ON_ERROR_STOP=1 -f db_dump.sql
|
||||
//
|
||||
// 这里只做存在性校验并给出明确提示,避免后续查询报出难懂的错。
|
||||
func testDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
cfg, err := config.Load("")
|
||||
@ -34,13 +41,11 @@ func testDB(t *testing.T) *gorm.DB {
|
||||
if err != nil {
|
||||
t.Fatalf("连接数据库失败: %v", err)
|
||||
}
|
||||
if err := database.AutoMigrate(db); err != nil {
|
||||
t.Fatalf("AutoMigrate 失败: %v", err)
|
||||
}
|
||||
if err := database.EnsureDedupSchema(db); err != nil {
|
||||
t.Fatalf("去重 schema 失败: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = database.Close(db) })
|
||||
|
||||
if err := db.Exec("SELECT 1 FROM brand_runway_images LIMIT 1").Error; err != nil {
|
||||
t.Fatalf("brand_runway_images 不可用(库结构未就绪?先 dbtool dump 再 psql -f 灌库): %v", err)
|
||||
}
|
||||
return db
|
||||
}
|
||||
|
||||
|
||||
@ -3,9 +3,10 @@
|
||||
-- pgvector 不支持 bigint -> vector 直接转换,需逐位拆成 0/1 向量(位 i -> 第 i 维,
|
||||
-- 与 Go 侧 phash.ToVectorBits 完全一致,保证旧数据与新写入向量在同一语义空间可比)。
|
||||
--
|
||||
-- 用法(任选其一):
|
||||
-- 1) 后端启动时由 AutoMigrate 自动执行(已内置,幂等,无需手动);
|
||||
-- 2) 手动执行: psql "postgres://fashion:fashion_dev_2026@localhost:5432/fashion" -f 03-migrate-phash-bigint-to-vector.sql
|
||||
-- 用法(仅升级旧库时需要):
|
||||
-- 本迁移原先由后端启动时的 AutoMigrate(migratePhashSQL)自动执行;AutoMigrate 已移除,
|
||||
-- 故升级老库改回手动执行:
|
||||
-- psql "postgres://fashion:fashion_dev_2026@localhost:5432/fashion" -f 03-migrate-phash-bigint-to-vector.sql
|
||||
--
|
||||
-- 已为 vector/text 类型的列自动跳过;可重复执行。
|
||||
|
||||
|
||||
@ -1,7 +1,7 @@
|
||||
// Command seed_users 预置内部账号(运营/编辑用,不开放公开注册)。
|
||||
//
|
||||
// 项目表结构由 GORM AutoMigrate 托管(服务启动时创建);
|
||||
// 本脚本仅在独立运行、尚未跑过服务时兜底建表,再 upsert 一个内部账号。
|
||||
// 项目表结构由 cmd/dbtool 导出的 SQL 维护(dbtool dump → psql -f),服务启动时不再自动建表;
|
||||
// 本脚本为「尚未灌库时也能独立跑」保留一份 users 兜底建表,再 upsert 一个内部账号。
|
||||
//
|
||||
// 用法:
|
||||
// go run ./scripts/seed_users
|
||||
@ -43,7 +43,7 @@ func main() {
|
||||
}
|
||||
}()
|
||||
|
||||
// 1. 兜底建表(与 User 模型保持一致;正式结构由服务启动时 AutoMigrate 创建,此处仅作独立运行兜底)
|
||||
// 1. 兜底建表(与 User 模型保持一致;正式结构由 dbtool dump 的 SQL 创建,此处仅作独立运行兜底)
|
||||
createSQL := `CREATE TABLE IF NOT EXISTS users (
|
||||
id bigserial PRIMARY KEY,
|
||||
created_at bigint NOT NULL DEFAULT 0,
|
||||
|
||||
Reference in New Issue
Block a user