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
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@

# Package Files #
*.jar
# 迁移工具发行包依赖这些内置插件完成 V2 插件元数据升级,必须随源码一起提交。
!plugins/*.jar
*.war
*.nar
*.ear
Expand Down
45 changes: 40 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,17 +1,52 @@
# DataEase 数据迁移工具

浏览器访问 `http://localhost:8080`,填写 DataEase 2.0 和 3.0 的 SSH、MySQL 及安装目录信息后执行迁移
浏览器访问 `http://localhost:9888`,填写 DataEase 2.0 和 3.0 的服务器地址、MySQL 及安装目录信息后执行完整迁移。远程服务器通过 SSH 操作,本地服务器直接操作本地目录

迁移顺序为:

1. 将源端 `data/i18n`、`font`、`exportData`、`map`、`geo`、`appearance` `static-resource` 打包并复制至目标端的 `data` 目录,并将发行包的三个内置插件 JAR 复制至目标端 `data/plugin`。
2. 通过 JAR 内置的 MySQL Connector/J 读取源库的结构、数据、视图、存储过程、函数、触发器及事件。
3. 删除目标端 JDBC URL 指向的数据库,使用 `utf8mb4` 与 `utf8mb4_0900_ai_ci` 重建,并通过 JDBC 写入迁移内容
1. 将源端实际存在的 `data/i18n`、`font`、`exportData`、`map`、`geo`、`appearance``static-resource`、`excel` 打包并复制至目标端的 `data` 目录,并将发行包内置的 V3 插件 JAR 复制至目标端 `data/plugin`。启用同步日志开关时,还会迁移 `logs/sync-task/task-handler-log`。不存在的可选目录会跳过;所有候选数据目录都不存在时任务会中止,避免错误安装路径产生空迁移
2. 优先使用发行包中与当前平台匹配的 `mysql`、`mysqldump` 迁移数据库;未找到可执行工具时,自动回退到 JAR 内置的 MySQL Connector/J,迁移表结构、数据、视图、存储过程、函数、触发器及事件。
3. 删除目标端 JDBC URL 指向的数据库,使用 `utf8mb4` 与 `utf8mb4_0900_ai_ci` 重建并写入迁移内容
4. 在目标端数据库执行内置的 `upgrade.sql` 升级脚本。
5. 按名称更新飞书多维表格、Apache Hive、达梦插件数据。源端存在其他插件时,任务日志会提示需升级的插件名称。
5. 在事务中按模块或名称兼容更新飞书多维表格、Apache Hive、达梦以及同步管理 PostgreSQL 源/目标插件数据。源端存在其他插件时,任务日志会提示需升级的插件名称。

仅支持 MySQL JDBC URL,例如 `jdbc:mysql://127.0.0.1:3306/dataease`。运行 JAR 的机器必须可通过 JDBC URL 直连源端和目标端 MySQL。填写的数据库用户需具备源库读取定义和数据的权限,以及目标库 `DROP`、`CREATE`、写入权限和执行升级脚本中 `ALTER`、`UPDATE`、`INSERT`、`DELETE` 等语句的权限。

源端和目标端安装目录必须填写对应服务器上的非根目录绝对路径;远程文件操作面向 Linux,本地文件操作支持 macOS/Linux。为兼容迁移工具回退到内置 JDBC 的情况,请预先创建一个可连接的全新空目标数据库;迁移开始后该数据库仍会被删除并重建。迁移期间应停止源 V2 的业务写入,并必须停止目标 V3 服务,避免文件快照、数据库数据或目标业务表被并发修改。

迁移任一阶段发生错误都会终止整个任务。页面日志会显示失败阶段、异常类型和底层根因,服务端日志会记录完整异常堆栈。同步任务参数为空、JSON 无效或缺少源/目标数据源对象时,日志还会逐条列出异常任务 ID、名称和具体原因,但不会输出可能包含数据库密码的任务参数原文。失败不会自动回滚已经完成的文件复制、目标库重建或 MySQL DDL,也不会修改 V2 源库;请先按日志修复 V2 数据,再使用全新目标数据库重新执行完整迁移,不要在失败目标库上继续补跑脚本。

## 本地迁移

把源端或目标端服务器地址填写为 `localhost`、`127.x`、`::1` 或运行迁移程序这台机器的网卡地址后,工具会直接归档、解压对应的本地安装目录并复制插件 JAR,不会建立 SSH 连接;SSH 端口、用户名和密码可以留空。远程地址仍会通过 SSH 执行相同操作,并要求填写有效的 SSH 配置。

源端和目标端会分别判断,因此支持本地到本地、本地到远程、远程到本地和远程到远程。连接方式不会改变迁移内容,四种组合都会依次执行服务文件迁移、数据库结构及数据复制、`upgrade.sql`、通用插件数据更新、同步管理专项数据转换和同步 PostgreSQL 插件数据更新。直接操作本地文件目前适用于 macOS/Linux,并要求本机提供 `/bin/sh` 和 `tar`。

## 同步任务日志复制开关

同步任务物理日志可能达到数百 MB 甚至更大,因此默认不复制。需要迁移时,在启动 JAR 时添加参数:

```bash
java -jar dataease-migration-1.0.0.jar --migration.files.copy-sync-task-logs=true
```

也可以使用环境变量:

```bash
MIGRATION_COPY_SYNC_TASK_LOGS=true java -jar dataease-migration-1.0.0.jar
```

开启后会把源端 `${安装目录}/logs/sync-task/task-handler-log` 单独归档,并合并解压至目标端同名目录,不会复制其他应用日志。日期目录和以 `per_sync_task_log.id` 命名的 `.log` 文件会保持不变;源端目录不存在时会记录日志并安全跳过。该配置是启动级参数,修改后需要重启迁移程序。

## 本机数据库测试

源库和目标库可以位于同一个本机 MySQL,但数据库名必须不同,例如:

- 源库:`jdbc:mysql://127.0.0.1:3306/dataease_v2_test`
- 目标库:`jdbc:mysql://127.0.0.1:3306/dataease_v3_test`

请先把待测数据导入源库并创建可连接的空目标库。目标库用户需要具备 `DROP`、`CREATE`、写入以及升级脚本所需的 `ALTER`、`UPDATE`、`INSERT`、`DELETE` 权限;执行测试时目标库仍会被删除并重建。若要验证 JDBC 批量复制优化,请确保本地 `tools/mysql/<平台>/bin` 中没有可用的 `mysql` 和 `mysqldump`,否则工具会优先走原生导出/导入路径。

## 自动选择数据库迁移工具

应用启动时按当前操作系统和 CPU 架构,在 `tools/mysql/<平台>/bin` 查找 `mysql` 和 `mysqldump`。两个工具均存在且可执行时,自动使用本地工具迁移;否则自动使用 JAR 内置 JDBC 迁移。
Expand Down
Binary file added plugins/postgresql-backend-sink-3.0.0.jar
Binary file not shown.
Binary file added plugins/postgresql-backend-source-3.0.0.jar
Binary file not shown.
16 changes: 10 additions & 6 deletions src/main/java/com/dataease/migration/model/ServerInfo.java
Original file line number Diff line number Diff line change
@@ -1,15 +1,19 @@
package com.dataease.migration.model;

import jakarta.validation.constraints.Max;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotBlank;

/**
* DataEase 服务端文件位置及连接信息。
*
* <p>host 和 installPath 对本地、远程迁移都必填;SSH 字段只在 host 指向远程机器时使用。
* 因为 Bean Validation 无法根据 host 是否属于本机做条件校验,username、password、port
* 刻意不声明全局非空/范围约束,改由 MigrationService 在任务入队前按连接方式校验。</p>
*/
public record ServerInfo(
@NotBlank(message = "服务器 IP 不能为空") String host,
@NotBlank(message = "服务器用户名不能为空") String username,
@NotBlank(message = "服务器密码不能为空") String password,
@Min(value = 1, message = "SSH 端口必须介于 1 到 65535")
@Max(value = 65535, message = "SSH 端口必须介于 1 到 65535") int port,
String username,
String password,
int port,
@NotBlank(message = "DataEase 安装目录不能为空") String installPath
) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -13,20 +13,37 @@
import java.sql.Statement;
import java.util.ArrayList;
import java.util.List;
import java.util.Locale;
import java.util.Properties;
import java.util.concurrent.TimeUnit;

