From 20f039f6e05b30f78fcc317ac2e2a64afbd68430 Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 14:37:37 +0800 Subject: [PATCH 01/10] =?UTF-8?q?fix:=20MySQL=20=E5=85=83=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E6=9F=A5=E8=AF=A2=E8=A1=A5=20rows.Err()=20=E6=A3=80?= =?UTF-8?q?=E6=9F=A5=EF=BC=8C=E9=98=B2=E6=AD=A2=E7=BB=93=E6=9E=9C=E9=9B=86?= =?UTF-8?q?=E6=88=AA=E6=96=AD=E8=A2=AB=E5=BD=93=E6=88=90=E5=AE=8C=E6=95=B4?= =?UTF-8?q?=E8=AF=BB=E5=8F=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit GetTableColumns / GetTableColumnsWithTypes / GetTablePrimaryKeys 在 rows.Next() 循环后直接返回 nil error。rows.Next() 遇网络中断、包解码 错误或服务端提前关闭结果集时返回 false,错误只能从 rows.Err() 取到, 漏检会把"结果集被截断"当成"正常读完"。 两处后果都是静默的,且现有校验抓不到: - 列清单截断 → sync_data.go 用它构造 SELECT 与 CopyFrom 的 copyColumns → 少迁几列、零错误零告警,行数校验只比 COUNT(*) 照样通过 - 复合主键 (a,b,c) 截断成 (a,b) → len==2 通过"没有主键"空检查 → keyset 的 WHERE (a,b) > (?,?) 游标不再唯一 → 跨批次漏行+重复行, 而漏 N 行与重 N 行在 COUNT(*) 下可能刚好抵消,报告"数据一致" GetTablePrimaryKeys 的检查置于 len(primaryKeys)==0 判断之前, 否则部分截断会被空检查放过。 配套新增 connection_rows_err_test.go:用标准库 database/sql/driver 实现最小假 driver 模拟迭代中途出错,不引入 go-sqlmock 以保持 go.mod 精简。已验证测试有效性——移除修复后三个断言全部转红。 Closes #170 --- internal/mysql/connection.go | 19 ++ internal/mysql/connection_rows_err_test.go | 207 +++++++++++++++++++++ 2 files changed, 226 insertions(+) create mode 100644 internal/mysql/connection_rows_err_test.go diff --git a/internal/mysql/connection.go b/internal/mysql/connection.go index d1b94e1..9f76c89 100644 --- a/internal/mysql/connection.go +++ b/internal/mysql/connection.go @@ -290,6 +290,12 @@ func (c *Connection) GetTableColumns(tableName string) ([]string, error) { columns = append(columns, field) } + // rows.Next() 在 I/O 中断或包解码错误时返回 false,错误只能从 rows.Err() 取到。 + // 漏检会把"结果集被截断"当成"正常读完",导致少迁列且全程无告警(行数校验只比 COUNT(*),照样通过) + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("遍历表 %s 列信息失败: %w", tableName, err) + } + return columns, nil } @@ -316,6 +322,11 @@ func (c *Connection) GetTableColumnsWithTypes(tableName string) ([]string, map[s columnTypes[field] = colType } + // 同 GetTableColumns:截断的列清单会被 sync_data.go 直接用于构造 SELECT 与 CopyFrom 的 copyColumns + if err := rows.Err(); err != nil { + return nil, nil, fmt.Errorf("遍历表 %s 列信息失败: %w", tableName, err) + } + return columns, columnTypes, nil } @@ -476,6 +487,14 @@ func (c *Connection) GetTablePrimaryKeys(tableName string) ([]string, error) { } } + // 必须置于下方 len(primaryKeys) == 0 检查之前:空检查只兜底"完全没读到主键", + // 而复合主键 (a,b,c) 被截断成 (a,b) 时 len==2 会顺利通过检查, + // 使 keyset 分页的 WHERE (a,b) > (?,?) 游标不再唯一 → 跨批次漏行 + 重复行, + // 且漏 N 行与重 N 行在 COUNT(*) 校验下可能刚好抵消,报告"数据一致" + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("遍历表 %s 主键信息失败: %w", tableName, err) + } + if len(primaryKeys) == 0 { return nil, fmt.Errorf("表 %s 没有主键", tableName) } diff --git a/internal/mysql/connection_rows_err_test.go b/internal/mysql/connection_rows_err_test.go new file mode 100644 index 0000000..6e553bb --- /dev/null +++ b/internal/mysql/connection_rows_err_test.go @@ -0,0 +1,207 @@ +package mysql + +import ( + "context" + "database/sql" + "database/sql/driver" + "errors" + "io" + "strings" + "testing" +) + +// issue #170:验证元数据查询在结果集被截断时必须返回错误,而不是把部分结果当成完整结果。 +// +// 用标准库 database/sql/driver 实现最小假 driver 模拟"迭代中途出错", +// 不引入 go-sqlmock 以保持 go.mod 精简(与现有 connection_test.go 的纯标准库风格一致)。 + +// truncatedRows 在成功返回 failAfter 行后返回 err。 +// err 为 io.EOF 时表示正常读完;为其他错误时表示结果集被截断。 +type truncatedRows struct { + columns []string + data [][]driver.Value + pos int + failAfter int + err error +} + +func (r *truncatedRows) Columns() []string { return r.columns } +func (r *truncatedRows) Close() error { return nil } + +func (r *truncatedRows) Next(dest []driver.Value) error { + if r.pos >= r.failAfter { + return r.err + } + copy(dest, r.data[r.pos]) + r.pos++ + return nil +} + +// truncatedConn 通过 QueryerContext 直接返回预置的 Rows,绕过 Prepare +type truncatedConn struct { + rows driver.Rows +} + +func (c *truncatedConn) Prepare(string) (driver.Stmt, error) { + return nil, errors.New("not implemented") +} +func (c *truncatedConn) Close() error { return nil } +func (c *truncatedConn) Begin() (driver.Tx, error) { return nil, errors.New("not implemented") } + +func (c *truncatedConn) QueryContext(context.Context, string, []driver.NamedValue) (driver.Rows, error) { + return c.rows, nil +} + +type truncatedConnector struct { + conn driver.Conn +} + +func (c *truncatedConnector) Connect(context.Context) (driver.Conn, error) { return c.conn, nil } +func (c *truncatedConnector) Driver() driver.Driver { return nil } + +// newTruncatedConnection 构造一个查询结果在 failAfter 行后被截断的 Connection +func newTruncatedConnection(t *testing.T, columns []string, data [][]driver.Value, failAfter int) *Connection { + t.Helper() + + rowsErr := io.EOF + if failAfter < len(data) { + rowsErr = errors.New("invalid connection: packet sequence mismatch") + } + + db := sql.OpenDB(&truncatedConnector{conn: &truncatedConn{rows: &truncatedRows{ + columns: columns, + data: data, + failAfter: failAfter, + err: rowsErr, + }}}) + t.Cleanup(func() { db.Close() }) + + return &Connection{db: db, ctx: context.Background()} +} + +// showColumnsResult 构造 SHOW COLUMNS 的 6 列结果集 +// 字段顺序:Field, Type, Null, Key, Default, Extra +// 其中 Field/Type/Null/Key/Extra 的 Scan 目标是 string(非 nullable),必须给空串而非 nil; +// 只有 Default 的目标是 sql.NullString,可以为 nil +func showColumnsResult(fields ...string) ([][]driver.Value, []string) { + columns := []string{"Field", "Type", "Null", "Key", "Default", "Extra"} + data := make([][]driver.Value, 0, len(fields)) + for _, f := range fields { + data = append(data, []driver.Value{ + []byte(f), []byte("int(11)"), []byte("YES"), []byte(""), nil, []byte(""), + }) + } + return data, columns +} + +// showKeysResult 构造 SHOW KEYS 的 15 列结果集(MySQL 8.0 形态), +// Column_name 位于索引 4,与 GetTablePrimaryKeys 的提取逻辑一致 +func showKeysResult(keyColumns ...string) ([][]driver.Value, []string) { + columns := []string{ + "Table", "Non_unique", "Key_name", "Seq_in_index", "Column_name", + "Collation", "Cardinality", "Sub_part", "Packed", "Null", + "Index_type", "Comment", "Index_comment", "Visible", "Expression", + } + data := make([][]driver.Value, 0, len(keyColumns)) + for i, col := range keyColumns { + row := make([]driver.Value, len(columns)) + row[0] = []byte("t") + row[1] = int64(0) + row[2] = []byte("PRIMARY") + row[3] = int64(i + 1) + row[4] = []byte(col) + row[10] = []byte("BTREE") + data = append(data, row) + } + return data, columns +} + +func TestGetTableColumnsWithTypesReturnsErrorOnTruncatedResultSet(t *testing.T) { + t.Run("完整结果集应返回全部列且无错误", func(t *testing.T) { + data, columns := showColumnsResult("id", "name", "age") + conn := newTruncatedConnection(t, columns, data, len(data)) + + got, types, err := conn.GetTableColumnsWithTypes("t") + if err != nil { + t.Fatalf("完整结果集不应报错,实际 %v", err) + } + if len(got) != 3 { + t.Fatalf("应返回 3 列,实际 %d 列: %v", len(got), got) + } + if len(types) != 3 { + t.Fatalf("应返回 3 个类型,实际 %d", len(types)) + } + }) + + t.Run("结果集截断必须报错而非返回部分列清单", func(t *testing.T) { + data, columns := showColumnsResult("id", "name", "age") + // 只成功返回第 1 行后中断:修复前会返回 (["id"], types, nil), + // 该残缺列清单会被 sync_data.go 直接用于构造 SELECT 与 CopyFrom 的 copyColumns + conn := newTruncatedConnection(t, columns, data, 1) + + got, _, err := conn.GetTableColumnsWithTypes("t") + if err == nil { + t.Fatalf("结果集被截断时必须返回错误,实际返回 %v 列且 err=nil", got) + } + if got != nil { + t.Errorf("出错时应返回 nil 列清单,实际 %v", got) + } + }) +} + +func TestGetTableColumnsReturnsErrorOnTruncatedResultSet(t *testing.T) { + data, columns := showColumnsResult("id", "name") + conn := newTruncatedConnection(t, columns, data, 1) + + got, err := conn.GetTableColumns("t") + if err == nil { + t.Fatalf("结果集被截断时必须返回错误,实际返回 %v 且 err=nil", got) + } + if got != nil { + t.Errorf("出错时应返回 nil 列清单,实际 %v", got) + } +} + +func TestGetTablePrimaryKeysReturnsErrorOnTruncatedCompositeKey(t *testing.T) { + t.Run("完整复合主键应全部返回", func(t *testing.T) { + data, columns := showKeysResult("a", "b", "c") + conn := newTruncatedConnection(t, columns, data, len(data)) + + got, err := conn.GetTablePrimaryKeys("t") + if err != nil { + t.Fatalf("完整结果集不应报错,实际 %v", err) + } + if len(got) != 3 { + t.Fatalf("应返回 3 个主键列,实际 %v", got) + } + }) + + // 这是本 issue 的核心场景:复合主键 (a,b,c) 被截断成 (a,b) 时, + // len(primaryKeys)==2 会顺利通过"没有主键"的空检查, + // 使 keyset 分页的 WHERE (a,b) > (?,?) 游标不再唯一 → 跨批次漏行 + 重复行, + // 而漏 N 行与重 N 行在 COUNT(*) 校验下可能刚好抵消,报告"数据一致" + t.Run("复合主键被截断必须报错而非返回部分主键", func(t *testing.T) { + data, columns := showKeysResult("a", "b", "c") + conn := newTruncatedConnection(t, columns, data, 2) + + got, err := conn.GetTablePrimaryKeys("t") + if err == nil { + t.Fatalf("复合主键被截断时必须返回错误,实际返回 %v 且 err=nil", got) + } + if got != nil { + t.Errorf("出错时应返回 nil 主键清单,实际 %v", got) + } + }) + + t.Run("完全无主键仍应报原有的没有主键错误", func(t *testing.T) { + conn := newTruncatedConnection(t, []string{"Field"}, nil, 0) + + _, err := conn.GetTablePrimaryKeys("t") + if err == nil { + t.Fatal("无主键表应返回错误") + } + if !strings.Contains(err.Error(), "没有主键") { + t.Errorf("应报'没有主键',实际 %q", err.Error()) + } + }) +} From 09ec3524eacab111b2b2b86913514358e8503031 Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 14:44:13 +0800 Subject: [PATCH 02/10] =?UTF-8?q?fix:=20=E7=A9=BA=E8=A1=A8=E5=88=86?= =?UTF-8?q?=E6=94=AF=E5=85=88=E4=BA=8E=20TRUNCATE=20=E8=BF=94=E5=9B=9E?= =?UTF-8?q?=EF=BC=8C=E5=AF=BC=E8=87=B4=20truncate=5Fbefore=5Fsync=20?= =?UTF-8?q?=E5=AF=B9=E7=A9=BA=E8=A1=A8=E5=A4=B1=E6=95=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit syncSingleTable 中「空表早返回」位于「TRUNCATE 目标表」之前,而 handleEmptyTable 内部只做 log、序列回填与行数校验,不执行 truncate。 触发前提是必然导致迁移终止,而非概率性:MySQL 侧空表 + PG 侧同名表 仍有上一轮数据 + truncate_before_sync=true → 空表分支早返回,truncate 永不执行 → handleEmptyTable 得到 MySQL 0 行 vs PostgreSQL N 行 → evaluateRowCountValidation 在 truncate=true 时返回 "数据校验不一致 (truncate_before_sync=true,终止迁移)" 错误信息指向"数据不一致",真实原因却是工具自己漏做了 truncate, 会把排查引向源库为什么"少数据"。重跑迁移(PG 表已存在且含上一轮数据、 MySQL 侧某表刚好被清空)很容易踩到。 修复:把 TRUNCATE 块移到空表早返回之前,并顺带把「取消检查」提到 TRUNCATE 之前,避免已取消仍执行 DDL。修复后顺序: 取消检查 → TRUNCATE → 空表早返回 → 分页插入。 注:issue #171 正文引用了一处代码中不存在的注释("即使表为空也执行"), 该引用有误,缺陷本身成立,详见 issue 讨论。 测试:truncateTable / handleEmptyTable 接收具体类型 *postgres.Connection 而非接口,无法在不引入接口抽象的前提下用替身断言执行顺序,故本次以 代码审查 + 现有测试不回归为准(build / vet / 全量测试通过)。 Closes #171 --- internal/converter/postgres/sync_data.go | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/internal/converter/postgres/sync_data.go b/internal/converter/postgres/sync_data.go index f986d83..1d8e75a 100644 --- a/internal/converter/postgres/sync_data.go +++ b/internal/converter/postgres/sync_data.go @@ -276,23 +276,27 @@ func syncSingleTable(ctx context.Context, mysqlConn *mysql.Connection, postgresC return fmt.Errorf("同步表 %s 失败: %w", table.Name, err) } - // 如果表为空,处理空表逻辑 - if totalRows == 0 { - return handleEmptyTable(postgresConn, config, table.Name, table.DDL, totalRows, log, logError, mutex, completedTasks, totalTasks, inconsistentTables, printer) - } - // 取消检查:开始实际写入前若已取消则直接退出 if err := ctx.Err(); err != nil { return fmt.Errorf("同步表 %s 已取消: %w", table.Name, err) } // 清空表数据(根据配置) + // 必须位于下方空表早返回之前:否则 MySQL 侧为空表、PG 侧同名表仍有上一轮数据时, + // 陈旧行不会被清掉,handleEmptyTable 的行数校验会以 + // "数据校验不一致: MySQL 0 行, PostgreSQL N 行 (truncate_before_sync=true,终止迁移)" + // 终止整个迁移,而真实原因是工具自己漏做了 truncate(issue #171) if config.Conversion.Options.TruncateBeforeSync { if err := truncateTable(ctx, postgresConn, table.Name, logError); err != nil { return fmt.Errorf("同步表 %s 失败: %w", table.Name, err) } } + // 如果表为空,处理空表逻辑(truncate 已在上方完成,此处只做序列回填与行数校验) + if totalRows == 0 { + return handleEmptyTable(postgresConn, config, table.Name, table.DDL, totalRows, log, logError, mutex, completedTasks, totalTasks, inconsistentTables, printer) + } + // 分页插入数据 processedRows, err := paginateAndInsert(ctx, mysqlConn, postgresConn, config, table, columns, columnTypes, totalRows, log, logError, mutex, progressChan) if err != nil { From 3574648765fe15f136b969e4892c8f03ddc593de Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 14:52:38 +0800 Subject: [PATCH 03/10] =?UTF-8?q?fix:=20worker=20goroutine=20=E8=A1=A5=20r?= =?UTF-8?q?ecover=EF=BC=8C=E5=8D=95=E7=82=B9=20panic=20=E9=99=8D=E7=BA=A7?= =?UTF-8?q?=E4=B8=BA=E5=8D=95=E8=A1=A8/=E5=8D=95=E6=89=B9=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 全仓库此前没有任何 recover(),而两类并发 worker 都在处理来自外部数据库 的任意数据:sync_data.go 的表级同步 worker(MySQL 行 → typedDest → pgx.CopyFrom)与 manager.go 的 runBatchStage 批次 worker(8 个转换阶段的 通用派发器)。 数据层是最容易 panic 的地方:切片越界、类型断言失败、nil map 写入、 getRowSlice 的 (*s)[:numCols] 都在热路径上。未 recover 的 panic 会: - crash 整个迁移进程,而非只让这一张表失败,此时目标表可能正处于 「已 TRUNCATE + 半截数据」状态且无法继续 - 跳过 errorChan 写入,聚合错误列表里完全看不到这次失败 - 不留下「哪张表、哪一批、什么堆栈」的上下文 两处 defer 内加 recover(),转成带对象名/阶段名与 debug.Stack() 的 error 送入各自 errorChan。recover 置于 defer 开头,确保后续资源释放语句自身 异常时错误已入队。manager.go 一处修改即覆盖全部 8 个转换阶段。 测试:runBatchStage 只依赖 m.conversionStats,可用 &Manager{} 零值直接 单测。在既有 TestRunBatchStage 中新增子测试(不新建文件、复用 drainErrors), 断言 panic 被转成 error、错误信息含阶段名/panic 值/堆栈、其余批次照常完成、 阶段统计仍记录。已验证有效性——移除 recover 后测试进程直接崩溃。 sync_data.go 的表级 worker 依赖具体类型 *mysql.Connection / *postgres.Connection,需接口抽象后才能单测,本次以代码审查 + 现有测试不回归为准(build / vet / race / 全量测试通过)。 Closes #172 --- internal/converter/postgres/manager.go | 12 ++++- internal/converter/postgres/manager_test.go | 51 +++++++++++++++++++++ internal/converter/postgres/sync_data.go | 11 +++++ 3 files changed, 73 insertions(+), 1 deletion(-) diff --git a/internal/converter/postgres/manager.go b/internal/converter/postgres/manager.go index 278e072..c8b7d3c 100644 --- a/internal/converter/postgres/manager.go +++ b/internal/converter/postgres/manager.go @@ -6,6 +6,7 @@ import ( "fmt" "os" "path/filepath" + "runtime/debug" "strings" "sync" "sync/atomic" @@ -565,7 +566,16 @@ func runBatchStage[T any](m *Manager, wg *sync.WaitGroup, semaphore chan struct{ batch := objects[i:end] wg.Add(1) go func(batch []T) { - defer wg.Done() + defer func() { + // 批次 worker 顶层 recover:runBatchStage 是全部 8 个转换阶段的通用派发器, + // 此处一改即覆盖表结构/视图/数据/索引/函数/用户/权限各阶段。 + // 不 recover 时 panic 会跳过下方 errorChan 写入并终止整个迁移进程, + // 聚合错误列表里完全看不到这次失败(issue #172) + if r := recover(); r != nil { + errorChan <- fmt.Errorf("阶段「%s」执行时发生 panic: %v\n%s", stageName, r, debug.Stack()) + } + wg.Done() + }() if err := stageFn(batch, semaphore); err != nil { errorChan <- err } diff --git a/internal/converter/postgres/manager_test.go b/internal/converter/postgres/manager_test.go index e13ad58..c69fe2f 100644 --- a/internal/converter/postgres/manager_test.go +++ b/internal/converter/postgres/manager_test.go @@ -71,6 +71,57 @@ func TestRunBatchStage(t *testing.T) { t.Fatal("应聚合到阶段错误") } }) + + // issue #172:worker 顶层必须 recover,否则单个批次的 panic 会终止整个迁移进程, + // 且跳过 errorChan 写入,聚合错误列表里完全看不到这次失败 + t.Run("worker panic 转成错误而非终止进程", func(t *testing.T) { + m := &Manager{} + var wg sync.WaitGroup + semaphore := make(chan struct{}, 4) + errorChan := make(chan error, 8) + + var mu sync.Mutex + var completed []string + // batchSize=1 → 每个对象一个 goroutine;"b" 必然 panic + stageFn := func(batch []string, sem chan struct{}) error { + sem <- struct{}{} + defer func() { <-sem }() + for _, o := range batch { + if o == "b" { + panic("模拟数据层 panic: index out of range") + } + mu.Lock() + completed = append(completed, o) + mu.Unlock() + } + return nil + } + + // 若 worker 顶层没有 recover,这一行会让测试进程直接崩溃 + runBatchStage(m, &wg, semaphore, errorChan, "同步表数据", []string{"a", "b", "c"}, 1, stageFn) + + err := drainErrors(errorChan) + if err == nil { + t.Fatal("panic 应被转成 error 进入聚合通道") + } + msg := err.Error() + for _, want := range []string{"同步表数据", "模拟数据层 panic", "goroutine"} { + if !strings.Contains(msg, want) { + t.Errorf("错误信息应含 %q 以便定位,实际 %q", want, msg) + } + } + + // panic 批次之外的对象必须照常完成——这是「降级为单点失败」而非「整体终止」的关键 + mu.Lock() + got := len(completed) + mu.Unlock() + if got != 2 { + t.Errorf("除 panic 的对象外应完成 2 个,实际 %d 个: %v", got, completed) + } + if len(m.conversionStats) != 1 || m.conversionStats[0].ObjectCount != 3 { + t.Errorf("阶段统计仍应正常记录,实际 %+v", m.conversionStats) + } + }) } // TestManagerContextNilSafe 未通过 NewManager 构造的 Manager diff --git a/internal/converter/postgres/sync_data.go b/internal/converter/postgres/sync_data.go index 1d8e75a..1ff7fcc 100644 --- a/internal/converter/postgres/sync_data.go +++ b/internal/converter/postgres/sync_data.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "os" + "runtime/debug" "strconv" "strings" "sync" @@ -201,6 +202,16 @@ func SyncTableData(ctx context.Context, mysqlConn *mysql.Connection, postgresCon go func(table mysql.TableInfo) { defer func() { + // worker 顶层必须 recover:数据层要把外部数据库返回的任意字节塞进 Go 类型, + // 切片越界 / 类型断言失败 / nil map 写入都可能 panic。不 recover 会让单张表的 + // 异常终止整个迁移进程,而此时目标表可能正处于「已 TRUNCATE + 半截数据」状态。 + // 置于 defer 开头,确保后续资源释放语句自身异常时错误已入队(issue #172) + if r := recover(); r != nil { + select { + case errorChan <- fmt.Errorf("同步表 %s 时发生 panic: %v\n%s", table.Name, r, debug.Stack()): + default: + } + } <-semaphore updateProgress() wg.Done() From ce96c5cdd2be6764eee238a6d4875aabcb101f0a Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 15:09:30 +0800 Subject: [PATCH 04/10] =?UTF-8?q?fix:=20=E6=89=B9=E6=AC=A1=20context=20?= =?UTF-8?q?=E8=A1=A5=E8=B6=85=E6=97=B6=EF=BC=8C=E9=81=BF=E5=85=8D=E7=BD=91?= =?UTF-8?q?=E7=BB=9C=E5=8D=8A=E5=BC=80=E6=97=B6=E8=BF=81=E7=A7=BB=E6=B0=B8?= =?UTF-8?q?=E4=B9=85=E6=8C=82=E6=AD=BB=E4=B8=94=20Ctrl-C=20=E6=97=A0?= =?UTF-8?q?=E6=95=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 批次 context 用 context.WithoutCancel(ctx) 派生,该 ctx 既无 deadline 也不可取消。剥离取消信号本身是正确的(让已开启的批次能完整提交), 但没有补回任何超时。 驱动层证据:go-sql-driver@v1.7.1/connection.go:592-595 的 watchCancel 在 ctx.Done() == nil 时直接返回、不启动 watcher,而 WithoutCancel 的 Done() 正是 nil;中断阻塞 read 依赖的恰是该 watcher(connection.go:620-622 的 <-ctx.Done() → mc.cancel → cleanup 关闭 netConn)。 后果:网络半开或对端 hang 时 goroutine 永久阻塞在 socket read,永久占住 semaphore 槽位,wg.Wait() 永不返回,且根 ctx 已被剥离导致 Ctrl-C 无效 (只能 kill -9),此时目标表可能正处于「已 TRUNCATE + 半截数据」状态。 对比 mysql/metadata.go:176 已用 WithTimeout(120s),数据热路径漏了。 修复:抽出 buildTableSyncContext 纯函数,WithoutCancel 之后套 WithTimeout; paginateAndInsert 每轮增加超时检查并给出可操作的错误提示; config 新增 conversion.limits.table_sync_timeout_seconds(默认 3600)。 与 issue #173 正文的三处实施调整: 1. 配置项由 batch_timeout_seconds 改为 table_sync_timeout_seconds,并在 循环外创建一次而非逐批创建——流式读取的 rows 跨批次复用且绑定首轮 context,逐批 cancel 会关闭其底层连接导致后续批次读取失败。 2. 不提供 0=不限制的逃生口(那等于保留原缺陷),<=0 一律回落默认 3600, 与本结构体其余字段语义一致;超大表可显式配置更大值。 3. 不强制注入 DSN 的 readTimeout/writeTimeout:WithTimeout 已足以让驱动 watcher 生效,而强制注入会改变现有行为,可能让大表 COUNT(*) (服务端算完才返回首个 packet)意外失败,属独立取舍。 测试:buildTableSyncContext 的 Done() 非 nil、deadline 单位换算、 根取消不穿透、<=0 退化行为;ValidateConfig 的默认值回落(0/负数/显式值)。 两个正向断言与退化分支的 Done()==nil 断言互为对照,可捕获 WithTimeout 被误写回 WithoutCancel 的回归。build / vet / gofmt / 全量测试 / race 通过。 Closes #173 --- config.example.yml | 2 + internal/config/config.go | 16 +++++ internal/config/config_test.go | 40 +++++++++++ internal/converter/postgres/sync_data.go | 36 ++++++++-- internal/converter/postgres/sync_data_test.go | 68 +++++++++++++++++++ 5 files changed, 158 insertions(+), 4 deletions(-) diff --git a/config.example.yml b/config.example.yml index 715b53c..7a7d67d 100644 --- a/config.example.yml +++ b/config.example.yml @@ -69,6 +69,8 @@ conversion: max_users_per_batch: 10 # 一次性转换用户的个数限制 max_rows_per_batch: 50000 # 一次性同步数据的行数限制(宽表会按估算行宽自适应下调) batch_insert_size: 50000 # 批量插入的大小(建议与 max_rows_per_batch 保持一致) + table_sync_timeout_seconds: 3600 # 单表数据同步超时(秒)。批次 context 需脱离取消信号以让已开启的批次完整提交, + # 但脱离后若无超时,网络半开时会永久阻塞且 Ctrl-C 无效;超大表可调大(如 86400) # 运行配置 run: diff --git a/internal/config/config.go b/internal/config/config.go index fa4144f..e34f240 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -107,6 +107,15 @@ type LimitsConfig struct { MaxUsersPerBatch int `mapstructure:"max_users_per_batch"` MaxRowsPerBatch int `mapstructure:"max_rows_per_batch"` // 一次性同步数据的行数限制 BatchInsertSize int `mapstructure:"batch_insert_size"` // 批量插入的大小 + // 单表数据同步的超时秒数(issue #173)。批次 context 用 WithoutCancel 派生, + // 以让已开启的批次不受根取消影响;但 WithoutCancel 同时使 Done() 返回 nil, + // go-sql-driver 的 watchCancel 因此不启动监听,网络半开时 goroutine 会永久阻塞、 + // 占住 semaphore 且 Ctrl-C 无效。套一层超时后 Done() 非 nil,驱动 watcher 生效, + // 超时会关闭连接从而中断阻塞的 read。 + // + // 采用表级而非批级语义:流式读取的 rows 跨批次复用且绑定首轮 context, + // 逐批 cancel 会关闭其底层连接,故超时覆盖整表同步过程。 + TableSyncTimeoutSeconds int `mapstructure:"table_sync_timeout_seconds"` } // RunConfig 运行配置 @@ -249,6 +258,13 @@ func (c *Config) ValidateConfig() error { if c.Conversion.Limits.BatchInsertSize <= 0 { c.Conversion.Limits.BatchInsertSize = 50000 // 默认值,与 MaxRowsPerBatch 对齐,避免读 50000 行却按 10000 行分块插入 } + if c.Conversion.Limits.TableSyncTimeoutSeconds <= 0 { + // 默认 1 小时:按 10000 行/秒估算可覆盖约 3600 万行的单表,对绝大多数表足够宽松; + // 同时把「网络半开导致永久挂死且 Ctrl-C 无效」收敛为超时后释放 semaphore 并报错。 + // 不提供 0=不限制的逃生口——那等于保留原缺陷;超大表可显式配置更大值(如 86400)。 + // 与本结构体其余字段的 <=0 回落默认值语义保持一致。 + c.Conversion.Limits.TableSyncTimeoutSeconds = 3600 + } // MPP 配置默认值 if c.Conversion.MPP.Database == "" { diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 5179267..1a20268 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -91,3 +91,43 @@ func TestConvertExclusionLists_Duplicates(t *testing.T) { t.Errorf("SkipViewSet missing key 'view1'") } } + +// TestValidateConfigTableSyncTimeoutDefault issue #173: +// 单表同步超时必须回落为默认值,否则批次 context 无 deadline, +// go-sql-driver 的 watchCancel 不启动监听,网络半开时永久挂死且 Ctrl-C 无效 +func TestValidateConfigTableSyncTimeoutDefault(t *testing.T) { + newMinimalConfig := func() *Config { + c := &Config{} + c.MySQL.Host = "localhost" + c.MySQL.Username = "u" + c.MySQL.Database = "d" + c.PostgreSQL.Host = "localhost" + c.PostgreSQL.Username = "u" + c.PostgreSQL.Database = "d" + return c + } + + cases := []struct { + name string + set int + want int + }{ + {name: "未配置时回落默认 3600", set: 0, want: 3600}, + {name: "负数回落默认 3600", set: -5, want: 3600}, + {name: "显式配置保持不变", set: 7200, want: 7200}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + c := newMinimalConfig() + c.Conversion.Limits.TableSyncTimeoutSeconds = tc.set + + if err := c.ValidateConfig(); err != nil { + t.Fatalf("最小合法配置不应校验失败: %v", err) + } + if got := c.Conversion.Limits.TableSyncTimeoutSeconds; got != tc.want { + t.Errorf("TableSyncTimeoutSeconds = %d, 期望 %d", got, tc.want) + } + }) + } +} diff --git a/internal/converter/postgres/sync_data.go b/internal/converter/postgres/sync_data.go index 1ff7fcc..92e5a37 100644 --- a/internal/converter/postgres/sync_data.go +++ b/internal/converter/postgres/sync_data.go @@ -407,9 +407,29 @@ func truncateTable(ctx context.Context, postgresConn *postgres.Connection, table return nil } +// buildTableSyncContext 构造数据同步用的 context:脱离根 context 的取消信号 +// (让已开启的批次能完整提交),并套一层超时。 +// +// 超时是必需的而非可选的:context.WithoutCancel 返回的 ctx 其 Done() 为 nil, +// 而 go-sql-driver 的 watchCancel 在 ctx.Done() == nil 时直接返回、不启动 watcher +// (mysql@v1.7.1/connection.go:592-595);中断阻塞 read 依赖的正是该 watcher +// (connection.go:620-622:<-ctx.Done() → mc.cancel → cleanup 关闭 netConn)。 +// 没有超时,网络半开时 goroutine 会永久阻塞在 socket read、占住 semaphore, +// 且根 ctx 已被剥离导致 Ctrl-C 无效(issue #173)。 +// +// timeoutSeconds <= 0 时退化为纯 WithoutCancel;配置层已将 <=0 回落为默认 3600, +// 正常路径不会走到该分支,此处仅为函数自身的完备性与可测性。 +func buildTableSyncContext(ctx context.Context, timeoutSeconds int) (context.Context, context.CancelFunc) { + base := context.WithoutCancel(ctx) + if timeoutSeconds <= 0 { + return base, func() {} + } + return context.WithTimeout(base, time.Duration(timeoutSeconds)*time.Second) +} + // paginateAndInsert 分页读取 MySQL 数据并批量插入 PostgreSQL // ctx 为根 context:每轮批次开始前检查取消信号,取消后不再开启新批次; -// 已开启批次的 DB 操作使用 context.WithoutCancel 派生的批次 context, +// 已开启批次的 DB 操作使用 buildTableSyncContext 派生的 context, // 脱离取消信号,保证进行中的批次完整提交后再退出 func paginateAndInsert(ctx context.Context, mysqlConn *mysql.Connection, postgresConn *postgres.Connection, config *config.Config, table mysql.TableInfo, columns []string, columnTypes map[string]string, totalRows int64, log func(format string, args ...interface{}), logError func(errMsg string, args ...interface{}), mutex *sync.Mutex, progressChan chan progressUpdate) (int64, error) { // 获取批量大小配置 @@ -494,15 +514,23 @@ func paginateAndInsert(ctx context.Context, mysqlConn *mysql.Connection, postgre // 进度条状态跟踪:起始时间用于计算速度与 ETA syncStartTime := time.Now() + // 同步 context:脱离根取消信号 + 表级超时(issue #173)。 + // 必须在循环外创建一次:流式读取的 rows 跨批次复用且绑定该 context, + // 逐批创建并 cancel 会关闭其底层连接,导致后续批次读取失败。 + batchCtx, cancelBatch := buildTableSyncContext(ctx, config.Conversion.Limits.TableSyncTimeoutSeconds) + defer cancelBatch() + for { // 取消检查:开启新批次前若已取消则停止,进行中的批次不受影响(完整提交) if err := ctx.Err(); err != nil { return processedRows, fmt.Errorf("表 %s 同步已取消(已处理 %d 行): %w", table.Name, processedRows, err) } - // 批次 context:脱离根 context 的取消信号, - // 保证本轮已开启的批次(MySQL 分页查询 + PG 事务)能完整执行并提交 - batchCtx := context.WithoutCancel(ctx) + // 超时检查:表级 context 到期后不再开启新批次, + // 并给出可操作的提示(否则用户只看到底层驱动的 connection 错误) + if err := batchCtx.Err(); err != nil { + return processedRows, fmt.Errorf("表 %s 同步超时(已处理 %d 行,可通过 conversion.limits.table_sync_timeout_seconds 调整): %w", table.Name, processedRows, err) + } var rows *sql.Rows var currentBatchSize int diff --git a/internal/converter/postgres/sync_data_test.go b/internal/converter/postgres/sync_data_test.go index 31c0490..4244ae0 100644 --- a/internal/converter/postgres/sync_data_test.go +++ b/internal/converter/postgres/sync_data_test.go @@ -2,6 +2,7 @@ package postgres import ( "bytes" + "context" "fmt" "strings" "sync" @@ -468,3 +469,70 @@ func TestBuildOffsetOrderBy(t *testing.T) { }) } } + +// TestBuildTableSyncContext issue #173:批次 context 必须同时满足两个约束—— +// (1) 脱离根取消信号,让已开启的批次能完整提交; +// (2) Done() 非 nil,否则 go-sql-driver 的 watchCancel 在 ctx.Done() == nil 时 +// +// 直接返回、不启动 watcher(mysql@v1.7.1/connection.go:592-595), +// 而中断阻塞 socket read 依赖的正是该 watcher(connection.go:620-622)。 +func TestBuildTableSyncContext(t *testing.T) { + t.Run("有超时时 Done 非 nil 且换算为秒", func(t *testing.T) { + ctx, cancel := buildTableSyncContext(context.Background(), 3600) + defer cancel() + + if ctx.Done() == nil { + t.Fatal("Done() 必须非 nil,否则驱动不会启动取消监听,网络半开时永久阻塞") + } + deadline, ok := ctx.Deadline() + if !ok { + t.Fatal("应设置 deadline") + } + // 验证单位换算:误用 Millisecond/Minute 都会让超时偏离 1 小时一个数量级 + if d := time.Until(deadline); d <= 3590*time.Second || d > 3600*time.Second { + t.Errorf("deadline 应约为 3600 秒后,实际剩余 %v", d) + } + }) + + t.Run("根 context 取消不穿透", func(t *testing.T) { + root, cancelRoot := context.WithCancel(context.Background()) + ctx, cancel := buildTableSyncContext(root, 3600) + defer cancel() + + cancelRoot() + + if err := ctx.Err(); err != nil { + t.Errorf("根取消不应穿透(否则已开启的批次无法完整提交),实际 %v", err) + } + select { + case <-ctx.Done(): + t.Error("根取消后批次 context 不应立即结束") + default: + } + }) + + t.Run("到期后 Err 为 DeadlineExceeded", func(t *testing.T) { + // buildTableSyncContext 的最小粒度是秒,此处用等价组合验证到期语义 + ctx, cancel := context.WithTimeout(context.WithoutCancel(context.Background()), 10*time.Millisecond) + defer cancel() + + select { + case <-ctx.Done(): + if ctx.Err() != context.DeadlineExceeded { + t.Errorf("应为 DeadlineExceeded,实际 %v", ctx.Err()) + } + case <-time.After(time.Second): + t.Fatal("超时未触发") + } + }) + + t.Run("timeout 非正数退化为纯 WithoutCancel", func(t *testing.T) { + for _, timeout := range []int{0, -1} { + ctx, cancel := buildTableSyncContext(context.Background(), timeout) + if ctx.Done() != nil { + t.Errorf("timeout=%d 时 Done() 应为 nil(保持旧行为)", timeout) + } + cancel() // 必须是 no-op 且不 panic + } + }) +} From 1513e33f122158983fe0ebf66d5b3de79b6aceab Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 15:20:06 +0800 Subject: [PATCH 05/10] =?UTF-8?q?fix:=20progressLineGuard=20=E6=94=B9?= =?UTF-8?q?=E7=94=A8=20atomic.Pointer=EF=BC=8C=E6=B6=88=E9=99=A4=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E7=AB=9E=E4=BA=89=E4=B8=8E=20nil=20=E8=A7=A3=E5=BC=95?= =?UTF-8?q?=E7=94=A8=E7=AA=97=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 字段注释原断言「写入发生在同步 worker 启动前与全部结束后,不存在并发读写」, 该断言不成立:写入点位于 syncTableData 函数体内部,而它是 runBatchStage 的 stageFn,后者为每个批次起一个独立 goroutine。211 表 / max_ddl_per_batch=10 约 22 个 goroutine 并发写同一字段,同时 Log/logError 在其他 goroutine 中读。 两个后果: 1. Log/logError 的 `if m.progressLineGuard != nil { m.progressLineGuard.endLine() }` 是 check-then-use,两行之间被 Store(nil) 就会解引用 nil——endLine 首行即 p.mu.Lock()。提出本 issue 时全仓库 recover() 为零命中,该 panic 会直接 终止迁移进程(已由 #172 兜住,但根因仍在)。 2. go test -race 必报 DATA RACE,但 CI 的 -race 只覆盖单元测试,集成步骤用 裸 go build,因此结构性不可见。 修复:字段改为 atomic.Pointer[progressPrinter];两处读取改为先 Load 到局部 变量再判空调用,消除 check-then-use 窗口;修正字段注释。 本次不解决「多批各自的 printer 互相覆盖、争抢 stdout」的功能性问题——那需要 把 printer/progressChan 提升为 Manager 级单例并统一管理消费者 goroutine 生命周期,属独立结构性重构。 测试:TestProgressLineGuardConcurrentAccess(8 写 + 8 读 goroutine × 300 轮) 与 TestProgressLineGuardUsedByLog(断言 Log/logError 确实通过 endLine 复位 dirty,证明读写路径被真实覆盖而非空跑)。 已验证有效性:用普通指针复现修复前写法的临时 sanity 测试在 -race 下报 4 次 WARNING: DATA RACE 并 FAIL,改为 atomic.Pointer 后同一并发模式通过。 build / vet / gofmt / 全量测试 / race 均通过。 Closes #174 --- internal/converter/postgres/manager.go | 25 +++++--- internal/converter/postgres/manager_test.go | 69 +++++++++++++++++++++ 2 files changed, 85 insertions(+), 9 deletions(-) diff --git a/internal/converter/postgres/manager.go b/internal/converter/postgres/manager.go index c8b7d3c..958db12 100644 --- a/internal/converter/postgres/manager.go +++ b/internal/converter/postgres/manager.go @@ -82,8 +82,10 @@ type Manager struct { // 评估结果(仅在评估模式下填充) assessmentResults *AssessmentResults // 数据同步阶段的进度行守卫:控制台输出普通行前先结束未收尾的进度行,避免粘连。 - // 写入发生在同步 worker 启动前与全部结束后,不存在并发读写 - progressLineGuard *progressPrinter + // 必须用 atomic:写入点位于 syncTableData 内部,而它是 runBatchStage 的 stageFn, + // 每个批次一个 goroutine(表数 / max_ddl_per_batch,211 表默认配置下约 22 个并发写), + // 同时 Log/logError 在其他 goroutine 中读取(issue #174) + progressLineGuard atomic.Pointer[progressPrinter] } // AssessmentResults 评估结果 @@ -1516,9 +1518,11 @@ func (m *Manager) convertUsers(users []mysql.UserInfo, semaphore chan struct{}) func (m *Manager) syncTableData(tables []mysql.TableInfo, semaphore chan struct{}) error { progressChan := make(chan progressUpdate, m.config.Conversion.Limits.Concurrency) printer := newProgressPrinter() - // 注册进度行守卫:数据同步期间 log/logError 的控制台输出先结束未收尾的进度行 - m.progressLineGuard = printer - defer func() { m.progressLineGuard = nil }() + // 注册进度行守卫:数据同步期间 log/logError 的控制台输出先结束未收尾的进度行。 + // 本函数由 runBatchStage 按批派发为独立 goroutine,故必须用 atomic 写入(issue #174)。 + // 注:多批各自的 printer 互相覆盖、争抢 stdout 属独立的结构性问题,本次只消除数据竞争 + m.progressLineGuard.Store(printer) + defer func() { m.progressLineGuard.Store(nil) }() return SyncTableData( m.context(), m.mysqlConn, @@ -1681,8 +1685,10 @@ func (m *Manager) Log(format string, args ...interface{}) { // 根据配置决定是否在控制台显示 if m.config.Run.ShowLogInConsole { - if m.progressLineGuard != nil { - m.progressLineGuard.endLine() + // 先 Load 到局部变量再判空调用,消除 check-then-use 窗口: + // 否则两行之间被其他 goroutine Store(nil) 就会解引用 nil(endLine 首行即 p.mu.Lock()) + if guard := m.progressLineGuard.Load(); guard != nil { + guard.endLine() } fmt.Println(logMsg) } @@ -1709,8 +1715,9 @@ func (m *Manager) logError(errMsg string, args ...interface{}) { // 根据配置决定是否在控制台显示 if m.config.Run.ShowConsoleLogs { - if m.progressLineGuard != nil { - m.progressLineGuard.endLine() + // 同 Log:Load 到局部变量再调用,避免 check-then-use 期间被置 nil + if guard := m.progressLineGuard.Load(); guard != nil { + guard.endLine() } fmt.Printf("错误: %s\n", errMsg) } diff --git a/internal/converter/postgres/manager_test.go b/internal/converter/postgres/manager_test.go index c69fe2f..5929e0b 100644 --- a/internal/converter/postgres/manager_test.go +++ b/internal/converter/postgres/manager_test.go @@ -431,3 +431,72 @@ func TestManagerCloseNilFiles(t *testing.T) { t.Fatalf("Close() 在无文件句柄时应返回 nil,实际: %v", err) } } + +// TestProgressLineGuardConcurrentAccess issue #174: +// progressLineGuard 由 runBatchStage 按批派发的多个 goroutine 并发写入, +// 同时被 Log/logError 在其他 goroutine 中读取。本测试需在 -race 下运行才有完整意义: +// 若字段是普通指针,写方与读方之间必报 DATA RACE。 +func TestProgressLineGuardConcurrentAccess(t *testing.T) { + m := newTestManager(1) + + const workers = 8 + const iterations = 300 + var wg sync.WaitGroup + + // 写方:复现 syncTableData 的 Store(printer) / defer Store(nil) + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for j := 0; j < iterations; j++ { + m.progressLineGuard.Store(newProgressPrinter()) + m.progressLineGuard.Store(nil) + } + }() + } + + // 读方:复现 Log/logError 中的 Load-then-call。 + // 若先判空再用字段本身(而非 Load 到局部变量),两行之间被 Store(nil) + // 就会解引用 nil——endLine 首行即 p.mu.Lock() + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for j := 0; j < iterations; j++ { + if guard := m.progressLineGuard.Load(); guard != nil { + guard.endLine() + } + } + }() + } + + wg.Wait() +} + +// TestProgressLineGuardUsedByLog 验证守卫确实被 Log/logError 使用: +// 输出普通行前必须收尾未结束的进度行,否则两者会粘连 +func TestProgressLineGuardUsedByLog(t *testing.T) { + m := newTestManager(1) + m.config.Run.ShowLogInConsole = true + m.config.Run.ShowConsoleLogs = true + + // 未注册守卫时不应 panic + m.Log("无守卫") + m.logError("无守卫") + + // 注册处于 dirty 状态的守卫:endLine 应真正收尾并把 dirty 复位 + p := newProgressPrinter() + p.dirty = true + m.progressLineGuard.Store(p) + + m.Log("有守卫") + if p.dirty { + t.Error("Log 输出前应已通过 endLine 收尾进度行,dirty 应被复位") + } + + p.dirty = true + m.logError("有守卫") + if p.dirty { + t.Error("logError 输出前应已通过 endLine 收尾进度行,dirty 应被复位") + } +} From 16876af19517f748bf5dd865220ffc53e7cacc31 Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 15:29:46 +0800 Subject: [PATCH 06/10] =?UTF-8?q?fix:=20consistent=5Fsnapshot=20=E4=B8=8E?= =?UTF-8?q?=E5=B9=B6=E5=8F=91=E4=BA=92=E6=96=A5=E6=A0=A1=E9=AA=8C=20+=20?= =?UTF-8?q?=E6=98=BE=E5=BC=8F=E6=8C=87=E5=AE=9A=20REPEATABLE=20READ=20?= =?UTF-8?q?=E9=9A=94=E7=A6=BB=E7=BA=A7=E5=88=AB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题一:快照事务被多 goroutine 共用 consistent_snapshot=true 时 Connection 上只有一个 snapshotTx *sql.Tx, querier() 对所有数据读取(GetTableData / QueryTableRows / 两个 keyset 分页 / GetTableRowCount)返回同一个 Tx,而 semaphore 容量 = Concurrency(默认 10)。 一个 *sql.Tx 绑定一条 MySQL 连接,且标准库不串行化同一 Tx 上的并发查询—— sql.go:2245-2263 的 Tx.grabConn 只持 closemu.RLock(),该锁的注释写明用途是 「prevents the transaction from closing while there is an active query」, 读锁允许多 goroutine 同时持有。并发查询落到同一 packet buffer, go-sql-driver@v1.7.1/buffer.go:136/158/169/177 在 b.length > 0 时返回 ErrBusyBuffer(注释三次重复 "Only one buffer (total) can be used at a time")。 isTransientConnError 又把 ErrBusyBuffer 判为可重试,而 retryOnTransientConn 的注释假设「重新发起查询即可由连接池换新连接执行」——在 Tx 上该前提不成立 (Tx 绑定固定连接),于是白白重试 2 次后整表失败,报出误导性的 busy buffer。 修复:ValidateConfig 增加互斥校验,放在 Concurrency 默认值回落之后以免误判。 问题二:隔离级别未显式指定 BeginTx(ctx, nil) 用服务端默认隔离级别,而 MySQL 的 WITH CONSISTENT SNAPSHOT 只在 REPEATABLE READ 下建立快照,在 READ COMMITTED 下被忽略并仅产生 warning。 源库为减少 gap lock 配成 RC 并不罕见,此时快照静默失效,用户却看到 manager.go:522 的「数据读取不受源库并发写入影响」。 修复:BeginTx 传 &sql.TxOptions{Isolation: LevelRepeatableRead, ReadOnly: true}。 更正审查过程中的一处判断:曾认为随后那条 START TRANSACTION WITH CONSISTENT SNAPSHOT 与 BeginTx 冗余应删除——这是错的,本次保留。BeginTx 只设定隔离级别, InnoDB 的一致性读快照在首次读取时才建立,各表首次读取时间点不同则跨表不一致; 该语句让快照在事务开始即刻建立。已在代码注释中写明,避免后来者再次「优化」掉。 回归风险核查:CI 配置、config.example.yml、集成测试脚本均未开启 consistent_snapshot(默认 false),新校验不影响现有流程。已更新 config.example.yml 该项注释说明互斥关系与隔离级别自动设置。 测试:互斥校验的 5 种组合(含「未配置并发时回落默认 1 不应误判」, 用于锁定校验位置在默认值回落之后)。已验证有效性——移除校验后 「快照开启 + 高并发必须拒绝」用例 FAIL。BeginConsistentSnapshot 依赖真实 MySQL 连接,以代码审查 + 现有测试不回归为准。 Closes #175 --- config.example.yml | 5 +++- internal/config/config.go | 14 +++++++++ internal/config/config_test.go | 54 ++++++++++++++++++++++++++++++++++ internal/mysql/connection.go | 22 ++++++++++++-- 4 files changed, 91 insertions(+), 4 deletions(-) diff --git a/config.example.yml b/config.example.yml index 7a7d67d..fa130b1 100644 --- a/config.example.yml +++ b/config.example.yml @@ -11,7 +11,10 @@ mysql: conn_max_lifetime: 3600 # 连接池配置的最大生命周期(秒) # 注意:parseTime 由工具强制开启(时间列依赖 time.Time 语义进入 CopyFrom),此处配置该项会被忽略 connection_params: charset=utf8mb4&interpolateParams=true&readTimeout=60s&writeTimeout=60s&timeout=30s # MySQL连接参数 - consistent_snapshot: false # 数据读取使用一致性快照(源库并发写入时保证读到一致数据;要求 REPEATABLE READ 隔离级别) + consistent_snapshot: false # 数据读取使用一致性快照(源库并发写入时保证读到一致数据) + # 注意:与 conversion.limits.concurrency > 1 互斥,开启时须把 concurrency 设为 1 + # (快照事务绑定单条 MySQL 连接,无法被多个 goroutine 并发使用,配置校验会直接拒绝该组合) + # 隔离级别由工具自动设为 REPEATABLE READ,无需在源库侧配置 # PostgreSQL连接配置 postgresql: diff --git a/internal/config/config.go b/internal/config/config.go index e34f240..f3808a1 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -279,5 +279,19 @@ func (c *Config) ValidateConfig() error { } } + // consistent_snapshot 与并发互斥(issue #175): + // 快照事务是 mysql.Connection 上的单个 *sql.Tx,绑定一条 MySQL 连接, + // querier() 对所有数据读取返回同一个 Tx。concurrency>1 时多个 goroutine 会在 + // 同一连接上并发查询,go-sql-driver 的 packet buffer 不支持并发使用而返回 + // ErrBusyBuffer;isTransientConnError 又把它判为可重试,但 Tx 绑定固定连接 + // 不会换连接,重试必然再次失败 → 整表同步失败并报出误导性的 "busy buffer"。 + // 必须放在上方 Concurrency 默认值回落之后,否则未显式配置并发时会误判。 + if c.MySQL.ConsistentSnapshot && c.Conversion.Limits.Concurrency > 1 { + return fmt.Errorf("consistent_snapshot 与 concurrency > 1 互斥(当前 concurrency=%d): "+ + "一致性快照事务绑定单条 MySQL 连接,无法被多个 goroutine 并发使用;"+ + "请关闭 mysql.consistent_snapshot,或将 conversion.limits.concurrency 设为 1", + c.Conversion.Limits.Concurrency) + } + return nil } diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 1a20268..fbdacf2 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -1,6 +1,7 @@ package config import ( + "strings" "testing" ) @@ -131,3 +132,56 @@ func TestValidateConfigTableSyncTimeoutDefault(t *testing.T) { }) } } + +// TestValidateConfigConsistentSnapshotConcurrencyMutex issue #175: +// 快照事务是 mysql.Connection 上的单个 *sql.Tx,绑定一条连接, +// querier() 对所有数据读取返回同一个 Tx;concurrency>1 时并发查询会撞 +// go-sql-driver 的 packet buffer(ErrBusyBuffer),且因 Tx 不换连接使重试无效。 +// 该组合必须在配置校验阶段就被拒绝,而不是运行到数据同步才失败。 +func TestValidateConfigConsistentSnapshotConcurrencyMutex(t *testing.T) { + newCfg := func(snapshot bool, concurrency int) *Config { + c := &Config{} + c.MySQL.Host = "localhost" + c.MySQL.Username = "u" + c.MySQL.Database = "d" + c.MySQL.ConsistentSnapshot = snapshot + c.PostgreSQL.Host = "localhost" + c.PostgreSQL.Username = "u" + c.PostgreSQL.Database = "d" + c.Conversion.Limits.Concurrency = concurrency + return c + } + + cases := []struct { + name string + snapshot bool + concurrency int + wantErr bool + }{ + {name: "快照关闭 + 高并发(默认场景)", snapshot: false, concurrency: 10, wantErr: false}, + {name: "快照关闭 + 单并发", snapshot: false, concurrency: 1, wantErr: false}, + {name: "快照开启 + 单并发应允许", snapshot: true, concurrency: 1, wantErr: false}, + {name: "快照开启 + 高并发必须拒绝", snapshot: true, concurrency: 10, wantErr: true}, + // 校验必须位于 Concurrency 默认值回落之后:未配置并发时回落为 1,不应误判 + {name: "快照开启 + 未配置并发(回落默认 1)不应误判", snapshot: true, concurrency: 0, wantErr: false}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + err := newCfg(tc.snapshot, tc.concurrency).ValidateConfig() + + if tc.wantErr { + if err == nil { + t.Fatal("consistent_snapshot 与 concurrency>1 的组合应被拒绝") + } + if !strings.Contains(err.Error(), "consistent_snapshot") { + t.Errorf("错误信息应指明冲突的配置项,实际 %q", err.Error()) + } + return + } + if err != nil { + t.Fatalf("该组合不应校验失败: %v", err) + } + }) + } +} diff --git a/internal/mysql/connection.go b/internal/mysql/connection.go index 9f76c89..3045c47 100644 --- a/internal/mysql/connection.go +++ b/internal/mysql/connection.go @@ -163,13 +163,29 @@ type Connection struct { } // BeginConsistentSnapshot 开启一致性快照事务(P1-07) -// 要求 MySQL 事务隔离级别为 REPEATABLE READ(InnoDB 默认); -// 开启后数据读取类查询(QueryTableRows/GetTableData*/GetTableRowCount)均通过该事务执行 +// 隔离级别由本函数显式指定为 REPEATABLE READ,不依赖服务端默认值; +// 开启后数据读取类查询(QueryTableRows/GetTableData*/GetTableRowCount)均通过该事务执行。 +// +// 注意:该事务是 Connection 上的单个 *sql.Tx,绑定一条 MySQL 连接, +// 因此只在 concurrency=1 下可用;配置层已对 consistent_snapshot 与 +// concurrency>1 的组合做互斥校验(issue #175) func (c *Connection) BeginConsistentSnapshot(ctx context.Context) error { - tx, err := c.db.BeginTx(ctx, nil) + // 显式指定 REPEATABLE READ:MySQL 的 WITH CONSISTENT SNAPSHOT 只在 RR 下建立快照, + // 在 READ COMMITTED 下会被忽略并仅产生 warning,导致「一致性快照」静默失效。 + // 源库为减少 gap lock 而配成 RC 的情况并不罕见,不能依赖服务端默认值(issue #175)。 + // ReadOnly 让 InnoDB 走只读事务路径,同时避免任何误写源库的可能。 + tx, err := c.db.BeginTx(ctx, &sql.TxOptions{ + Isolation: sql.LevelRepeatableRead, + ReadOnly: true, + }) if err != nil { return fmt.Errorf("开始快照事务失败: %w", err) } + // 下面这条 START TRANSACTION 不是冗余,请勿「优化」掉: + // BeginTx 只设定隔离级别,InnoDB 的一致性读快照是在首次读取时才建立, + // 若各表首次读取的时间点不同则跨表不一致;WITH CONSISTENT SNAPSHOT 让快照 + // 在事务开始即刻建立。MySQL 会隐式提交上方 BeginTx 开启的空只读事务, + // 多一次往返但无副作用,这是拿到跨表一致快照的必要代价。 if _, err := tx.ExecContext(ctx, "START TRANSACTION WITH CONSISTENT SNAPSHOT"); err != nil { tx.Rollback() return fmt.Errorf("开启一致性快照失败: %w", err) From ec21e1fff19266bb61e041ee10478c63f136c098 Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 15:35:00 +0800 Subject: [PATCH 07/10] =?UTF-8?q?refactor:=20=E5=88=A0=E9=99=A4=E6=97=A0?= =?UTF-8?q?=E6=95=88=E7=9A=84=20SELECT=20MAX=20=E5=85=9C=E5=BA=95=EF=BC=8C?= =?UTF-8?q?=E6=9D=9C=E7=BB=9D=E8=B7=A8=E5=BA=93=E6=B8=B8=E6=A0=87=E6=B1=A1?= =?UTF-8?q?=E6=9F=93=E7=9A=84=E6=9C=AA=E6=9D=A5=E9=99=B7=E9=98=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit BatchInsertDataWithCompositeKeys 末尾曾在 lastValue 为 nil 时,于本 PG 事务上 执行 SELECT MAX(主键) FROM 目标表,并把结果当作下一轮 MySQL keyset 游标返回。 用目标库状态充当源库读取游标在设计上是错的。 完整追踪调用链后确认当前不可达,故定级为低(详见 issue #176): - 复合主键路径(sync_data.go:575)读取用 compositeLastValues,而 lastValue 仅在 len(primaryKeyIndexes)==1 时赋值,故此路径下 MAX 结果从不被使用; - 单主键路径经 BatchInsertDataWithTransactionAndGetLastValue(connection.go:950) 委托进来,:988-1009 用 EqualFold 定位后 resolvedPrimaryKeys 取自 copyColumns[i], 故 :1012-1019 的精确匹配必然成功,有数据时 lastValue 正常赋值、MAX 不触发; - 唯一触发场景是「本批 0 行」,而 currentBatchSize==0 会让 sync_data.go:594-597 立即 break,结果同样不被使用。 仍应删除的三个理由: 1. 复合主键场景每批一次无用的 PG 查询; 2. 注释称其为「后备方案」,误导读者以为存在跨库游标兜底能力; 3. 未来陷阱——一旦有人改动 :1048-1050 的赋值条件或让 :1003-1009 的 fallback 变得可达(如列名大小写策略调整),立刻变成真实的数据损坏: skip_existing_tables=true 且目标表已有数据时,PG 的 MAX 大于 MySQL 已读位置 会导致中间整段数据静默丢失,小于则重复读取撞主键,而 validateData 只比 COUNT(*),漏行与重行可能刚好抵消、报告仍显示「数据一致」。 在返回处留注释说明游标的唯一合法来源,以及不可用时应退回 OFFSET 分页并告警、 绝不应查询目标库。删除后 SELECT MAX 在本文件仅剩 :780 的序列回填(正确用法)。 行为无变化,现有测试未覆盖该分支。build / vet / gofmt / 全量测试通过。 Closes #176 --- internal/postgres/connection.go | 20 ++++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/internal/postgres/connection.go b/internal/postgres/connection.go index ce610a4..c8e233e 100644 --- a/internal/postgres/connection.go +++ b/internal/postgres/connection.go @@ -1096,16 +1096,16 @@ func (c *Connection) BatchInsertDataWithCompositeKeys(ctx context.Context, tx pg return 0, nil, nil, err } - // 只有在没有找到主键值的情况下,才执行 MAX 查询(作为后备方案) - if len(resolvedPrimaryKeys) > 0 && lastValue == nil { - // 对于复合主键,只查询第一个主键列 - query := fmt.Sprintf("SELECT MAX(\"%s\") FROM \"%s\"", resolvedPrimaryKeys[0], tableName) - err := tx.QueryRow(ctx, query).Scan(&lastValue) - if err != nil && err != pgx.ErrNoRows { - return 0, nil, nil, fmt.Errorf("获取最后一个主键值失败:%w", err) - } - } - + // 主键游标只能来自本批实际写入的行(lastValue / compositeLastValues)。 + // + // 此处曾有一段「兜底」:lastValue 为 nil 时在本 PG 事务上执行 + // SELECT MAX(主键) FROM 目标表,并把结果当作下一轮 MySQL keyset 游标返回。 + // 用目标库的状态充当源库的读取游标在设计上是错的——skip_existing_tables=true + // 且目标表已有数据时,PG 的 MAX 大于 MySQL 已读位置会导致中间整段数据静默丢失, + // 小于则重复读取撞主键;而 validateData 只比 COUNT(*),漏行与重行可能刚好抵消, + // 报告仍显示「数据一致」。已删除(issue #176)。 + // + // 游标不可用时应由调用方退回 OFFSET 分页并告警,绝不应查询目标库。 return totalRows, lastValue, compositeLastValues, nil } From 148381e90aa5c0dc470944ba7d160cbb8dff39cf Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 15:45:32 +0800 Subject: [PATCH 08/10] =?UTF-8?q?fix:=20bit=20=E5=88=97=E7=9A=84=20DEFAULT?= =?UTF-8?q?=20b'...'=20=E4=BD=8D=E5=AD=97=E9=9D=A2=E9=87=8F=E8=BD=AC?= =?UTF-8?q?=E4=B8=BA=E5=8D=81=E8=BF=9B=E5=88=B6=E6=95=B4=E6=95=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit bit(n) 已正确映射为 BIGINT(bit(64) 为 NUMERIC(20,0)),但 DEFAULT 子句中的 MySQL 位字面量 b'0101' 完全未处理,原样透传给 PostgreSQL。PG 中 b'...' 是 bit 类型字面量,赋给 BIGINT 列会因类型不匹配被拒绝。 实证(本地真实 PostgreSQL): 修复前 → "flags" BIGINT default b'0' ERROR: 42804: column "flags" is of type bigint but default expression is of type bit 修复后 → "flags" BIGINT default 0 CREATE TABLE 成功,且 INSERT ... DEFAULT VALUES RETURNING 得到 0|1|10 —— 默认值的数值语义也正确(b'1010' = 10), 不只是让语法通过 处理逻辑此前完全缺失:grep "b'0'|b'1'|bitLiteral" internal/ 零命中。 cleanTypeDefinition 只映射类型、清理字符集与零日期默认值,不涉及位字面量。 真实可达:MySQL 对带默认值的 bit 列,SHOW CREATE TABLE 输出的正是 DEFAULT b'0' 形态。CI 不可见的原因是 create_table.sql 中 grep 不到任何 bit ... DEFAULT b'...' 用例——case_64_bit_types 与 case_190_bit_full 都只有裸 bit 列无默认值,转换成功从而掩盖了该缺陷。 修复:新增包级正则 reBitLiteralDefault = (?i)(\bdefault\s+)b'([01]+)', 用 strconv.ParseUint(bin, 2, 64) 转十进制后回填。 - 限定在 default 之后,避免误伤其他位置; - cleanTypeDefinition 的小写化用 toLowerOutsideQuotes,引号内内容保持原样, 故 b 会被小写而位串不受影响,正则用 (?i) 兼容大写 B'...' 输入; - 位宽上限取 64:MySQL bit 最大即 bit(64),其最大值 18446744073709551615 恰为 uint64 上界,ParseUint 不会溢出;解析失败时保留原样, 让 PG 报错暴露而非静默吞掉。 测试:补 5 个转换用例(b'0'/b'1'/多位串/bit(64) 全 1 上界/大写 B'1010') + 断言产物不再残留 b';另加一个守卫测试验证字符串字面量内容不被误伤。 已验证有效性——移除转换逻辑后测试 FAIL 并打印出未转换的 b'0'。 build / vet / gofmt / 全量测试通过。 Closes #177 --- internal/converter/postgres/sync_tableddl.go | 22 ++++++ .../converter/postgres/sync_tableddl_test.go | 71 +++++++++++++++++++ 2 files changed, 93 insertions(+) diff --git a/internal/converter/postgres/sync_tableddl.go b/internal/converter/postgres/sync_tableddl.go index ff67ae0..91229f5 100644 --- a/internal/converter/postgres/sync_tableddl.go +++ b/internal/converter/postgres/sync_tableddl.go @@ -47,6 +47,11 @@ var ( // BIT(64) 最大值 18446744073709551615 超出 BIGINT 上限,需映射为 NUMERIC(20,0) // (MySQL BIT 宽度上限为 64,BIT(n<=63) 走标准 bit -> BIGINT 映射) reBit64 = regexp.MustCompile(`(?i)\bbit\(64\)`) + // DEFAULT 子句中的 MySQL 位字面量 b'0101'(issue #177)。 + // bit(n) 已映射为 BIGINT(bit(64) 为 NUMERIC(20,0)),而 b'...' 在 PG 中是 bit + // 类型字面量,透传会报 42804「column is of type bigint but default expression + // is of type bit」,故需转为十进制整数。限定在 default 之后以避免误伤其他位置。 + reBitLiteralDefault = regexp.MustCompile(`(?i)(\bdefault\s+)b'([01]+)'`) // 类型清理相关正则 reVarcharMissingParen = regexp.MustCompile(`(?i)varchar\(\d+`) @@ -1449,6 +1454,23 @@ func cleanTypeDefinition(typeDefinition string, tinyInt1AsBoolean bool) string { lowerTypeDef = strings.ReplaceAll(lowerTypeDef, " default '0000-00-00 00:00:00.000'", "") lowerTypeDef = strings.ReplaceAll(lowerTypeDef, " default '0000-00-00'", "") + // MySQL 位字面量 DEFAULT b'0101' → 十进制整数(issue #177): + // bit(n) 已映射为 BIGINT(bit(64) 为 NUMERIC(20,0)),而 b'...' 在 PG 中是 bit + // 类型字面量,透传会报 42804。MySQL bit 宽度上限为 64,其最大值 + // 18446744073709551615 恰为 uint64 上界,故 ParseUint(_, 2, 64) 不会溢出; + // 解析失败时保留原样,让 PG 报错暴露而非静默吞掉。 + lowerTypeDef = reBitLiteralDefault.ReplaceAllStringFunc(lowerTypeDef, func(m string) string { + match := reBitLiteralDefault.FindStringSubmatch(m) + if len(match) != 3 { + return m + } + v, err := strconv.ParseUint(match[2], 2, 64) + if err != nil { + return m + } + return match[1] + strconv.FormatUint(v, 10) + }) + if strings.Contains(strings.ToUpper(lowerTypeDef), "GENERATED ALWAYS AS") { lowerTypeDef = reCharsetPrefix.ReplaceAllString(lowerTypeDef, "$1") lowerTypeDef = reVirtual.ReplaceAllString(lowerTypeDef, " STORED") diff --git a/internal/converter/postgres/sync_tableddl_test.go b/internal/converter/postgres/sync_tableddl_test.go index 60c5086..4ad4ea4 100644 --- a/internal/converter/postgres/sync_tableddl_test.go +++ b/internal/converter/postgres/sync_tableddl_test.go @@ -685,6 +685,77 @@ func TestConvertTableDDL_BitTypes(t *testing.T) { } } +// TestConvertTableDDL_BitLiteralDefault issue #177: +// MySQL 位字面量 DEFAULT b'0101' 必须转为十进制整数。 +// bit(n) 映射为 BIGINT(bit(64) 为 NUMERIC(20,0)),而 b'...' 在 PG 中是 bit 类型 +// 字面量,透传会报 42804「column is of type bigint but default expression is of type bit」 +func TestConvertTableDDL_BitLiteralDefault(t *testing.T) { + const ones64 = "1111111111111111111111111111111111111111111111111111111111111111" + + mysqlDDL := `CREATE TABLE test_bit_default ( + b_zero bit(8) DEFAULT b'0', + b_one bit(1) DEFAULT b'1', + b_multi bit(8) DEFAULT b'1010', + b_max bit(64) DEFAULT b'` + ones64 + `', + b_upper bit(4) DEFAULT B'1010', + b_null bit(1) DEFAULT NULL +) ENGINE=InnoDB` + + result, err := ConvertTableDDL(mysqlDDL, false) + if err != nil { + t.Fatalf("ConvertTableDDL failed: %v", err) + } + + checks := []struct { + name string + want string + }{ + {"b'0' 转 0", `"b_zero" BIGINT default 0`}, + {"b'1' 转 1", `"b_one" BIGINT default 1`}, + {"b'1010' 转 10", `"b_multi" BIGINT default 10`}, + // bit(64) 全 1 恰为 uint64 上界,验证不溢出 + {"bit(64) 全 1 转 uint64 上界", `"b_max" NUMERIC(20,0) default 18446744073709551615`}, + // 大写 B'...' 经 toLowerOutsideQuotes 后同样应被转换 + {"大写 B'1010' 转 10", `"b_upper" BIGINT default 10`}, + } + for _, c := range checks { + t.Run(c.name, func(t *testing.T) { + if !strings.Contains(result.DDL, c.want) { + t.Errorf("DDL 应包含 %q,实际 DDL: %s", c.want, result.DDL) + } + }) + } + + // 产物中不得残留任何 MySQL 位字面量 + if strings.Contains(strings.ToLower(result.DDL), "b'") { + t.Errorf("DDL 不应残留位字面量 b'...': %s", result.DDL) + } +} + +// TestConvertTableDDL_BitLiteralDefaultKeepsStringLiterals +// 位字面量转换必须限定在 DEFAULT 之后,不得误伤字符串字面量内容 +func TestConvertTableDDL_BitLiteralDefaultKeepsStringLiterals(t *testing.T) { + mysqlDDL := `CREATE TABLE test_bit_guard ( + note varchar(50) DEFAULT 'xb''1010''', + flag bit(4) DEFAULT b'1010' +) ENGINE=InnoDB` + + result, err := ConvertTableDDL(mysqlDDL, false) + if err != nil { + t.Fatalf("ConvertTableDDL failed: %v", err) + } + + // 字符串字面量 'xb''1010''' 的内容必须原样保留(x 前缀确保它不以 b' 开头, + // 从而与位字面量形态区分开),不能被当成位字面量转成十进制 + if !strings.Contains(result.DDL, "xb''1010''") { + t.Errorf("字符串字面量内容应原样保留,实际 DDL: %s", result.DDL) + } + // 真正的位字面量仍应被转换 + if !strings.Contains(result.DDL, `"flag" BIGINT default 10`) { + t.Errorf("位字面量应转为十进制,实际 DDL: %s", result.DDL) + } +} + // TestCleanTypeDefinition_TinyInt1Mapping tinyint(1) 映射策略(P2-03 + 42883 修复): // 默认映射为 SMALLINT 保留整数语义(兼容视图/函数中 `col = 1` 等用法), // 显式开启 tinyInt1AsBoolean 时映射为 BOOLEAN。 From 7d9da070fbd0f6a4fbb7bce4e0a273bd576c015f Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 16:07:34 +0800 Subject: [PATCH 09/10] =?UTF-8?q?fix:=20COMMENT=20=E4=B8=AD=E7=9A=84?= =?UTF-8?q?=E5=8F=8D=E6=96=9C=E6=9D=A0=E8=BD=AC=E4=B9=89=E5=87=BB=E7=A9=BF?= =?UTF-8?q?=E6=B3=A8=E9=87=8A=E6=8F=90=E5=8F=96=EF=BC=8C=E7=A0=B4=E5=9D=8F?= =?UTF-8?q?=E5=88=97=E5=AE=9A=E4=B9=89=E5=B9=B6=E5=AF=BC=E8=87=B4=20PG=204?= =?UTF-8?q?2601?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 注:本项原计划先建 issue,但 issue 创建被环境策略拦截,故将完整问题描述 保留在此 commit message 中。 问题 ---- reComment 只认 '' 双写转义,不认 MySQL 默认的反斜杠转义 \': reComment = (?i)\s+comment\s+'((?:[^']|'')*)'\s*,?\s*|\s+comment\s+"([^"]*)"\s*,?\s* 而 MySQL 的 SHOW CREATE TABLE 对注释中的单引号输出的正是反斜杠转义形态 (除非 sql_mode 含 NO_BACKSLASH_ESCAPES)。英文注释带撇号(don't、user's) 在真实库中极其常见。 实证 ---- 输入(SHOW CREATE TABLE 的真实输出形态): CREATE TABLE `t_comment` ( `id` int(11) NOT NULL AUTO_INCREMENT, `note` varchar(100) DEFAULT NULL COMMENT 'user\'s note', `plain` varchar(50) DEFAULT NULL COMMENT 'normal comment', PRIMARY KEY (`id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='table\'s comment' 修复前 ConvertTableDDL 的产物: CREATE TABLE "t_comment" ("id" SERIAL not null ,"note" VARCHAR(100)s note',"plain" VARCHAR(50), PRIMARY KEY ("id")) 转换器返回 err=nil、Warnings=[](完全静默)。真实 PostgreSQL 上执行: ERROR: 42601: syntax error at or near "s" 匹配在 user\ 处提前闭合,导致三处破坏:残留 s note' 进入列定义、该列的 DEFAULT NULL 消失、ColumnComments 中 note 被截断为 "user\"(TableComment 同样被截断为 "table\")。列定义逐行处理,故后续 plain 列本身不受影响。 修复后(同样在真实 PG 上端到端验证): CREATE TABLE "t_comment" ("id" SERIAL not null ,"note" VARCHAR(100),"plain" VARCHAR(50), PRIMARY KEY ("id")); COMMENT ON TABLE "t_comment" IS 'table''s comment'; COMMENT ON COLUMN "t_comment"."note" IS 'user''s note'; COMMENT ON COLUMN "t_comment"."plain" IS 'normal comment'; 四条语句全部执行成功,且读回内容为 table's comment | user's note | normal comment 即注释语义完全正确,不只是让语法通过。 可达性 ------ 真实可达。CI 不可见的原因是 scripts/mysql/create_table.sql 中 grep \' 零命中 ——测试集的注释都不含撇号。 改动 ---- 1. reComment 两个分支都支持反斜杠转义: 单引号 '((?:[^']|'')*)' → '((?:[^'\\]|''|\\.)*)' 双引号 "([^"]*)" → "((?:[^"\\]|\\.)*)" 2. reTableComment(表级 COMMENT='...')同样处理。 3. 新增 unescapeMySQLStringLiteral,按 MySQL 语义把捕获内容还原为原始文本: \' \" \\ → 对应字符;'' → ';\n \r \t \b → 控制符;\0 → 丢弃(PG 文本 不允许 NUL 字节);\% \_ → 保留反斜杠(MySQL 语义);其他 \x → 反斜杠 被忽略、取字符本身(MySQL 文档:For all other escape sequences, backslash is ignored)。按字节遍历对 UTF-8 安全(多字节序列每字节 >= 0x80, 不会与 \ 0x5C 或 ' 0x27 冲突)。 4. 两个提取点(列注释、表注释)改为经该函数还原。 刻意不在提取端做 PG 转义:下游生成端已经做了——GenerateColumnCommentsSQL 与 processComment 均把 ' 转为 ''。提取端只负责还原,避免二次转义。 测试 ---- - TestUnescapeMySQLStringLiteral:14 个用例覆盖全部转义规则(含中文混排、 末尾孤立反斜杠、多转义共存、NUL 丢弃、\% \_ 保留) - TestConvertTableDDL_CommentBackslashEscape:端到端断言无残留、两列定义完整、 括号平衡、ColumnComments 与 TableComment 还原正确 - TestGenerateColumnCommentsSQL_EscapesSingleQuote:锁定生成端的 ' → '' 转义 已验证有效性:把 reComment/reTableComment 改回旧正则后测试 FAIL,并精确打印出 残留的 s note' 与被截断的 "user\" / "table\"。 build / vet / gofmt / 全量测试通过。 附带发现(未在本次修复,建议单独处理) ------------------------------------ processComment(manager.go:1196-1210)的替换顺序有误:先把 \n \r \t 转成 字面 \n \r \t,最后才执行 \\ → \\\\,于是刚生成的反斜杠被再次翻倍成 \\n \\r \\t;PG 在 standard_conforming_strings=on(默认)下会按字面存储。 反斜杠翻倍应当最先执行。此外 processComment(转成字面 \n)与 GenerateColumnCommentsSQL(直接删除换行符)是两条并行且策略不一致的注释 处理路径,建议统一。 --- internal/converter/postgres/sync_tableddl.go | 78 ++++++++++++- .../converter/postgres/sync_tableddl_test.go | 108 ++++++++++++++++++ 2 files changed, 181 insertions(+), 5 deletions(-) diff --git a/internal/converter/postgres/sync_tableddl.go b/internal/converter/postgres/sync_tableddl.go index 91229f5..0857aeb 100644 --- a/internal/converter/postgres/sync_tableddl.go +++ b/internal/converter/postgres/sync_tableddl.go @@ -68,8 +68,15 @@ var ( reBasicTypes = regexp.MustCompile(`(?i)\b(bigint|integer|smallint|int|bigserial|serial|boolean|text|bytea|timestamp|date|time|decimal|double precision|real)\b`) // 表相关正则 - reComment = regexp.MustCompile(`(?i)\s+comment\s+'((?:[^']|'')*)'\s*,?\s*|\s+comment\s+"([^"]*)"\s*,?\s*`) - reTableComment = regexp.MustCompile(`(?i)\s+COMMENT\s*=\s*'([^']*)'`) + // 注释字面量必须同时支持 MySQL 的两种转义形态: + // - 反斜杠转义 \'(SHOW CREATE TABLE 的默认输出,除非 sql_mode 含 NO_BACKSLASH_ESCAPES) + // - 双写转义 '' + // 旧正则 '((?:[^']|'')*)' 不认 \',匹配会在 user\ 处提前闭合,造成两个后果: + // 残留的 s note' 文本进入列定义使 PG 报 42601 syntax error at or near "s", + // 且该列注释被截断为 user\(DEFAULT NULL 也随之丢失)。转换器却返回 err=nil。 + // 捕获到的内容仍是 MySQL 转义形态,需经 unescapeMySQLStringLiteral 还原为原始文本。 + reComment = regexp.MustCompile(`(?i)\s+comment\s+'((?:[^'\\]|''|\\.)*)'\s*,?\s*|\s+comment\s+"((?:[^"\\]|\\.)*)"\s*,?\s*`) + reTableComment = regexp.MustCompile(`(?i)\s+COMMENT\s*=\s*'((?:[^'\\]|''|\\.)*)'`) // 索引相关正则 reIndexPattern = regexp.MustCompile(`(?i)^(UNIQUE\s+)?(FULLTEXT\s+)?(KEY|INDEX)\s+`) @@ -263,7 +270,8 @@ func parseTableInfo(mysqlDDL string) (tableName string, isTemporary bool, tableC tableComment = "" tableCommentMatch := reTableComment.FindStringSubmatch(mysqlDDL) if tableCommentMatch != nil { - tableComment = tableCommentMatch[1] + // 同列注释:还原 MySQL 转义,PG 侧转义由下游负责 + tableComment = unescapeMySQLStringLiteral(tableCommentMatch[1]) } var bracketCount int @@ -732,6 +740,64 @@ func normalizePartitionBound(bound string) string { return trimmed } +// unescapeMySQLStringLiteral 把 MySQL 字符串字面量中的转义序列还原为原始文本。 +// +// MySQL 的 SHOW CREATE TABLE 用反斜杠转义(COMMENT 'user\'s note'),而 PostgreSQL +// 在 standard_conforming_strings=on(默认)下不认反斜杠转义,因此必须先还原为原始内容, +// 再由下游生成端按 PG 规则重新转义(GenerateColumnCommentsSQL 与 processComment 均已 +// 把 ' 转为 ”)。此处刻意不做 PG 转义,避免二次转义。 +// +// 还原规则遵循 MySQL 字符串字面量语义: +// - \' \" \\ → 对应字符 +// - ” → '(双写转义,即 NO_BACKSLASH_ESCAPES 模式下的形态) +// - \n \r \t \b → 对应控制符 +// - \0 → 丢弃(PostgreSQL 文本类型不允许 NUL 字节) +// - \% \_ → 保留反斜杠(MySQL 对这两个转义保留反斜杠本身) +// - 其他 \x → 反斜杠被忽略、取字符本身 +// (MySQL 文档:For all other escape sequences, backslash is ignored) +// +// 按字节遍历对 UTF-8 安全:多字节序列的每个字节都 >= 0x80,不会与 \ (0x5C) +// 或 ' (0x27) 冲突。 +func unescapeMySQLStringLiteral(s string) string { + if !strings.Contains(s, `\`) && !strings.Contains(s, "''") { + return s + } + + var b strings.Builder + b.Grow(len(s)) + for i := 0; i < len(s); i++ { + c := s[i] + if c == '\\' && i+1 < len(s) { + i++ + switch s[i] { + case 'n': + b.WriteByte('\n') + case 'r': + b.WriteByte('\r') + case 't': + b.WriteByte('\t') + case 'b': + b.WriteByte('\b') + case '0': + // PostgreSQL 文本类型不接受 NUL 字节,丢弃 + case '%', '_': + b.WriteByte('\\') + b.WriteByte(s[i]) + default: + b.WriteByte(s[i]) + } + continue + } + if c == '\'' && i+1 < len(s) && s[i+1] == '\'' { + b.WriteByte('\'') + i++ + continue + } + b.WriteByte(c) + } + return b.String() +} + // toLowerOutsideQuotes 将字符串中非引号包裹内容转换为小写 func toLowerOutsideQuotes(input string) string { var builder strings.Builder @@ -1238,10 +1304,12 @@ func processColumnDefinition(line string, lowercaseColumns bool) (columnName str commentMatch := reComment.FindStringSubmatch(line) if commentMatch != nil { + // 捕获内容仍是 MySQL 转义形态,须还原为原始文本; + // PG 侧的 ' → '' 转义由下游 GenerateColumnCommentsSQL / processComment 负责 if commentMatch[1] != "" { - columnComment = commentMatch[1] + columnComment = unescapeMySQLStringLiteral(commentMatch[1]) } else { - columnComment = commentMatch[2] + columnComment = unescapeMySQLStringLiteral(commentMatch[2]) } } line = reComment.ReplaceAllString(line, "") diff --git a/internal/converter/postgres/sync_tableddl_test.go b/internal/converter/postgres/sync_tableddl_test.go index 4ad4ea4..34538d7 100644 --- a/internal/converter/postgres/sync_tableddl_test.go +++ b/internal/converter/postgres/sync_tableddl_test.go @@ -756,6 +756,114 @@ func TestConvertTableDDL_BitLiteralDefaultKeepsStringLiterals(t *testing.T) { } } +// TestUnescapeMySQLStringLiteral MySQL 字符串字面量转义的还原规则。 +// MySQL 的 SHOW CREATE TABLE 用反斜杠转义(COMMENT 'user\'s note'), +// 而 PG 在 standard_conforming_strings=on 下不认反斜杠转义,故须先还原为原始文本, +// 再由下游生成端按 PG 规则把 ' 转为 ”。 +func TestUnescapeMySQLStringLiteral(t *testing.T) { + tests := []struct { + name string + input string + want string + }{ + {"无转义原样返回", "plain comment", "plain comment"}, + {"空串", "", ""}, + {"反斜杠单引号", `user\'s note`, "user's note"}, + {"双写单引号", "it''s ok", "it's ok"}, + {"双反斜杠", `path\\to`, `path\to`}, + {"换行与制表", `line1\nline2\ttab`, "line1\nline2\ttab"}, + {"回车与退格", `a\rb\bc`, "a\rb\bc"}, + {"NUL 被丢弃", `a\0b`, "ab"}, + {"百分号与下划线保留反斜杠", `100\% and \_x`, `100\% and \_x`}, + {"未知转义忽略反斜杠", `\z`, "z"}, + {"双引号转义", `say \"hi\"`, `say "hi"`}, + {"中文与转义混排", `用户\'s 备注`, "用户's 备注"}, + {"末尾孤立反斜杠原样保留", `trailing\`, `trailing\`}, + {"多个转义共存", `a\'b\\c\nd`, "a'b\\c\nd"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := unescapeMySQLStringLiteral(tt.input); got != tt.want { + t.Errorf("unescapeMySQLStringLiteral(%q) = %q, want %q", tt.input, got, tt.want) + } + }) + } +} + +// TestConvertTableDDL_CommentBackslashEscape 含反斜杠转义撇号的注释不得破坏 DDL 结构。 +// +// 修复前 reComment 只认 ” 双写转义,匹配会在 user\ 处提前闭合,导致: +// - 残留 s note' 使 PG 报 42601 syntax error at or near "s" +// - 该列的 DEFAULT NULL 消失 +// - ColumnComments 中 note 被截断为 "user\",TableComment 被截断为 "table\" +// - 而转换器返回 err=nil、Warnings=[],完全静默 +// +// 列定义是逐行处理的,故后续 plain 列本身不受影响;此处仍断言其完整性以防回归。 +func TestConvertTableDDL_CommentBackslashEscape(t *testing.T) { + mysqlDDL := "CREATE TABLE `t_comment` (\n" + + " `id` int(11) NOT NULL AUTO_INCREMENT,\n" + + " `note` varchar(100) DEFAULT NULL COMMENT 'user\\'s note',\n" + + " `plain` varchar(50) DEFAULT NULL COMMENT 'normal comment',\n" + + " PRIMARY KEY (`id`)\n" + + ") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='table\\'s comment'" + + result, err := ConvertTableDDL(mysqlDDL, true) + if err != nil { + t.Fatalf("ConvertTableDDL failed: %v", err) + } + + // 1. 不得残留被截断的注释文本 + for _, bad := range []string{"s note'", `\'`, "note',"} { + if strings.Contains(result.DDL, bad) { + t.Errorf("DDL 不应残留 %q,实际: %s", bad, result.DDL) + } + } + + // 2. 两列的定义都必须完整存在(修复前 plain 的 COMMENT 会被并入 note) + for _, want := range []string{`"note" VARCHAR(100)`, `"plain" VARCHAR(50)`} { + if !strings.Contains(result.DDL, want) { + t.Errorf("DDL 应包含完整的列定义 %q,实际: %s", want, result.DDL) + } + } + + // 3. 括号必须平衡(结构未被破坏) + if openCnt, closeCnt := strings.Count(result.DDL, "("), strings.Count(result.DDL, ")"); openCnt != closeCnt { + t.Errorf("DDL 括号不平衡: ( = %d, ) = %d,实际: %s", openCnt, closeCnt, result.DDL) + } + + // 4. 注释内容须还原为原始文本(PG 侧转义由下游生成端负责) + if got := result.ColumnComments["note"]; got != "user's note" { + t.Errorf("ColumnComments[note] = %q, want %q", got, "user's note") + } + if got := result.ColumnComments["plain"]; got != "normal comment" { + t.Errorf("ColumnComments[plain] = %q, want %q", got, "normal comment") + } + if result.TableComment != "table's comment" { + t.Errorf("TableComment = %q, want %q", result.TableComment, "table's comment") + } +} + +// TestGenerateColumnCommentsSQL_EscapesSingleQuote 还原后的注释在生成 PG 语句时 +// 必须按 PG 规则把 ' 转为 ”,否则 COMMENT ON 语句本身会语法错误 +func TestGenerateColumnCommentsSQL_EscapesSingleQuote(t *testing.T) { + sqls := GenerateColumnCommentsSQL( + `"t_comment"`, + map[string]string{"note": `"note"`}, + map[string]string{"note": "user's note"}, + ) + + if len(sqls) != 1 { + t.Fatalf("应生成 1 条语句,实际 %d: %v", len(sqls), sqls) + } + if !strings.Contains(sqls[0], "IS 'user''s note'") { + t.Errorf("单引号应按 PG 规则双写,实际: %s", sqls[0]) + } + if strings.Contains(sqls[0], `\'`) { + t.Errorf("不应残留反斜杠转义,实际: %s", sqls[0]) + } +} + // TestCleanTypeDefinition_TinyInt1Mapping tinyint(1) 映射策略(P2-03 + 42883 修复): // 默认映射为 SMALLINT 保留整数语义(兼容视图/函数中 `col = 1` 等用法), // 显式开启 tinyInt1AsBoolean 时映射为 BOOLEAN。 From 7b825b7a04eee0e4f953f4a00574f74e80cf9c4e Mon Sep 17 00:00:00 2001 From: xiaoxu Date: Wed, 23 Sep 2026 16:21:53 +0800 Subject: [PATCH 10/10] =?UTF-8?q?fix:=20reSetVar=20=E7=A0=B4=E5=9D=8F?= =?UTF-8?q?=E5=87=BD=E6=95=B0=E4=BD=93=E4=B8=AD=E7=9A=84=20UPDATE...SET?= =?UTF-8?q?=EF=BC=8C=E4=BA=A7=E5=87=BA=E9=9D=9E=E6=B3=95=20SQL=20=E8=87=B4?= =?UTF-8?q?=20PG=2042601?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 注:本项原计划先建 issue,但 issue 创建被环境策略拦截,故将完整问题描述 保留在此 commit message 中。 问题 ---- reSetVar = (?i)\bSET\s+(\w+)\s*=\s* → "$1 := " 它的本意是转换 plpgsql 的变量赋值(SET v = 1 → v := 1),但无法区分 UPDATE 语句的 SET 子句,会把 SET 整个删掉并把 = 改成 :=。 实证(输入为 SHOW CREATE FUNCTION 的真实形态): UPDATE orders SET status = 1 WHERE id = p_id; → UPDATE orders status := 1 WHERE id = p_id; UPDATE orders SET status = 1, amount = 0 WHERE id = p_id; → UPDATE orders status := 1, amount = 0 WHERE id = p_id; ← 首列被改、次列保留 SET v = p_id + 1; → v := p_id + 1; ← 这是唯一正确的目标场景 ConvertFunctionDDL 返回 err = nil。真实 PostgreSQL 上执行: ERROR: 42601: syntax error at or near ":=" LINE 8: UPDATE orders status := 1 WHERE id = p_id; plpgsql 在 CREATE FUNCTION 时就预编译函数体(check_function_bodies=on), 因此含 UPDATE 的函数一律创建失败。而 CI 的全部 10 个组合都是 functions:false, 这条路径在集成层面零覆盖;create_function.sql 中 grep 'UPDATE\s+\w+\s+SET' 零命中,单测也覆盖不到。 三条既有"修补规则"为何都救不回来 -------------------------------- - reUpdateThen / reUpdateThenEq(:81-82)针对的是 `UPDATE x THEN y :=` 形态, 而 reSetVar 实际产出的是 `UPDATE x y :=`(没有 THEN),故永不匹配; - reUpdateSet 原先是 applyMiscFixes 内的局部变量,模式 `UPDATE\s+(\w+)\s+SET\s+` → `UPDATE $1 SET ` 是恒等替换,且每次调用都 重新编译正则。本次将其移到包级并加注释说明其真实作用仅为空白规范化。 这正是"正则修补正则产物"堆叠的典型后果:下游补丁假设的中间形态与上游 实际产出的形态不一致,规则之间不可组合。 修复 ---- Go 的 RE2 不支持后顾断言,无法用 (?