363 lines
15 KiB
Java
363 lines
15 KiB
Java
import java.io.InputStream;
|
||
import java.nio.charset.StandardCharsets;
|
||
import java.nio.file.Files;
|
||
import java.nio.file.Path;
|
||
import java.sql.Connection;
|
||
import java.sql.DriverManager;
|
||
import java.sql.PreparedStatement;
|
||
import java.sql.ResultSet;
|
||
import java.sql.Statement;
|
||
import java.util.ArrayList;
|
||
import java.util.HashSet;
|
||
import java.util.LinkedHashMap;
|
||
import java.util.List;
|
||
import java.util.Map;
|
||
import java.util.Properties;
|
||
import java.util.Set;
|
||
|
||
/**
|
||
* 其他数据迁移:base_otherdata 下 64 个字典模块的数据导入 + b_from / b_grade 重建
|
||
*
|
||
* 数据来源:旧库 G3HY2025;目标:新库 FMS(连接信息取 config/dbconfigs)
|
||
*
|
||
* 流程:
|
||
* 1. 从 s_module 读出 base_otherdata 下的数据模块(b_save_table);
|
||
* 2. 逐表导入:旧库读行 → 生成雪花 b_id → 按列映射写入新表;
|
||
* 列映射:同名对应;审计字段 b_inputuser_id→b_created_by 等;b_name_c / b_vessel_name → b_name;
|
||
* b_canuse 固定 1、b_xh 按旧编号顺序生成 10/20/…、b_i18n 留空;
|
||
* 3. b_from / b_grade:先读出旧行 → drop + create(雪花主键)→ 重导 →
|
||
* 转换 b_othercompany.b_from_id / b_type_id 的引用值为新雪花;
|
||
* 4. 输出映射文件 tools/migration/otherdata-id-map.tsv(表名、旧编号、新雪花)留档。
|
||
*
|
||
* 用法(在 fms-api 下运行;整体一个事务,失败全部回滚):
|
||
* java -cp "tools/migration;tools/migration/mssql-jdbc-13.4.0.jre11.jar" OtherDataMigrate --dry-run
|
||
* java -cp "tools/migration;tools/migration/mssql-jdbc-13.4.0.jre11.jar" OtherDataMigrate
|
||
*/
|
||
public final class OtherDataMigrate {
|
||
|
||
/** 与后端 IdGenerator 相同的时间基准和位布局;WorkerId 取 63(迁移专用,后端用 1,互不冲突) */
|
||
private static final long BASE_TIME = 1582136402000L;
|
||
private static final long WORKER_ID = 63L;
|
||
private static final int SEQ_BITS = 6;
|
||
private static final int WORKER_BITS = 6;
|
||
|
||
private static long lastMs = -1L;
|
||
private static int seq = 0;
|
||
|
||
/** 审计字段映射(新列 → 旧列) */
|
||
private static final Map<String, String> COLUMN_ALIAS = Map.of(
|
||
"b_created_by", "b_inputuser_id",
|
||
"b_created_at", "b_inputdatetime",
|
||
"b_updated_by", "b_updateuser_id",
|
||
"b_updated_at", "b_updatedatetime");
|
||
|
||
/** 固定值列(不取旧数据) */
|
||
private static final Set<String> FIXED_COLUMNS = Set.of("b_i18n", "b_canuse", "b_xh");
|
||
|
||
/** 需要雪花化重建的表 → 建表 DDL(列名与旧结构一致,仅主键改雪花) */
|
||
private static final Map<String, String> REBUILD_DDL = Map.of(
|
||
"b_from", "create table dbo.b_from ("
|
||
+ "b_id bigint not null primary key, b_name nvarchar(50) null, b_i18n varchar(150) null, "
|
||
+ "b_canuse tinyint not null default 1, b_xh int not null default 0, b_bz nvarchar(50) null, "
|
||
+ "b_created_by varchar(50) null, b_created_at datetime2 null, "
|
||
+ "b_updated_by varchar(50) null, b_updated_at datetime2 null)",
|
||
"b_grade", "create table dbo.b_grade ("
|
||
+ "b_id bigint not null primary key, b_name nvarchar(50) null, b_i18n varchar(150) null, "
|
||
+ "b_canuse tinyint not null default 1, b_xh int not null default 0, b_bz nvarchar(50) null, "
|
||
+ "b_created_by varchar(50) null, b_created_at datetime2 null, "
|
||
+ "b_updated_by varchar(50) null, b_updated_at datetime2 null)");
|
||
|
||
public static void main(String[] args) throws Exception {
|
||
boolean dryRun = args.length > 0 && "--dry-run".equals(args[0]);
|
||
Path apiRoot = Path.of("").toAbsolutePath().normalize();
|
||
Path mapFile = apiRoot.resolve("tools/migration/otherdata-id-map.tsv");
|
||
|
||
StringBuilder mapText = new StringBuilder("表名\t旧编号\t新ID\n");
|
||
int totalRows = 0;
|
||
|
||
try (Connection old = connect(apiRoot, "G3HY2025");
|
||
Connection fms = connect(apiRoot, "G3HD")) {
|
||
fms.setAutoCommit(false);
|
||
try {
|
||
List<String> tables = loadTables(fms);
|
||
System.out.println("模块数: " + tables.size() + (dryRun ? "(dry-run)" : ""));
|
||
|
||
Map<String, Long> fromRefs = new LinkedHashMap<>();
|
||
Map<String, Long> gradeRefs = new LinkedHashMap<>();
|
||
|
||
for (String table : tables) {
|
||
Map<String, Long> mapping = migrateTable(old, fms, table, mapText, dryRun);
|
||
totalRows += mapping.size();
|
||
if ("b_from".equals(table)) {
|
||
fromRefs.putAll(mapping);
|
||
}
|
||
if ("b_grade".equals(table)) {
|
||
gradeRefs.putAll(mapping);
|
||
}
|
||
System.out.println(" " + pad(table, 28) + "旧库 " + mapping.size() + " 行");
|
||
}
|
||
|
||
if (!dryRun) {
|
||
int n1 = convertRefs(fms, "b_othercompany", "b_from_id", fromRefs);
|
||
int n2 = convertRefs(fms, "b_othercompany", "b_type_id", gradeRefs);
|
||
System.out.println("引用转换: b_othercompany.b_from_id " + n1 + " 行, b_type_id " + n2 + " 行");
|
||
fms.commit();
|
||
System.out.println("完成:共导入 " + totalRows + " 行");
|
||
} else {
|
||
fms.rollback();
|
||
System.out.println("[dry-run] 已回滚,共 " + totalRows + " 行待导入");
|
||
}
|
||
} catch (Exception e) {
|
||
fms.rollback();
|
||
throw e;
|
||
}
|
||
}
|
||
|
||
if (!dryRun) {
|
||
Files.writeString(mapFile, mapText.toString(), StandardCharsets.UTF_8);
|
||
System.out.println("映射已写入: " + mapFile);
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------- 模块清单
|
||
|
||
private static List<String> loadTables(Connection fms) throws Exception {
|
||
List<String> tables = new ArrayList<>();
|
||
try (Statement statement = fms.createStatement();
|
||
ResultSet rows = statement.executeQuery(
|
||
"select b_save_table from dbo.s_module where b_path like '/base/base_otherdata/%' "
|
||
+ "and b_module_type = 'data' order by b_xh, b_id")) {
|
||
while (rows.next()) {
|
||
tables.add(rows.getString(1));
|
||
}
|
||
}
|
||
return tables;
|
||
}
|
||
|
||
// ---------------------------------------------------------------- 单表迁移
|
||
|
||
/** 导入一张表,返回「旧编号 → 新雪花」映射 */
|
||
private static Map<String, Long> migrateTable(Connection old, Connection fms, String table,
|
||
StringBuilder mapText, boolean dryRun) throws Exception {
|
||
// 新表列(跳过 b_id,它由雪花生成)
|
||
record Target(String name, String typeName, String source, Object fixed) {}
|
||
List<Target> targets = new ArrayList<>();
|
||
try (Statement statement = fms.createStatement();
|
||
ResultSet rows = statement.executeQuery(
|
||
"select c.name, ty.name from sys.columns c join sys.types ty on ty.user_type_id = c.user_type_id "
|
||
+ "where c.object_id = object_id('dbo." + table + "') order by c.column_id")) {
|
||
while (rows.next()) {
|
||
String name = rows.getString(1);
|
||
String typeName = rows.getString(2);
|
||
if (!"b_id".equals(name)) {
|
||
targets.add(new Target(name, typeName, null, null));
|
||
}
|
||
}
|
||
}
|
||
|
||
// 旧表列
|
||
Set<String> oldColumns = new HashSet<>();
|
||
try (Statement statement = old.createStatement();
|
||
ResultSet rows = statement.executeQuery(
|
||
"select name from sys.columns where object_id = object_id('dbo." + table + "')")) {
|
||
while (rows.next()) {
|
||
oldColumns.add(rows.getString(1));
|
||
}
|
||
}
|
||
if (oldColumns.isEmpty()) {
|
||
throw new IllegalStateException("旧库缺少表 dbo." + table);
|
||
}
|
||
|
||
// 解析每个新列的取值来源
|
||
List<Target> resolved = new ArrayList<>();
|
||
for (Target target : targets) {
|
||
if (FIXED_COLUMNS.contains(target.name())) {
|
||
resolved.add(new Target(target.name(), target.typeName(), null, "@fixed"));
|
||
} else {
|
||
String source = resolveSource(target.name(), oldColumns);
|
||
resolved.add(new Target(target.name(), target.typeName(), source, null));
|
||
}
|
||
}
|
||
|
||
// 读取旧数据
|
||
StringBuilder select = new StringBuilder("select b_id");
|
||
for (Target target : resolved) {
|
||
if (target.source() != null) {
|
||
select.append(", [").append(target.source()).append("]");
|
||
}
|
||
}
|
||
select.append(" from dbo.[").append(table).append("] order by b_id");
|
||
|
||
List<Object[]> rowsData = new ArrayList<>();
|
||
List<String> oldIds = new ArrayList<>();
|
||
try (Statement statement = old.createStatement();
|
||
ResultSet rows = statement.executeQuery(select.toString())) {
|
||
int columnCount = rows.getMetaData().getColumnCount();
|
||
while (rows.next()) {
|
||
oldIds.add(String.valueOf(rows.getObject(1)));
|
||
Object[] values = new Object[columnCount - 1];
|
||
for (int i = 2; i <= columnCount; i++) {
|
||
values[i - 2] = rows.getObject(i);
|
||
}
|
||
rowsData.add(values);
|
||
}
|
||
}
|
||
|
||
if (dryRun) {
|
||
return toMap(oldIds, rowsData);
|
||
}
|
||
|
||
// 重建表(b_from / b_grade)或清空目标表(其余)
|
||
try (Statement statement = fms.createStatement()) {
|
||
if (REBUILD_DDL.containsKey(table)) {
|
||
statement.execute("drop table if exists dbo." + table);
|
||
statement.execute(REBUILD_DDL.get(table));
|
||
} else {
|
||
statement.execute("delete from dbo." + table);
|
||
}
|
||
}
|
||
|
||
// 插入
|
||
List<String> names = new ArrayList<>();
|
||
for (Target target : resolved) {
|
||
names.add("[" + target.name() + "]");
|
||
}
|
||
String insert = "insert into dbo.[" + table + "] (b_id, " + String.join(", ", names) + ") values (?"
|
||
+ ", ?".repeat(names.size()) + ")";
|
||
Map<String, Long> mapping = new LinkedHashMap<>();
|
||
try (PreparedStatement statement = fms.prepareStatement(insert)) {
|
||
int rowIndex = 0;
|
||
for (int row = 0; row < rowsData.size(); row++) {
|
||
rowIndex++;
|
||
long newId = nextId();
|
||
String oldId = oldIds.get(row);
|
||
mapping.put(oldId, newId);
|
||
mapText.append(table).append('\t').append(oldId).append('\t').append(newId).append('\n');
|
||
|
||
statement.setLong(1, newId);
|
||
Object[] values = rowsData.get(row);
|
||
int parameter = 2;
|
||
int sourceIndex = 0;
|
||
for (Target target : resolved) {
|
||
if (target.source() != null) {
|
||
setValue(statement, parameter, target.typeName(), values[sourceIndex++]);
|
||
} else if ("b_xh".equals(target.name())) {
|
||
statement.setInt(parameter, rowIndex * 10);
|
||
} else if ("b_canuse".equals(target.name())) {
|
||
statement.setInt(parameter, 1);
|
||
} else {
|
||
statement.setObject(parameter, null);
|
||
}
|
||
parameter++;
|
||
}
|
||
statement.addBatch();
|
||
}
|
||
statement.executeBatch();
|
||
}
|
||
return mapping;
|
||
}
|
||
|
||
/** 新列 → 旧列;返回 null 表示没有对应旧列 */
|
||
private static String resolveSource(String newColumn, Set<String> oldColumns) {
|
||
String alias = COLUMN_ALIAS.get(newColumn);
|
||
if (alias != null && oldColumns.contains(alias)) {
|
||
return alias;
|
||
}
|
||
if (oldColumns.contains(newColumn)) {
|
||
return newColumn;
|
||
}
|
||
if ("b_name".equals(newColumn)) {
|
||
if (oldColumns.contains("b_name_c")) {
|
||
return "b_name_c";
|
||
}
|
||
if (oldColumns.contains("b_vessel_name")) {
|
||
return "b_vessel_name";
|
||
}
|
||
}
|
||
return null;
|
||
}
|
||
|
||
private static void setValue(PreparedStatement statement, int index, String typeName, Object value) throws Exception {
|
||
if ("tinyint".equals(typeName)) {
|
||
if (value == null) {
|
||
statement.setInt(index, 0);
|
||
return;
|
||
}
|
||
if (value instanceof Number number) {
|
||
statement.setInt(index, number.intValue() == 0 ? 0 : 1);
|
||
return;
|
||
}
|
||
String text = String.valueOf(value).trim();
|
||
statement.setInt(index, "1".equals(text) || "true".equalsIgnoreCase(text) ? 1 : 0);
|
||
return;
|
||
}
|
||
statement.setObject(index, value);
|
||
}
|
||
|
||
private static Map<String, Long> toMap(List<String> oldIds, List<Object[]> rowsData) {
|
||
Map<String, Long> mapping = new LinkedHashMap<>();
|
||
for (String oldId : oldIds) {
|
||
mapping.put(oldId, 0L);
|
||
}
|
||
return mapping;
|
||
}
|
||
|
||
// ---------------------------------------------------------------- 引用转换
|
||
|
||
private static int convertRefs(Connection fms, String table, String column, Map<String, Long> mapping) throws Exception {
|
||
if (mapping.isEmpty()) {
|
||
return 0;
|
||
}
|
||
int count = 0;
|
||
try (PreparedStatement statement = fms.prepareStatement(
|
||
"update dbo." + table + " set " + column + " = ? where " + column + " = ?")) {
|
||
for (Map.Entry<String, Long> entry : mapping.entrySet()) {
|
||
statement.setString(1, String.valueOf(entry.getValue()));
|
||
statement.setString(2, entry.getKey());
|
||
count += statement.executeUpdate();
|
||
}
|
||
}
|
||
return count;
|
||
}
|
||
|
||
// ---------------------------------------------------------------- 基础设施
|
||
|
||
private static Connection connect(Path apiRoot, String org) throws Exception {
|
||
Properties properties = new Properties();
|
||
try (InputStream input = Files.newInputStream(apiRoot.resolve("config/dbconfigs/" + org + ".properties"))) {
|
||
properties.load(input);
|
||
}
|
||
return DriverManager.getConnection(
|
||
properties.getProperty("url"),
|
||
properties.getProperty("username"),
|
||
properties.getProperty("password"));
|
||
}
|
||
|
||
/** 简易雪花:与后端同基准/位布局,同一毫秒内序列递增,回拨时沿用上次时间 */
|
||
private static synchronized long nextId() {
|
||
long now = System.currentTimeMillis();
|
||
if (now < lastMs) {
|
||
now = lastMs;
|
||
}
|
||
if (now == lastMs) {
|
||
seq++;
|
||
if (seq > 63) {
|
||
now++;
|
||
seq = 0;
|
||
}
|
||
} else {
|
||
seq = 0;
|
||
}
|
||
lastMs = now;
|
||
return ((now - BASE_TIME) << (WORKER_BITS + SEQ_BITS)) | (WORKER_ID << SEQ_BITS) | seq;
|
||
}
|
||
|
||
private static String pad(String value, int width) {
|
||
StringBuilder text = new StringBuilder(value == null ? "" : value);
|
||
while (text.length() < width) {
|
||
text.append(' ');
|
||
}
|
||
return text.toString();
|
||
}
|
||
}
|