/**
* 没有可用 mysql/mysqldump 时使用的 JDBC 回退迁移器,负责复制 MySQL 结构和数据。
* 大表数据采用流式读取、批量改写和分段提交,避免一次性加载全表或逐行网络往返。
*/
@Component
public class JdbcDatabaseMigrator implements DatabaseMigrator {
/**
* 单批行数需要同时兼顾吞吐和内存:Connector/J 会将一批改写成多值 INSERT,而同步日志表
* 可能含有较大的 TEXT/LONGTEXT,批次过大会显著抬高迁移进程及 MySQL 的瞬时内存占用。
*/
private static final int BATCH_SIZE = 500;
/**
* 每 20 批提交一次,减少逐批提交的 fsync 开销;用 BATCH_SIZE 计算可保证提交前没有待执行批次。
*/
private static final int COMMIT_INTERVAL = BATCH_SIZE * 20;

@Override
public void migrate(DatabaseInfo source, DatabaseInfo target, MigrationJob job) throws SQLException {
DatabaseConnection sourceConnection = DatabaseConnection.fromJdbcUrl(source.jdbcUrl());
DatabaseConnection targetConnection = DatabaseConnection.fromJdbcUrl(target.jdbcUrl());
try (Connection sourceDb = connect(source); Connection targetDb = connect(target)) {
try (Connection sourceDb = connect(source, false); Connection targetDb = connect(target, true)) {
sourceDb.setAutoCommit(false);
recreateDatabase(targetDb, targetConnection, job);
targetDb.setAutoCommit(false);
execute(targetDb, "SET FOREIGN_KEY_CHECKS = 0");
job.log("JDBC 批量写入已启用:每批 " + BATCH_SIZE + " 行,每 "
+ COMMIT_INTERVAL + " 行提交一次。");
try {
List<String> tables = findObjects(sourceDb, sourceConnection.database(), "BASE TABLE");
job.log("正在迁移 " + tables.size() + " 个数据表结构。");
Expand Down Expand Up @@ -55,8 +72,24 @@ public void migrate(DatabaseInfo source, DatabaseInfo target, MigrationJob job)
}
}

private Connection connect(DatabaseInfo database) throws SQLException {
return DriverManager.getConnection(database.jdbcUrl(), database.username(), database.password());
private Connection connect(DatabaseInfo database, boolean optimizeBatchWrites) throws SQLException {
return DriverManager.getConnection(database.jdbcUrl(), connectionProperties(database, optimizeBatchWrites));
}

/**
* 只为目标连接启用批量改写。源连接负责流式读取,无需写入优化;目标连接关闭服务端预编译后,
* Connector/J 才能稳定地把 executeBatch 改写为一条多值 INSERT。单独传 Properties 还会覆盖
* JDBC URL 中的同名配置,避免页面输入的参数意外关闭迁移工具的批量策略。
*/
static Properties connectionProperties(DatabaseInfo database, boolean optimizeBatchWrites) {
Properties properties = new Properties();
properties.setProperty("user", database.username());
properties.setProperty("password", database.password());
if (optimizeBatchWrites) {
properties.setProperty("rewriteBatchedStatements", "true");
properties.setProperty("useServerPrepStmts", "false");
}
return properties;
}

private void recreateDatabase(Connection target, DatabaseConnection connection, MigrationJob job) throws SQLException {
Expand Down Expand Up @@ -99,31 +132,60 @@ private void copyTableData(Connection source, Connection target, String sourceDa
String insert = "INSERT INTO " + targetTable + " (" + quotedColumns + ") VALUES (" + placeholders + ")";

job.log("正在迁移数据表:" + table + "。");
long startedNanos = System.nanoTime();
try (Statement read = source.createStatement(ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY);
PreparedStatement write = target.prepareStatement(insert)) {
read.setFetchSize(Integer.MIN_VALUE);
try (ResultSet rows = read.executeQuery(select)) {
int rowCount = 0;
int pendingBatchRows = 0;
while (rows.next()) {
for (int index = 0; index < columns.size(); index++) {
write.setObject(index + 1, rows.getObject(index + 1));
}
write.addBatch();
rowCount++;
if (rowCount % BATCH_SIZE == 0) {
pendingBatchRows++;
if (pendingBatchRows == BATCH_SIZE) {
write.executeBatch();
write.clearBatch();
pendingBatchRows = 0;
}
if (rowCount % COMMIT_INTERVAL == 0) {
// COMMIT_INTERVAL 是 BATCH_SIZE 的整数倍,此时所有 addBatch 内容均已落库。
target.commit();
job.log("数据表 " + table + " 已迁移 " + rowCount + " 行,平均 "
+ rowsPerSecond(rowCount, startedNanos) + " 行/秒。");
}
}
if (rowCount % BATCH_SIZE != 0) {
if (pendingBatchRows > 0) {
write.executeBatch();
write.clearBatch();
}
// 整万行已在循环内提交;这里只提交最后不足一个提交区间的尾批,避免空提交。
if (rowCount % COMMIT_INTERVAL != 0) {
target.commit();
}
job.log("数据表 " + table + " 迁移完成,共 " + rowCount + " 行。");
job.log("数据表 " + table + " 迁移完成,共 " + rowCount + " 行,耗时 "
+ formatElapsed(startedNanos) + ",平均 " + rowsPerSecond(rowCount, startedNanos)
+ " 行/秒。");
}
}
}

private long rowsPerSecond(int rowCount, long startedNanos) {
long elapsedNanos = Math.max(1, System.nanoTime() - startedNanos);
return Math.round(rowCount * 1_000_000_000.0 / elapsedNanos);
}

private String formatElapsed(long startedNanos) {
long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startedNanos);
if (elapsedMillis < 1_000) {
return elapsedMillis + " 毫秒";
}
return String.format(Locale.ROOT, "%.2f 秒", elapsedMillis / 1_000.0);
}

private List<String> migratableColumns(Connection connection, String database, String table) throws SQLException {
List<String> columns = new ArrayList<>();
try (Statement statement = connection.createStatement();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ public final class MigrationJob {
private final List<SseEmitter> subscribers = new CopyOnWriteArrayList<>();
private final Instant createdAt = Instant.now();
private boolean completed;
private boolean succeeded;

public void log(String message) {
String line = "[" + Instant.now() + "] " + message;
Expand All @@ -40,12 +41,13 @@ public void subscribe(SseEmitter emitter) {
complete(emitter);
}

public void complete() {
public void complete(boolean succeeded) {
List<SseEmitter> activeSubscribers;
synchronized (this) {
if (completed) {
return;
}
this.succeeded = succeeded;
completed = true;
activeSubscribers = new ArrayList<>(subscribers);
}
Expand All @@ -66,7 +68,10 @@ private void send(SseEmitter emitter, String name, String data) {

private void complete(SseEmitter emitter) {
try {
emitter.send(SseEmitter.event().name("complete").data("迁移任务结束"));
String message = succeeded
? "迁移任务成功完成"
: "迁移任务已失败,请根据上方异常明细处理后,使用全新目标库重新迁移";
emitter.send(SseEmitter.event().name("complete").data(message));
} catch (IOException ignored) {
// The browser may have disconnected before the final event.
} finally {
Expand Down
Loading