feat(sync): 数据同步支持差异对比、行级选择与实时进度日志

- 新增差异分析/预览接口与前端预览抽屉(插入/更新/删除)
  - 支持按表勾选插入/更新/删除(删除默认不勾选)
  - 支持按主键选择行级同步;无主键/复合主键表跳过并提示
  - 同步过程实时输出中文日志与进度条,便于定位失败步骤
This commit is contained in:
杨国锋
2026-02-03 17:37:41 +08:00
parent 80fbfd6365
commit e3bf160072
24 changed files with 2087 additions and 280 deletions

View File

@@ -88,19 +88,7 @@ func (c *CustomDB) Query(query string) ([]map[string]interface{}, []string, erro
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}

View File

@@ -119,19 +119,7 @@ func (d *DamengDB) Query(query string) ([]map[string]interface{}, []string, erro
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}

View File

@@ -150,19 +150,7 @@ func (k *KingbaseDB) Query(query string) ([]map[string]interface{}, []string, er
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}

View File

@@ -48,7 +48,7 @@ func (m *MySQLDB) Connect(config connection.ConnectionConfig) error {
}
m.conn = db
m.pingTimeout = getConnectTimeout(config)
// Force verification
if err := m.Ping(); err != nil {
return fmt.Errorf("连接建立后验证失败:%w", err)
@@ -107,19 +107,7 @@ func (m *MySQLDB) Query(query string) ([]map[string]interface{}, []string, error
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}
@@ -159,12 +147,12 @@ func (m *MySQLDB) GetTables(dbName string) ([]string, error) {
if dbName != "" {
query = fmt.Sprintf("SHOW TABLES FROM `%s`", dbName)
}
data, _, err := m.Query(query)
if err != nil {
return nil, err
}
var tables []string
for _, row := range data {
for _, v := range row {
@@ -185,7 +173,7 @@ func (m *MySQLDB) GetCreateStatement(dbName, tableName string) (string, error) {
if err != nil {
return "", err
}
if len(data) > 0 {
if val, ok := data[0]["Create Table"]; ok {
return fmt.Sprintf("%v", val), nil
@@ -215,12 +203,12 @@ func (m *MySQLDB) GetColumns(dbName, tableName string) ([]connection.ColumnDefin
Extra: fmt.Sprintf("%v", row["Extra"]),
Comment: fmt.Sprintf("%v", row["Comment"]),
}
if row["Default"] != nil {
d := fmt.Sprintf("%v", row["Default"])
col.Default = &d
}
columns = append(columns, col)
}
return columns, nil
@@ -248,14 +236,14 @@ func (m *MySQLDB) GetIndexes(dbName, tableName string) ([]connection.IndexDefini
}
}
seq := 0
if val, ok := row["Seq_in_index"]; ok {
seq := 0
if val, ok := row["Seq_in_index"]; ok {
if f, ok := val.(float64); ok {
seq = int(f)
} else if i, ok := val.(int64); ok {
seq = int(i)
}
}
}
idx := connection.IndexDefinition{
Name: fmt.Sprintf("%v", row["Key_name"]),
@@ -345,12 +333,12 @@ func (m *MySQLDB) ApplyChanges(tableName string, changes connection.ChangeSet) e
for _, update := range changes.Updates {
var sets []string
var args []interface{}
for k, v := range update.Values {
sets = append(sets, fmt.Sprintf("`%s` = ?", k))
args = append(args, v)
}
if len(sets) == 0 {
continue
}
@@ -360,7 +348,7 @@ func (m *MySQLDB) ApplyChanges(tableName string, changes connection.ChangeSet) e
wheres = append(wheres, fmt.Sprintf("`%s` = ?", k))
args = append(args, v)
}
if len(wheres) == 0 {
return fmt.Errorf("update requires keys")
}
@@ -376,13 +364,13 @@ func (m *MySQLDB) ApplyChanges(tableName string, changes connection.ChangeSet) e
var cols []string
var placeholders []string
var args []interface{}
for k, v := range row {
cols = append(cols, fmt.Sprintf("`%s`", k))
placeholders = append(placeholders, "?")
args = append(args, v)
}
if len(cols) == 0 {
continue
}

View File

@@ -125,19 +125,7 @@ func (o *OracleDB) Query(query string) ([]map[string]interface{}, []string, erro
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}

View File

@@ -48,7 +48,7 @@ func (p *PostgresDB) Connect(config connection.ConnectionConfig) error {
}
p.conn = db
p.pingTimeout = getConnectTimeout(config)
// Force verification
if err := p.Ping(); err != nil {
return fmt.Errorf("连接建立后验证失败:%w", err)
@@ -81,8 +81,7 @@ func (p *PostgresDB) Query(query string) ([]map[string]interface{}, []string, er
return nil, nil, fmt.Errorf("connection not open")
}
rows, err := p.conn.Query(query)
rows, err := p.conn.Query(query)
if err != nil {
return nil, nil, err
}
@@ -108,19 +107,7 @@ rows, err := p.conn.Query(query)
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}
@@ -159,7 +146,7 @@ func (p *PostgresDB) GetTables(dbName string) ([]string, error) {
if err != nil {
return nil, err
}
var tables []string
for _, row := range data {
schema, okSchema := row["schemaname"]

View File

@@ -0,0 +1,58 @@
package db
import (
"encoding/hex"
"unicode"
"unicode/utf8"
)
// normalizeQueryValue normalizes driver-returned values for UI/JSON transport.
// 当前主要处理 []byte如果是可读文本则转为 string否则转为十六进制字符串避免前端出现“空白值”。
func normalizeQueryValue(v interface{}) interface{} {
if b, ok := v.([]byte); ok {
return bytesToReadableString(b)
}
return v
}
func bytesToReadableString(b []byte) interface{} {
if b == nil {
return nil
}
if len(b) == 0 {
return ""
}
if utf8.Valid(b) {
s := string(b)
if isMostlyPrintable(s) {
return s
}
}
return "0x" + hex.EncodeToString(b)
}
func isMostlyPrintable(s string) bool {
if s == "" {
return true
}
total := 0
printable := 0
for _, r := range s {
total++
switch r {
case '\n', '\r', '\t':
printable++
continue
default:
}
if unicode.IsPrint(r) {
printable++
}
}
// 允许少量不可见字符,避免把正常文本误判为二进制。
return printable*100 >= total*90
}

View File

@@ -17,14 +17,14 @@ type SQLiteDB struct {
}
func (s *SQLiteDB) Connect(config connection.ConnectionConfig) error {
dsn := config.Host
dsn := config.Host
db, err := sql.Open("sqlite", dsn)
if err != nil {
return fmt.Errorf("打开数据库连接失败:%w", err)
}
s.conn = db
s.pingTimeout = getConnectTimeout(config)
// Force verification
if err := s.Ping(); err != nil {
return fmt.Errorf("连接建立后验证失败:%w", err)
@@ -83,19 +83,7 @@ func (s *SQLiteDB) Query(query string) ([]map[string]interface{}, []string, erro
entry := make(map[string]interface{})
for i, col := range columns {
var v interface{}
val := values[i]
b, ok := val.([]byte)
if ok {
if b == nil {
v = nil
} else {
v = string(b)
}
} else {
v = val
}
entry[col] = v
entry[col] = normalizeQueryValue(values[i])
}
resultData = append(resultData, entry)
}
@@ -124,7 +112,7 @@ func (s *SQLiteDB) GetTables(dbName string) ([]string, error) {
if err != nil {
return nil, err
}
var tables []string
for _, row := range data {
if val, ok := row["name"]; ok {