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 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 FIXED_COLUMNS = Set.of("b_i18n", "b_canuse", "b_xh"); /** 需要雪花化重建的表 → 建表 DDL(列名与旧结构一致,仅主键改雪花) */ private static final Map 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 tables = loadTables(fms); System.out.println("模块数: " + tables.size() + (dryRun ? "(dry-run)" : "")); Map fromRefs = new LinkedHashMap<>(); Map gradeRefs = new LinkedHashMap<>(); for (String table : tables) { Map 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 loadTables(Connection fms) throws Exception { List 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 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 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 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 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 rowsData = new ArrayList<>(); List 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 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 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 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 toMap(List oldIds, List rowsData) { Map mapping = new LinkedHashMap<>(); for (String oldId : oldIds) { mapping.put(oldId, 0L); } return mapping; } // ---------------------------------------------------------------- 引用转换 private static int convertRefs(Connection fms, String table, String column, Map 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 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(); } }