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
7 changes: 6 additions & 1 deletion config.example.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -69,6 +72,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:
Expand Down
30 changes: 30 additions & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 运行配置
Expand Down Expand Up @@ -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 == "" {
Expand All @@ -263,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
}
94 changes: 94 additions & 0 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package config

import (
"strings"
"testing"
)

Expand Down Expand Up @@ -91,3 +92,96 @@ 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)
}
})
}
}

// 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)
}
})
}
}
37 changes: 27 additions & 10 deletions internal/converter/postgres/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"fmt"
"os"
"path/filepath"
"runtime/debug"
"strings"
"sync"
"sync/atomic"
Expand Down Expand Up @@ -81,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 评估结果
Expand Down Expand Up @@ -565,7 +568,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
}
Expand Down Expand Up @@ -1506,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,
Expand Down Expand Up @@ -1671,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)
}
Expand All @@ -1699,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)
}
Expand Down
120 changes: 120 additions & 0 deletions internal/converter/postgres/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -380,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 应被复位")
}
}
Loading
Loading