Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ mysql2pg
config.yml
*.log
*.exe
ddl/
# gitlab ci
.cache
.golangci.yml
Expand Down
4 changes: 3 additions & 1 deletion config.example.yml
Original file line number Diff line number Diff line change
Expand Up @@ -77,4 +77,6 @@ run:
enable_file_logging: true # 是否启用文件日志
log_file_path: ./conversion.log # 日志文件保存路径
show_console_logs: true # 是否在控制台显示日志信息
show_log_in_console: false # 是否在控制台显示Log日志输出
show_log_in_console: false # 是否在控制台显示Log日志输出
enable_ddl_output: false # 是否把转换后的 PostgreSQL DDL(表/分区/索引/视图/函数等)导出到文件。多线程时可能影响性能
ddl_output_file_path: ./ddl/pgsql_ddl.sql # DDL 导出文件路径(相对路径以运行目录为准,目录不存在会自动创建)
4 changes: 4 additions & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,10 @@ type RunConfig struct {
LogFilePath string `mapstructure:"log_file_path"`
ShowConsoleLogs bool `mapstructure:"show_console_logs"`
ShowLogInConsole bool `mapstructure:"show_log_in_console"`
// EnableDDLOutput 是否把转换后的 PostgreSQL DDL(表/索引/视图/函数等)导出到文件
EnableDDLOutput bool `mapstructure:"enable_ddl_output"`
// DDLOutputFilePath 转换后 DDL 的导出文件路径(相对路径以运行目录为准)
DDLOutputFilePath string `mapstructure:"ddl_output_file_path"`
}

// LoadConfig 加载配置文件
Expand Down
97 changes: 86 additions & 11 deletions internal/converter/postgres/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,10 @@ package postgres

