Files
2026-09-14 22:30:13 +08:00

363 lines
15 KiB
Java
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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();
}
}