|
@@ -0,0 +1,463 @@
|
|
|
|
|
+package com.jzg.service.impl;
|
|
|
|
|
+
|
|
|
|
|
+import com.jzg.commons.util.idgen.IdGenerate;
|
|
|
|
|
+import com.jzg.mapper.MigrateDataMapper;
|
|
|
|
|
+import com.jzg.service.MigrateDataService;
|
|
|
|
|
+import lombok.extern.slf4j.Slf4j;
|
|
|
|
|
+import org.springframework.beans.factory.annotation.Autowired;
|
|
|
|
|
+import org.springframework.jdbc.core.JdbcTemplate;
|
|
|
|
|
+import org.springframework.stereotype.Service;
|
|
|
|
|
+import org.springframework.transaction.annotation.Transactional;
|
|
|
|
|
+
|
|
|
|
|
+import java.util.*;
|
|
|
|
|
+import java.util.stream.Collectors;
|
|
|
|
|
+
|
|
|
|
|
+/**
|
|
|
|
|
+ * 数据迁移服务实现类
|
|
|
|
|
+ */
|
|
|
|
|
+@Service
|
|
|
|
|
+@Slf4j
|
|
|
|
|
+public class MigrateDataServiceImpl implements MigrateDataService {
|
|
|
|
|
+
|
|
|
|
|
+ @Autowired
|
|
|
|
|
+ private MigrateDataMapper migrateDataMapper;
|
|
|
|
|
+ @Autowired
|
|
|
|
|
+ private JdbcTemplate jdbcTemplate;
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 为数据生成雪花ID
|
|
|
|
|
+ * @param dataMap 数据Map
|
|
|
|
|
+ * @param primaryKeyField 主键字段名
|
|
|
|
|
+ */
|
|
|
|
|
+ private void generateSnowflakeId(Map<String, Object> dataMap, String primaryKeyField) {
|
|
|
|
|
+ if (primaryKeyField != null && !primaryKeyField.isEmpty()) {
|
|
|
|
|
+ String snowflakeId = IdGenerate.nextId();
|
|
|
|
|
+ dataMap.put(primaryKeyField, snowflakeId);
|
|
|
|
|
+ log.debug("为主键字段 {} 生成雪花ID: {}", primaryKeyField, snowflakeId);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 将数据列表中的值转换为合适的类型,避免JSON类型转换问题
|
|
|
|
|
+ * 保留布尔值和数字类型的原始值,只将其他类型转换为字符串
|
|
|
|
|
+ * @param dataList 数据列表
|
|
|
|
|
+ * @return 转换后的数据列表
|
|
|
|
|
+ */
|
|
|
|
|
+ private List<Map<String, Object>> convertValuesToString(List<Map<String, Object>> dataList) {
|
|
|
|
|
+ return dataList.stream()
|
|
|
|
|
+ .map(row -> {
|
|
|
|
|
+ Map<String, Object> newRow = new HashMap<>();
|
|
|
|
|
+ row.forEach((key, value) -> {
|
|
|
|
|
+ if (value == null) {
|
|
|
|
|
+ newRow.put(key, null);
|
|
|
|
|
+ } else if (value instanceof Boolean || value instanceof Number) {
|
|
|
|
|
+ // 保留布尔值和数字类型的原始值
|
|
|
|
|
+ newRow.put(key, value);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ // 其他类型转换为字符串
|
|
|
|
|
+ newRow.put(key, value.toString());
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ return newRow;
|
|
|
|
|
+ })
|
|
|
|
|
+ .collect(Collectors.toList());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional(rollbackFor = Exception.class)
|
|
|
|
|
+ public int migrateData(String sourceTableName, String targetTableName, String whereCondition) {
|
|
|
|
|
+ return migrateData(sourceTableName, targetTableName, whereCondition, null, false, null, false);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 带雪花ID生成的数据迁移方法
|
|
|
|
|
+ */
|
|
|
|
|
+ public int migrateData(String sourceTableName, String targetTableName, String whereCondition,
|
|
|
|
|
+ String primaryKeyField, boolean generateSnowflakeId,
|
|
|
|
|
+ String sourcePrimaryKeyField, boolean saveOriginalId) {
|
|
|
|
|
+ log.info("开始迁移数据: 源表={}, 目标表={}, 条件={}, 主键字段={}, 生成雪花ID={}, 源主键字段={}, 保存原始ID={}",
|
|
|
|
|
+ sourceTableName, targetTableName, whereCondition, primaryKeyField, generateSnowflakeId,
|
|
|
|
|
+ sourcePrimaryKeyField, saveOriginalId);
|
|
|
|
|
+
|
|
|
|
|
+ // 检查源表和目标表是否存在
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(sourceTableName)) {
|
|
|
|
|
+ throw new RuntimeException("源表 " + sourceTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(targetTableName)) {
|
|
|
|
|
+ throw new RuntimeException("目标表 " + targetTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要保存原始ID,检查并添加original_id列
|
|
|
|
|
+ if (saveOriginalId && sourcePrimaryKeyField != null && !sourcePrimaryKeyField.isEmpty()) {
|
|
|
|
|
+ if (!migrateDataMapper.checkColumnExists(targetTableName, "original_id")) {
|
|
|
|
|
+ migrateDataMapper.addOriginalIdColumn(targetTableName);
|
|
|
|
|
+ log.info("已为目标表 {} 添加 original_id 列", targetTableName);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 从源表查询数据
|
|
|
|
|
+ List<Map<String, Object>> sourceData = migrateDataMapper.selectFromSourceTable(sourceTableName, whereCondition);
|
|
|
|
|
+ if (sourceData == null || sourceData.isEmpty()) {
|
|
|
|
|
+ log.info("源表中没有符合条件的数据");
|
|
|
|
|
+ return 0;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 获取源表结构,确保所有字段都被包含
|
|
|
|
|
+ List<Map<String, Object>> sourceTableStructure = migrateDataMapper.getTableStructure(sourceTableName);
|
|
|
|
|
+ Set<String> allColumns = sourceTableStructure.stream()
|
|
|
|
|
+ .map(col -> col.get("columnName").toString())
|
|
|
|
|
+ .collect(Collectors.toSet());
|
|
|
|
|
+
|
|
|
|
|
+ // 确保每条记录都包含所有字段(即使值为null)
|
|
|
|
|
+ for (Map<String, Object> row : sourceData) {
|
|
|
|
|
+ for (String column : allColumns) {
|
|
|
|
|
+ if (!row.containsKey(column)) {
|
|
|
|
|
+ row.put(column, null);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 将所有值转换为String类型,避免JSON类型转换问题
|
|
|
|
|
+ sourceData = convertValuesToString(sourceData);
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要保存原始ID,为每条记录添加original_id字段
|
|
|
|
|
+ if (saveOriginalId && sourcePrimaryKeyField != null && !sourcePrimaryKeyField.isEmpty()) {
|
|
|
|
|
+ sourceData.forEach(row -> {
|
|
|
|
|
+ Object originalId = row.get(sourcePrimaryKeyField);
|
|
|
|
|
+ if (originalId != null) {
|
|
|
|
|
+ row.put("original_id", originalId.toString());
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要生成雪花ID,则为每条记录生成
|
|
|
|
|
+ if (generateSnowflakeId && primaryKeyField != null && !primaryKeyField.isEmpty()) {
|
|
|
|
|
+ sourceData.forEach(row -> generateSnowflakeId(row, primaryKeyField));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 批量插入到目标表
|
|
|
|
|
+ int batchSize = 1000;
|
|
|
|
|
+ int totalMigrated = 0;
|
|
|
|
|
+
|
|
|
|
|
+ for (int i = 0; i < sourceData.size(); i += batchSize) {
|
|
|
|
|
+ int end = Math.min(i + batchSize, sourceData.size());
|
|
|
|
|
+ List<Map<String, Object>> batch = sourceData.subList(i, end);
|
|
|
|
|
+
|
|
|
|
|
+ // 使用JdbcTemplate执行批量插入
|
|
|
|
|
+ if (!batch.isEmpty()) {
|
|
|
|
|
+ // 获取列名
|
|
|
|
|
+ Set<String> columns = batch.get(0).keySet();
|
|
|
|
|
+ String columnNames = String.join(", ", columns);
|
|
|
|
|
+
|
|
|
|
|
+ // 构建占位符
|
|
|
|
|
+ String placeholders = String.join(", ", Collections.nCopies(columns.size(), "?"));
|
|
|
|
|
+
|
|
|
|
|
+ // 构建SQL语句
|
|
|
|
|
+ String sql = "INSERT INTO " + targetTableName + " (" + columnNames + ") VALUES (" + placeholders + ")";
|
|
|
|
|
+
|
|
|
|
|
+ // 准备参数
|
|
|
|
|
+ List<Object[]> args = batch.stream()
|
|
|
|
|
+ .map(row -> columns.stream()
|
|
|
|
|
+ .map(col -> row.get(col))
|
|
|
|
|
+ .toArray())
|
|
|
|
|
+ .collect(Collectors.toList());
|
|
|
|
|
+
|
|
|
|
|
+ // 执行批量插入
|
|
|
|
|
+ int[] results = jdbcTemplate.batchUpdate(sql, args);
|
|
|
|
|
+ totalMigrated += Arrays.stream(results).sum();
|
|
|
|
|
+ log.info("已迁移 {}/{} 条记录", totalMigrated, sourceData.size());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ log.info("数据迁移完成: 共迁移 {} 条记录", totalMigrated);
|
|
|
|
|
+ return totalMigrated;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional(rollbackFor = Exception.class)
|
|
|
|
|
+ public int migrateDataWithMapping(String sourceTableName, String targetTableName,
|
|
|
|
|
+ Map<String, String> fieldMapping, String whereCondition) {
|
|
|
|
|
+ return migrateDataWithMapping(sourceTableName, targetTableName, fieldMapping, whereCondition, null, false, null, false);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 带雪花ID生成的字段映射迁移方法
|
|
|
|
|
+ */
|
|
|
|
|
+ public int migrateDataWithMapping(String sourceTableName, String targetTableName,
|
|
|
|
|
+ Map<String, String> fieldMapping, String whereCondition,
|
|
|
|
|
+ String primaryKeyField, boolean generateSnowflakeId,
|
|
|
|
|
+ String sourcePrimaryKeyField, boolean saveOriginalId) {
|
|
|
|
|
+ log.info("开始迁移数据(带字段映射): 源表={}, 目标表={}, 映射关系={}, 条件={}, 主键字段={}, 生成雪花ID={}, 源主键字段={}, 保存原始ID={}",
|
|
|
|
|
+ sourceTableName, targetTableName, fieldMapping, whereCondition, primaryKeyField, generateSnowflakeId,
|
|
|
|
|
+ sourcePrimaryKeyField, saveOriginalId);
|
|
|
|
|
+
|
|
|
|
|
+ // 检查源表和目标表是否存在
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(sourceTableName)) {
|
|
|
|
|
+ throw new RuntimeException("源表 " + sourceTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(targetTableName)) {
|
|
|
|
|
+ throw new RuntimeException("目标表 " + targetTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要保存原始ID,检查并添加original_id列
|
|
|
|
|
+ if (saveOriginalId && sourcePrimaryKeyField != null && !sourcePrimaryKeyField.isEmpty()) {
|
|
|
|
|
+ if (!migrateDataMapper.checkColumnExists(targetTableName, "original_id")) {
|
|
|
|
|
+ migrateDataMapper.addOriginalIdColumn(targetTableName);
|
|
|
|
|
+ log.info("已为目标表 {} 添加 original_id 列", targetTableName);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 从源表查询数据
|
|
|
|
|
+ List<Map<String, Object>> sourceData = migrateDataMapper.selectFromSourceTable(sourceTableName, whereCondition);
|
|
|
|
|
+ if (sourceData == null || sourceData.isEmpty()) {
|
|
|
|
|
+ log.info("源表中没有符合条件的数据");
|
|
|
|
|
+ return 0;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 获取源表结构,确保所有字段都被包含
|
|
|
|
|
+ List<Map<String, Object>> sourceTableStructure = migrateDataMapper.getTableStructure(sourceTableName);
|
|
|
|
|
+ Set<String> allColumns = sourceTableStructure.stream()
|
|
|
|
|
+ .map(col -> col.get("columnName").toString())
|
|
|
|
|
+ .collect(Collectors.toSet());
|
|
|
|
|
+
|
|
|
|
|
+ // 确保每条记录都包含所有字段(即使值为null)
|
|
|
|
|
+ for (Map<String, Object> row : sourceData) {
|
|
|
|
|
+ for (String column : allColumns) {
|
|
|
|
|
+ if (!row.containsKey(column)) {
|
|
|
|
|
+ row.put(column, null);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 应用字段映射
|
|
|
|
|
+ List<Map<String, Object>> targetData = sourceData.stream()
|
|
|
|
|
+ .map(sourceRow -> {
|
|
|
|
|
+ Map<String, Object> targetRow = new HashMap<>();
|
|
|
|
|
+ fieldMapping.forEach((sourceField, targetField) -> {
|
|
|
|
|
+ // 直接获取值,包括null值
|
|
|
|
|
+ targetRow.put(targetField, sourceRow.get(sourceField));
|
|
|
|
|
+ });
|
|
|
|
|
+ return targetRow;
|
|
|
|
|
+ })
|
|
|
|
|
+ .collect(Collectors.toList());
|
|
|
|
|
+
|
|
|
|
|
+ // 将所有值转换为String类型,避免JSON类型转换问题
|
|
|
|
|
+ targetData = convertValuesToString(targetData);
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要保存原始ID,为每条记录添加original_id字段
|
|
|
|
|
+ if (saveOriginalId && sourcePrimaryKeyField != null && !sourcePrimaryKeyField.isEmpty()) {
|
|
|
|
|
+ targetData.forEach(row -> {
|
|
|
|
|
+ Object originalId = row.get(sourcePrimaryKeyField);
|
|
|
|
|
+ if (originalId != null) {
|
|
|
|
|
+ row.put("original_id", originalId.toString());
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要生成雪花ID,则为每条记录生成
|
|
|
|
|
+ if (generateSnowflakeId && primaryKeyField != null && !primaryKeyField.isEmpty()) {
|
|
|
|
|
+ targetData.forEach(row -> generateSnowflakeId(row, primaryKeyField));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 批量插入到目标表
|
|
|
|
|
+ int batchSize = 1000;
|
|
|
|
|
+ int totalMigrated = 0;
|
|
|
|
|
+
|
|
|
|
|
+ for (int i = 0; i < targetData.size(); i += batchSize) {
|
|
|
|
|
+ int end = Math.min(i + batchSize, targetData.size());
|
|
|
|
|
+ List<Map<String, Object>> batch = targetData.subList(i, end);
|
|
|
|
|
+
|
|
|
|
|
+ // 使用JdbcTemplate执行批量插入
|
|
|
|
|
+ if (!batch.isEmpty()) {
|
|
|
|
|
+ // 获取列名
|
|
|
|
|
+ Set<String> columns = batch.get(0).keySet();
|
|
|
|
|
+ String columnNames = String.join(", ", columns);
|
|
|
|
|
+
|
|
|
|
|
+ // 构建占位符
|
|
|
|
|
+ String placeholders = String.join(", ", Collections.nCopies(columns.size(), "?"));
|
|
|
|
|
+
|
|
|
|
|
+ // 构建SQL语句
|
|
|
|
|
+ String sql = "INSERT INTO " + targetTableName + " (" + columnNames + ") VALUES (" + placeholders + ")";
|
|
|
|
|
+
|
|
|
|
|
+ // 准备参数
|
|
|
|
|
+ List<Object[]> args = batch.stream()
|
|
|
|
|
+ .map(row -> columns.stream()
|
|
|
|
|
+ .map(col -> row.get(col))
|
|
|
|
|
+ .toArray())
|
|
|
|
|
+ .collect(Collectors.toList());
|
|
|
|
|
+
|
|
|
|
|
+ // 执行批量插入
|
|
|
|
|
+ int[] results = jdbcTemplate.batchUpdate(sql, args);
|
|
|
|
|
+ totalMigrated += Arrays.stream(results).sum();
|
|
|
|
|
+ log.info("已迁移 {}/{} 条记录", totalMigrated, targetData.size());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ log.info("数据迁移完成: 共迁移 {} 条记录", totalMigrated);
|
|
|
|
|
+ return totalMigrated;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional(rollbackFor = Exception.class)
|
|
|
|
|
+ public int migrateDataWithFields(String sourceTableName, String targetTableName,
|
|
|
|
|
+ List<String> fields, String whereCondition) {
|
|
|
|
|
+ return migrateDataWithFields(sourceTableName, targetTableName, fields, whereCondition, null, false, null, false);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 带雪花ID生成的指定字段迁移方法
|
|
|
|
|
+ */
|
|
|
|
|
+ public int migrateDataWithFields(String sourceTableName, String targetTableName,
|
|
|
|
|
+ List<String> fields, String whereCondition,
|
|
|
|
|
+ String primaryKeyField, boolean generateSnowflakeId,
|
|
|
|
|
+ String sourcePrimaryKeyField, boolean saveOriginalId) {
|
|
|
|
|
+ log.info("开始迁移数据(指定字段): 源表={}, 目标表={}, 字段={}, 条件={}, 主键字段={}, 生成雪花ID={}, 源主键字段={}, 保存原始ID={}",
|
|
|
|
|
+ sourceTableName, targetTableName, fields, whereCondition, primaryKeyField, generateSnowflakeId,
|
|
|
|
|
+ sourcePrimaryKeyField, saveOriginalId);
|
|
|
|
|
+
|
|
|
|
|
+ // 检查源表和目标表是否存在
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(sourceTableName)) {
|
|
|
|
|
+ throw new RuntimeException("源表 " + sourceTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(targetTableName)) {
|
|
|
|
|
+ throw new RuntimeException("目标表 " + targetTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要保存原始ID,检查并添加original_id列
|
|
|
|
|
+ if (saveOriginalId && sourcePrimaryKeyField != null && !sourcePrimaryKeyField.isEmpty()) {
|
|
|
|
|
+ if (!migrateDataMapper.checkColumnExists(targetTableName, "original_id")) {
|
|
|
|
|
+ migrateDataMapper.addOriginalIdColumn(targetTableName);
|
|
|
|
|
+ log.info("已为目标表 {} 添加 original_id 列", targetTableName);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 从源表查询数据
|
|
|
|
|
+ List<Map<String, Object>> sourceData = migrateDataMapper.selectFromSourceTable(sourceTableName, whereCondition);
|
|
|
|
|
+ if (sourceData == null || sourceData.isEmpty()) {
|
|
|
|
|
+ log.info("源表中没有符合条件的数据");
|
|
|
|
|
+ return 0;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 获取源表结构,确保所有字段都被包含
|
|
|
|
|
+ List<Map<String, Object>> sourceTableStructure = migrateDataMapper.getTableStructure(sourceTableName);
|
|
|
|
|
+ Set<String> allColumns = sourceTableStructure.stream()
|
|
|
|
|
+ .map(col -> col.get("columnName").toString())
|
|
|
|
|
+ .collect(Collectors.toSet());
|
|
|
|
|
+
|
|
|
|
|
+ // 确保每条记录都包含所有字段(即使值为null)
|
|
|
|
|
+ for (Map<String, Object> row : sourceData) {
|
|
|
|
|
+ for (String column : allColumns) {
|
|
|
|
|
+ if (!row.containsKey(column)) {
|
|
|
|
|
+ row.put(column, null);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 过滤字段
|
|
|
|
|
+ List<Map<String, Object>> targetData = sourceData.stream()
|
|
|
|
|
+ .map(sourceRow -> {
|
|
|
|
|
+ Map<String, Object> targetRow = new HashMap<>();
|
|
|
|
|
+ fields.forEach(field -> {
|
|
|
|
|
+ // 直接获取值,包括null值
|
|
|
|
|
+ targetRow.put(field, sourceRow.get(field));
|
|
|
|
|
+ });
|
|
|
|
|
+ return targetRow;
|
|
|
|
|
+ })
|
|
|
|
|
+ .collect(Collectors.toList());
|
|
|
|
|
+
|
|
|
|
|
+ // 将所有值转换为String类型,避免JSON类型转换问题
|
|
|
|
|
+ targetData = convertValuesToString(targetData);
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要保存原始ID,为每条记录添加original_id字段
|
|
|
|
|
+ if (saveOriginalId && sourcePrimaryKeyField != null && !sourcePrimaryKeyField.isEmpty()) {
|
|
|
|
|
+ targetData.forEach(row -> {
|
|
|
|
|
+ Object originalId = row.get(sourcePrimaryKeyField);
|
|
|
|
|
+ if (originalId != null) {
|
|
|
|
|
+ row.put("original_id", originalId.toString());
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 如果需要生成雪花ID,则为每条记录生成
|
|
|
|
|
+ if (generateSnowflakeId && primaryKeyField != null && !primaryKeyField.isEmpty()) {
|
|
|
|
|
+ targetData.forEach(row -> generateSnowflakeId(row, primaryKeyField));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 批量插入到目标表
|
|
|
|
|
+ int batchSize = 1000;
|
|
|
|
|
+ int totalMigrated = 0;
|
|
|
|
|
+
|
|
|
|
|
+ for (int i = 0; i < targetData.size(); i += batchSize) {
|
|
|
|
|
+ int end = Math.min(i + batchSize, targetData.size());
|
|
|
|
|
+ List<Map<String, Object>> batch = targetData.subList(i, end);
|
|
|
|
|
+
|
|
|
|
|
+ // 使用JdbcTemplate执行批量插入
|
|
|
|
|
+ if (!batch.isEmpty()) {
|
|
|
|
|
+ // 获取列名
|
|
|
|
|
+ Set<String> columns = batch.get(0).keySet();
|
|
|
|
|
+ String columnNames = String.join(", ", columns);
|
|
|
|
|
+
|
|
|
|
|
+ // 构建占位符
|
|
|
|
|
+ String placeholders = String.join(", ", Collections.nCopies(columns.size(), "?"));
|
|
|
|
|
+
|
|
|
|
|
+ // 构建SQL语句
|
|
|
|
|
+ String sql = "INSERT INTO " + targetTableName + " (" + columnNames + ") VALUES (" + placeholders + ")";
|
|
|
|
|
+
|
|
|
|
|
+ // 准备参数
|
|
|
|
|
+ List<Object[]> args = batch.stream()
|
|
|
|
|
+ .map(row -> columns.stream()
|
|
|
|
|
+ .map(col -> row.get(col))
|
|
|
|
|
+ .toArray())
|
|
|
|
|
+ .collect(Collectors.toList());
|
|
|
|
|
+
|
|
|
|
|
+ // 执行批量插入
|
|
|
|
|
+ int[] results = jdbcTemplate.batchUpdate(sql, args);
|
|
|
|
|
+ totalMigrated += Arrays.stream(results).sum();
|
|
|
|
|
+ log.info("已迁移 {}/{} 条记录", totalMigrated, targetData.size());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ log.info("数据迁移完成: 共迁移 {} 条记录", totalMigrated);
|
|
|
|
|
+ return totalMigrated;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional(rollbackFor = Exception.class)
|
|
|
|
|
+ public int updateTargetData(String targetTableName, Map<String, Object> dataMap, String whereCondition) {
|
|
|
|
|
+ log.info("更新目标表数据: 表={}, 数据={}, 条件={}", targetTableName, dataMap, whereCondition);
|
|
|
|
|
+
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(targetTableName)) {
|
|
|
|
|
+ throw new RuntimeException("目标表 " + targetTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return migrateDataMapper.updateTargetTable(targetTableName, dataMap, whereCondition);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional(rollbackFor = Exception.class)
|
|
|
|
|
+ public int deleteTargetData(String targetTableName, String whereCondition) {
|
|
|
|
|
+ log.info("删除目标表数据: 表={}, 条件={}", targetTableName, whereCondition);
|
|
|
|
|
+
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(targetTableName)) {
|
|
|
|
|
+ throw new RuntimeException("目标表 " + targetTableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return migrateDataMapper.deleteFromTargetTable(targetTableName, whereCondition);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ public List<Map<String, Object>> getTableStructure(String tableName) {
|
|
|
|
|
+ log.info("获取表结构: 表={}", tableName);
|
|
|
|
|
+
|
|
|
|
|
+ if (!migrateDataMapper.checkTableExists(tableName)) {
|
|
|
|
|
+ throw new RuntimeException("表 " + tableName + " 不存在");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return migrateDataMapper.getTableStructure(tableName);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ public boolean checkTableExists(String tableName) {
|
|
|
|
|
+ return migrateDataMapper.checkTableExists(tableName);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|