import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"sync"
"sync/atomic"
Expand Down Expand Up @@ -47,11 +49,13 @@ func (c *ConversionContext) shouldUseJsonArrayInsert() bool {

// Manager 转换管理器
type Manager struct {
mysqlConn *mysql.Connection
postgresConn *postgres.Connection
config *config.Config
errorLogFile *os.File
logFile *os.File
mysqlConn *mysql.Connection
postgresConn *postgres.Connection
config *config.Config
errorLogFile *os.File
logFile *os.File
// 转换后 DDL 导出文件(run.enable_ddl_output=true 时打开)
ddlOutputFile *os.File
totalTasks int
completedTasks atomic.Int64
mutex sync.Mutex
Expand Down Expand Up @@ -145,21 +149,66 @@ func (m *Manager) context() context.Context {
return m.ctx
}

// openDDLExportFile 打开 run.ddl_output_file_path 指定的导出文件。
// 每次运行覆盖重建;文件所在目录不存在时自动创建
func (m *Manager) openDDLExportFile() error {
if m.config == nil || !m.config.Run.EnableDDLOutput || m.config.Run.DDLOutputFilePath == "" {
return nil
}
path := m.config.Run.DDLOutputFilePath
if dir := filepath.Dir(path); dir != "" && dir != "." {
if err := os.MkdirAll(dir, 0755); err != nil {
return fmt.Errorf("创建 DDL 导出目录失败: %w", err)
}
}
f, err := os.OpenFile(path, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0644)
if err != nil {
return fmt.Errorf("打开 DDL 导出文件失败: %w", err)
}
m.ddlOutputFile = f
return nil
}

// exportDDLToFile 把转换后生成的 DDL 写入导出文件(若已启用)。
// 各转换阶段以多个 goroutine 并发执行,统一经 m.mutex 串行化写入
func (m *Manager) exportDDLToFile(category, objectName, ddl string) {
if m.ddlOutputFile == nil {
return
}
ddl = strings.TrimSpace(ddl)
if ddl == "" {
return
}
ddl = strings.TrimSuffix(ddl, ";")
m.mutex.Lock()
defer m.mutex.Unlock()
fmt.Fprintf(m.ddlOutputFile, "-- ===== [%s] %s =====\n%s;\n\n", category, objectName, ddl)
}

// Close 关闭转换管理器
// 关闭打开的日志文件
// 关闭打开的日志文件、DDL 导出文件和错误日志文件;
// 任一文件关闭失败都会保留并返回,互不覆盖
func (m *Manager) Close() error {
var err error
var errs []error
if m.logFile != nil {
if closeErr := m.logFile.Close(); closeErr != nil {
err = closeErr
errs = append(errs, fmt.Errorf("关闭日志文件失败: %w", closeErr))
}
}

if closeErr := m.errorLogFile.Close(); closeErr != nil && err == nil {
err = closeErr
if m.ddlOutputFile != nil {
if closeErr := m.ddlOutputFile.Close(); closeErr != nil {
errs = append(errs, fmt.Errorf("关闭 DDL 导出文件失败: %w", closeErr))
}
}

if m.errorLogFile != nil {
if closeErr := m.errorLogFile.Close(); closeErr != nil {
errs = append(errs, fmt.Errorf("关闭错误日志文件失败: %w", closeErr))
}
}

return err
return errors.Join(errs...)
}

// Run 执行完整的转换流程
Expand All @@ -170,6 +219,14 @@ func (m *Manager) Run() error {
return fmt.Errorf("转换已取消: %w", err)
}

// run.enable_ddl_output=true 时打开导出文件
if err := m.openDDLExportFile(); err != nil {
return err
}
if m.ddlOutputFile != nil {
m.Log("转换后的 PostgreSQL DDL 将导出到文件: %s", m.config.Run.DDLOutputFilePath)
}

m.Log("表MySQL 的DDL、数据、view、索引、函数、用户和权限的转换到 PostgreSQL ...")

// 检查是否启用了表列表功能
Expand Down Expand Up @@ -933,6 +990,9 @@ func (m *Manager) convertViews(views []mysql.ViewInfo, semaphore chan struct{})
return err
}

// 导出转换后的视图 DDL
m.exportDDLToFile("视图", view.ViewName, pgViewDDL)

// 执行创建视图的SQL语句
if err := m.postgresConn.ExecuteDDL(pgViewDDL, view.ViewDefinition); err != nil {
errMsg := fmt.Sprintf("创建表视图 %s 失败: %v", view.ViewName, err)
Expand Down Expand Up @@ -1005,6 +1065,15 @@ func (m *Manager) convertTables(tables []mysql.TableInfo, semaphore chan struct{
m.RecordConversionWarning("表结构", table.Name, w)
}

// 导出转换后的表结构 DDL(含分区子表与 CHECK 约束)
m.exportDDLToFile("表结构", table.Name, pgResult.DDL)
for _, partitionDDL := range pgResult.PartitionDDLs {
m.exportDDLToFile("表结构", table.Name+" 分区子表", partitionDDL)
}
for _, checkDDL := range pgResult.CheckConstraints {
m.exportDDLToFile("表结构", table.Name+" CHECK 约束", checkDDL)
}

// 先检查表是否存在
tableExists, err := m.postgresConn.TableExists(table.Name)
if err != nil {
Expand Down Expand Up @@ -1235,6 +1304,9 @@ func (m *Manager) convertFunctions(functions []mysql.FunctionInfo, semaphore cha
return err
}

// 导出转换后的函数/存储过程 DDL
m.exportDDLToFile("函数/存储过程", function.Name, pgDDL)

if err := m.postgresConn.ExecuteDDL(pgDDL, function.DDL); err != nil {
errMsg := fmt.Sprintf("执行函数 %s DDL失败: %v", function.Name, err)
m.logError(errMsg)
Expand Down Expand Up @@ -1335,6 +1407,9 @@ func (m *Manager) convertIndexes(indexes []mysql.IndexInfo, semaphore chan struc
continue
}

// 导出转换后的索引 DDL
m.exportDDLToFile("索引", fmt.Sprintf("%s.%s", index.Table, lowercaseIndexName), pgDDL)

// 执行DDL语句
if err := m.postgresConn.ExecuteDDL(pgDDL); err != nil {
// 检查是否是索引已存在的错误
Expand Down
76 changes: 76 additions & 0 deletions internal/converter/postgres/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,9 @@ import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"sync"
"testing"

Expand Down Expand Up @@ -304,3 +307,76 @@ func TestConvertFunctions_SkippedFunctionsCountedOnce(t *testing.T) {
t.Errorf("completedTasks = %d, want %d(被跳过的函数只应计数一次,不应双计数)", got, want)
}
}

// TestDDLExportToFile run.enable_ddl_output=true:目录自动创建、语句按分类写入、
// 本身带分号的 DDL 不会产生双分号
func TestDDLExportToFile(t *testing.T) {
outPath := filepath.Join(t.TempDir(), "ddl", "pgsql_ddl.sql")
m := &Manager{config: &config.Config{Run: config.RunConfig{
EnableDDLOutput: true,
DDLOutputFilePath: outPath,
}}}

if err := m.openDDLExportFile(); err != nil {
t.Fatalf("openDDLExportFile() 失败: %v", err)
}
m.exportDDLToFile("表结构", "sys_user", `CREATE TABLE "sys_user" (id BIGSERIAL PRIMARY KEY)`)
m.exportDDLToFile("表结构", "orders 分区子表", `CREATE TABLE "orders_p1" PARTITION OF "orders" FOR VALUES FROM (0) TO (100);`)
m.exportDDLToFile("索引", "sys_user.idx_username", `CREATE INDEX "idx_username" ON "sys_user" ("username");`)
m.exportDDLToFile("视图", "v_user", `CREATE OR REPLACE VIEW "v_user" AS SELECT 1`)

if err := m.ddlOutputFile.Close(); err != nil {
t.Fatalf("关闭 DDL 导出文件失败: %v", err)
}

data, err := os.ReadFile(outPath)
if err != nil {
t.Fatalf("读取 DDL 导出文件失败: %v", err)
}
content := string(data)
for _, want := range []string{
"-- ===== [表结构] sys_user =====",
`CREATE TABLE "sys_user" (id BIGSERIAL PRIMARY KEY);`,
"-- ===== [表结构] orders 分区子表 =====",
`CREATE TABLE "orders_p1" PARTITION OF "orders" FOR VALUES FROM (0) TO (100);`,
"-- ===== [索引] sys_user.idx_username =====",
`CREATE INDEX "idx_username" ON "sys_user" ("username");`,
"-- ===== [视图] v_user =====",
`CREATE OR REPLACE VIEW "v_user" AS SELECT 1;`,
} {
if !strings.Contains(content, want) {
t.Errorf("DDL 导出文件缺少内容 %q,实际内容:\n%s", want, content)
}
}
if strings.Contains(content, ";;") {
t.Errorf("DDL 导出文件出现双分号,实际内容:\n%s", content)
}
}

// TestDDLExportDisabled run.enable_ddl_output=false 时不应创建导出文件
func TestDDLExportDisabled(t *testing.T) {
outPath := filepath.Join(t.TempDir(), "should_not_exist.sql")
m := &Manager{config: &config.Config{Run: config.RunConfig{
EnableDDLOutput: false,
DDLOutputFilePath: outPath,
}}}

if err := m.openDDLExportFile(); err != nil {
t.Fatalf("openDDLExportFile() 失败: %v", err)
}
if m.ddlOutputFile != nil {
t.Fatal("enable_ddl_output=false 时不应打开导出文件")
}
if _, err := os.Stat(outPath); !os.IsNotExist(err) {
t.Fatalf("enable_ddl_output=false 时不应创建导出文件 %s", outPath)
}
}

// TestManagerCloseNilFiles Close 在文件句柄均为 nil(如测试或部分初始化路径)
// 时不应 panic,且应返回 nil
func TestManagerCloseNilFiles(t *testing.T) {
m := &Manager{}
if err := m.Close(); err != nil {
t.Fatalf("Close() 在无文件句柄时应返回 nil,实际: %v", err)
}
}
4 changes: 3 additions & 1 deletion internal/postgres/connection.go
Original file line number Diff line number Diff line change
Expand Up @@ -990,14 +990,16 @@ func (c *Connection) BatchInsertDataWithCompositeKeys(ctx context.Context, tx pg
if primaryKey == "" {
continue
}
found := false
for i, col := range columns {
if strings.EqualFold(col, primaryKey) {
resolvedPrimaryKeys = append(resolvedPrimaryKeys, copyColumns[i])
found = true
break
}
}
// fallback: 如果没找到,使用转换后的主键名
if len(resolvedPrimaryKeys) < len(primaryKeys) {
if !found {
resolvedPK := primaryKey
if lowercaseColumns {
resolvedPK = strings.ToLower(primaryKey)
Expand Down