wangzhijun пре 3 месеци
родитељ
комит
5d95992857

+ 185 - 0
platform/src/main/java/com/jzg/controller/MigrateDateController.java

@@ -0,0 +1,185 @@
+package com.jzg.controller;
+
+import com.jzg.commons.core.base.BaseController;
+import com.jzg.commons.core.page.HttpResult;
+import com.jzg.entity.dto.MigrateDataParam;
+import com.jzg.service.MigrateDataService;
+import io.swagger.v3.oas.annotations.Operation;
+import io.swagger.v3.oas.annotations.tags.Tag;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.validation.annotation.Validated;
+import org.springframework.web.bind.annotation.*;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 数据迁移控制器
+ */
+@RestController
+@RequestMapping("/migrate/data")
+@Tag(name = "数据迁移管理")
+@Slf4j
+public class MigrateDateController extends BaseController {
+
+    @Autowired
+    private MigrateDataService migrateDataService;
+
+    /**
+     * 从源表迁移数据到目标表
+     * @param param 迁移参数
+     * @return 迁移结果
+     */
+    @Operation(summary = "迁移数据")
+    @PostMapping("/migrate")
+    public HttpResult migrateData(@Validated @RequestBody MigrateDataParam param) {
+        try {
+            // 如果需要先清空目标表
+            if (Boolean.TRUE.equals(param.getClearTargetBeforeMigrate())) {
+                String clearCondition = param.getClearTargetCondition();
+                int deleteCount = migrateDataService.deleteTargetData(
+                    param.getTargetTableName(),
+                    clearCondition != null ? clearCondition : "1=1"
+                );
+                log.info("已清空目标表 {} 中 {} 条数据", param.getTargetTableName(), deleteCount);
+            }
+
+            // 执行数据迁移
+            int migratedCount;
+            String primaryKeyField = param.getPrimaryKeyField();
+            boolean generateSnowflakeId = Boolean.TRUE.equals(param.getGenerateSnowflakeId());
+            String sourcePrimaryKeyField = param.getSourcePrimaryKeyField();
+            boolean saveOriginalId = Boolean.TRUE.equals(param.getSaveOriginalId());
+
+            if (param.getFieldMapping() != null && !param.getFieldMapping().isEmpty()) {
+                // 使用字段映射迁移
+                migratedCount = migrateDataService.migrateDataWithMapping(
+                    param.getSourceTableName(),
+                    param.getTargetTableName(),
+                    param.getFieldMapping(),
+                    param.getWhereCondition(),
+                    primaryKeyField,
+                    generateSnowflakeId,
+                    sourcePrimaryKeyField,
+                    saveOriginalId
+                );
+            } else if (param.getFields() != null && !param.getFields().isEmpty()) {
+                // 迁移指定字段
+                migratedCount = migrateDataService.migrateDataWithFields(
+                    param.getSourceTableName(),
+                    param.getTargetTableName(),
+                    param.getFields(),
+                    param.getWhereCondition(),
+                    primaryKeyField,
+                    generateSnowflakeId,
+                    sourcePrimaryKeyField,
+                    saveOriginalId
+                );
+            } else {
+                // 迁移所有字段
+                migratedCount = migrateDataService.migrateData(
+                    param.getSourceTableName(),
+                    param.getTargetTableName(),
+                    param.getWhereCondition(),
+                    primaryKeyField,
+                    generateSnowflakeId,
+                    sourcePrimaryKeyField,
+                    saveOriginalId
+                );
+            }
+
+            return HttpResult.ok("数据迁移成功,共迁移 " + migratedCount + " 条记录");
+        } catch (Exception e) {
+            log.error("数据迁移失败", e);
+            return HttpResult.error("数据迁移失败: " + e.getMessage());
+        }
+    }
+
+    /**
+     * 更新目标表数据
+     * @param targetTableName 目标表名
+     * @param dataMap 需要更新的数据
+     * @param whereCondition 更新条件
+     * @return 更新结果
+     */
+    @Operation(summary = "更新目标表数据")
+    @PutMapping("/update/{targetTableName}")
+    public HttpResult updateTargetData(
+            @PathVariable String targetTableName,
+            @RequestBody Map<String, Object> dataMap,
+            @RequestParam(required = false) String whereCondition) {
+        try {
+            if (whereCondition == null || whereCondition.trim().isEmpty()) {
+                return HttpResult.error("更新条件不能为空");
+            }
+
+            int count = migrateDataService.updateTargetData(targetTableName, dataMap, whereCondition);
+            return HttpResult.ok("更新成功,影响 " + count + " 条记录");
+        } catch (Exception e) {
+            log.error("更新数据失败", e);
+            return HttpResult.error("更新失败: " + e.getMessage());
+        }
+    }
+
+    /**
+     * 删除目标表数据
+     * @param targetTableName 目标表名
+     * @param whereCondition 删除条件
+     * @return 删除结果
+     */
+    @Operation(summary = "删除目标表数据")
+    @DeleteMapping("/delete/{targetTableName}")
+    public HttpResult deleteTargetData(
+            @PathVariable String targetTableName,
+            @RequestParam String whereCondition) {
+        try {
+            if (whereCondition == null || whereCondition.trim().isEmpty()) {
+                return HttpResult.error("删除条件不能为空");
+            }
+
+            int count = migrateDataService.deleteTargetData(targetTableName, whereCondition);
+            return HttpResult.ok("删除成功,影响 " + count + " 条记录");
+        } catch (Exception e) {
+            log.error("删除数据失败", e);
+            return HttpResult.error("删除失败: " + e.getMessage());
+        }
+    }
+
+    /**
+     * 获取表结构信息
+     * @param tableName 表名
+     * @return 表结构信息
+     */
+    @Operation(summary = "获取表结构信息")
+    @GetMapping("/structure/{tableName}")
+    public HttpResult getTableStructure(@PathVariable String tableName) {
+        try {
+            List<Map<String, Object>> structure = migrateDataService.getTableStructure(tableName);
+            return HttpResult.ok(structure);
+        } catch (Exception e) {
+            log.error("获取表结构失败", e);
+            return HttpResult.error("获取表结构失败: " + e.getMessage());
+        }
+    }
+
+    /**
+     * 检查表是否存在
+     * @param tableName 表名
+     * @return 检查结果
+     */
+    @Operation(summary = "检查表是否存在")
+    @GetMapping("/exists/{tableName}")
+    public HttpResult checkTableExists(@PathVariable String tableName) {
+        try {
+            boolean exists = migrateDataService.checkTableExists(tableName);
+            return HttpResult.ok(Map.of(
+                "tableName", tableName,
+                "exists", exists
+            ));
+        } catch (Exception e) {
+            log.error("检查表是否存在失败", e);
+            return HttpResult.error("检查失败: " + e.getMessage());
+        }
+    }
+}

+ 49 - 0
platform/src/main/java/com/jzg/entity/dto/MigrateDataParam.java

@@ -0,0 +1,49 @@
+
+package com.jzg.entity.dto;
+
+import io.swagger.v3.oas.annotations.media.Schema;
+import lombok.Data;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 数据迁移参数
+ */
+@Data
+@Schema(description = "数据迁移参数")
+public class MigrateDataParam {
+
+    @Schema(description = "源表名", required = true)
+    private String sourceTableName;
+
+    @Schema(description = "目标表名", required = true)
+    private String targetTableName;
+
+    @Schema(description = "查询条件(SQL WHERE子句,不包含WHERE关键字)")
+    private String whereCondition;
+
+    @Schema(description = "需要迁移的字段列表(为空则迁移所有字段)")
+    private List<String> fields;
+
+    @Schema(description = "字段映射关系(源字段名 -> 目标字段名)")
+    private Map<String, String> fieldMapping;
+
+    @Schema(description = "是否先清空目标表数据")
+    private Boolean clearTargetBeforeMigrate = false;
+
+    @Schema(description = "清空目标表的条件(SQL WHERE子句,不包含WHERE关键字)")
+    private String clearTargetCondition;
+
+    @Schema(description = "目标表的主键字段名")
+    private String primaryKeyField;
+
+    @Schema(description = "是否为主键生成雪花ID(默认false)")
+    private Boolean generateSnowflakeId = false;
+
+    @Schema(description = "是否保存原始ID到original_id字段(默认false)")
+    private Boolean saveOriginalId = false;
+
+    @Schema(description = "源表的主键字段名(用于保存到original_id)")
+    private String sourcePrimaryKeyField;
+}

+ 101 - 0
platform/src/main/java/com/jzg/mapper/MigrateDataMapper.java

@@ -0,0 +1,101 @@
+
+package com.jzg.mapper;
+
+import com.jzg.commons.annotation.IgnoreTenant;
+import org.apache.ibatis.annotations.*;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 数据迁移Mapper
+ */
+@Mapper
+public interface MigrateDataMapper {
+
+    /**
+     * 从源表查询数据(使用JDBC直接查询,避免MyBatis Plus类型处理器干扰)
+     * @param sourceTableName 源表名
+     * @param whereCondition 查询条件(可为空)
+     * @return 查询结果
+     */
+    @IgnoreTenant
+    List<Map<String, Object>> selectFromSourceTable(@Param("sourceTableName") String sourceTableName, 
+                                                     @Param("whereCondition") String whereCondition);
+
+    /**
+     * 插入数据到目标表
+     * @param targetTableName 目标表名
+     * @param dataMap 数据Map
+     * @return 影响行数
+     */
+    @IgnoreTenant
+    int insertToTargetTable(@Param("targetTableName") String targetTableName, 
+                           @Param("dataMap") Map<String, Object> dataMap);
+
+    /**
+     * 批量插入数据到目标表
+     * @param targetTableName 目标表名
+     * @param dataList 数据列表
+     * @return 影响行数
+     */
+    @IgnoreTenant
+    int batchInsertToTargetTable(@Param("targetTableName") String targetTableName, 
+                                @Param("dataList") List<Map<String, Object>> dataList);
+
+    /**
+     * 更新目标表数据
+     * @param targetTableName 目标表名
+     * @param dataMap 数据Map
+     * @param whereCondition 更新条件
+     * @return 影响行数
+     */
+    @IgnoreTenant
+    int updateTargetTable(@Param("targetTableName") String targetTableName,
+                         @Param("dataMap") Map<String, Object> dataMap,
+                         @Param("whereCondition") String whereCondition);
+
+    /**
+     * 删除目标表数据
+     * @param targetTableName 目标表名
+     * @param whereCondition 删除条件
+     * @return 影响行数
+     */
+    @IgnoreTenant
+    int deleteFromTargetTable(@Param("targetTableName") String targetTableName,
+                             @Param("whereCondition") String whereCondition);
+
+    /**
+     * 检查目标表是否存在
+     * @param tableName 表名
+     * @return 存在返回true
+     */
+    @IgnoreTenant
+    boolean checkTableExists(@Param("tableName") String tableName);
+
+    /**
+     * 检查表中是否存在指定列
+     * @param tableName 表名
+     * @param columnName 列名
+     * @return 存在返回true
+     */
+    @IgnoreTenant
+    boolean checkColumnExists(@Param("tableName") String tableName, @Param("columnName") String columnName);
+
+    /**
+     * 为表添加original_id列
+     * @param tableName 表名
+     * @return 影响行数
+     */
+    @IgnoreTenant
+    int addOriginalIdColumn(@Param("tableName") String tableName);
+
+    /**
+     * 获取表结构信息
+     * @param tableName 表名
+     * @return 表字段信息
+     */
+    @IgnoreTenant
+    List<Map<String, Object>> getTableStructure(@Param("tableName") String tableName);
+}

+ 121 - 0
platform/src/main/java/com/jzg/service/MigrateDataService.java

@@ -0,0 +1,121 @@
+package com.jzg.service;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 数据迁移服务接口
+ */
+public interface MigrateDataService {
+
+    /**
+     * 从源表迁移数据到目标表
+     * @param sourceTableName 源表名
+     * @param targetTableName 目标表名
+     * @param whereCondition 查询条件(可为空)
+     * @return 迁移的记录数
+     */
+    int migrateData(String sourceTableName, String targetTableName, String whereCondition);
+
+    /**
+     * 从源表迁移数据到目标表(支持雪花ID生成)
+     * @param sourceTableName 源表名
+     * @param targetTableName 目标表名
+     * @param whereCondition 查询条件(可为空)
+     * @param primaryKeyField 主键字段名
+     * @param generateSnowflakeId 是否生成雪花ID
+     * @param sourcePrimaryKeyField 源表主键字段名
+     * @param saveOriginalId 是否保存原始ID
+     * @return 迁移的记录数
+     */
+    int migrateData(String sourceTableName, String targetTableName, String whereCondition,
+                    String primaryKeyField, boolean generateSnowflakeId,
+                    String sourcePrimaryKeyField, boolean saveOriginalId);
+
+    /**
+     * 从源表迁移数据到目标表(支持字段映射)
+     * @param sourceTableName 源表名
+     * @param targetTableName 目标表名
+     * @param fieldMapping 字段映射关系(源字段->目标字段)
+     * @param whereCondition 查询条件(可为空)
+     * @return 迁移的记录数
+     */
+    int migrateDataWithMapping(String sourceTableName, String targetTableName,
+                              Map<String, String> fieldMapping, String whereCondition);
+
+    /**
+     * 从源表迁移数据到目标表(支持字段映射和雪花ID生成)
+     * @param sourceTableName 源表名
+     * @param targetTableName 目标表名
+     * @param fieldMapping 字段映射关系(源字段->目标字段)
+     * @param whereCondition 查询条件(可为空)
+     * @param primaryKeyField 主键字段名
+     * @param generateSnowflakeId 是否生成雪花ID
+     * @param sourcePrimaryKeyField 源表主键字段名
+     * @param saveOriginalId 是否保存原始ID
+     * @return 迁移的记录数
+     */
+    int migrateDataWithMapping(String sourceTableName, String targetTableName,
+                              Map<String, String> fieldMapping, String whereCondition,
+                              String primaryKeyField, boolean generateSnowflakeId,
+                              String sourcePrimaryKeyField, boolean saveOriginalId);
+
+    /**
+     * 从源表迁移指定字段的数据到目标表
+     * @param sourceTableName 源表名
+     * @param targetTableName 目标表名
+     * @param fields 需要迁移的字段列表
+     * @param whereCondition 查询条件(可为空)
+     * @return 迁移的记录数
+     */
+    int migrateDataWithFields(String sourceTableName, String targetTableName,
+                             List<String> fields, String whereCondition);
+
+    /**
+     * 从源表迁移指定字段的数据到目标表(支持雪花ID生成)
+     * @param sourceTableName 源表名
+     * @param targetTableName 目标表名
+     * @param fields 需要迁移的字段列表
+     * @param whereCondition 查询条件(可为空)
+     * @param primaryKeyField 主键字段名
+     * @param generateSnowflakeId 是否生成雪花ID
+     * @param sourcePrimaryKeyField 源表主键字段名
+     * @param saveOriginalId 是否保存原始ID
+     * @return 迁移的记录数
+     */
+    int migrateDataWithFields(String sourceTableName, String targetTableName,
+                             List<String> fields, String whereCondition,
+                             String primaryKeyField, boolean generateSnowflakeId,
+                             String sourcePrimaryKeyField, boolean saveOriginalId);
+
+    /**
+     * 更新目标表数据
+     * @param targetTableName 目标表名
+     * @param dataMap 需要更新的数据
+     * @param whereCondition 更新条件
+     * @return 影响的记录数
+     */
+    int updateTargetData(String targetTableName, Map<String, Object> dataMap, String whereCondition);
+
+    /**
+     * 删除目标表数据
+     * @param targetTableName 目标表名
+     * @param whereCondition 删除条件
+     * @return 影响的记录数
+     */
+    int deleteTargetData(String targetTableName, String whereCondition);
+
+    /**
+     * 获取表结构信息
+     * @param tableName 表名
+     * @return 表结构信息
+     */
+    List<Map<String, Object>> getTableStructure(String tableName);
+
+    /**
+     * 检查表是否存在
+     * @param tableName 表名
+     * @return 存在返回true
+     */
+    boolean checkTableExists(String tableName);
+}

+ 463 - 0
platform/src/main/java/com/jzg/service/impl/MigrateDataServiceImpl.java

@@ -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);
+    }
+}

+ 102 - 0
platform/src/main/resources/mapper/MigrateDataMapper.xml

@@ -0,0 +1,102 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
+<mapper namespace="com.jzg.mapper.MigrateDataMapper">
+
+    <!-- 从源表查询数据 -->
+    <select id="selectFromSourceTable" resultType="java.util.HashMap">
+        SELECT * FROM ${sourceTableName}
+        <if test="whereCondition != null and whereCondition != ''">
+            WHERE ${whereCondition}
+        </if>
+    </select>
+
+    <!-- 插入数据到目标表 -->
+    <insert id="insertToTargetTable">
+        INSERT INTO ${targetTableName}
+        <foreach collection="dataMap.keys" item="key" open="(" separator="," close=")">
+            ${key}
+        </foreach>
+        VALUES
+        <foreach collection="dataMap.keys" item="key" open="(" separator="," close=")">
+            #{dataMap[${key}]}
+        </foreach>
+    </insert>
+
+    <!-- 批量插入数据到目标表 -->
+    <insert id="batchInsertToTargetTable">
+        INSERT INTO ${targetTableName}
+        <trim prefix="(" suffix=")" suffixOverrides=",">
+            <foreach collection="dataList[0].keys" item="key">
+                ${key},
+            </foreach>
+        </trim>
+        VALUES
+        <foreach collection="dataList" item="item" separator=",">
+            <trim prefix="(" suffix=")" suffixOverrides=",">
+                <foreach collection="item.keys" item="key">
+                    #{item[${key}]},
+                </foreach>
+            </trim>
+        </foreach>
+    </insert>
+
+    <!-- 更新目标表数据 -->
+    <update id="updateTargetTable">
+        UPDATE ${targetTableName}
+        <set>
+            <foreach collection="dataMap.keys" item="key" separator=",">
+                ${key} = #{dataMap[${key}]}
+            </foreach>
+        </set>
+        <if test="whereCondition != null and whereCondition != ''">
+            WHERE ${whereCondition}
+        </if>
+    </update>
+
+    <!-- 删除目标表数据 -->
+    <delete id="deleteFromTargetTable">
+        DELETE FROM ${targetTableName}
+        <if test="whereCondition != null and whereCondition != ''">
+            WHERE ${whereCondition}
+        </if>
+    </delete>
+
+    <!-- 检查表是否存在 -->
+    <select id="checkTableExists" resultType="boolean">
+        SELECT COUNT(*) > 0
+        FROM information_schema.tables
+        WHERE table_schema = DATABASE()
+        AND table_name = #{tableName}
+    </select>
+
+    <!-- 获取表结构信息 -->
+    <select id="getTableStructure" resultType="java.util.HashMap">
+        SELECT 
+            COLUMN_NAME as columnName,
+            DATA_TYPE as dataType,
+            IS_NULLABLE as isNullable,
+            COLUMN_KEY as columnKey,
+            COLUMN_DEFAULT as columnDefault,
+            COLUMN_COMMENT as columnComment
+        FROM information_schema.columns
+        WHERE table_schema = DATABASE()
+        AND table_name = #{tableName}
+        ORDER BY ORDINAL_POSITION
+    </select>
+
+    <!-- 检查表中是否存在指定列 -->
+    <select id="checkColumnExists" resultType="boolean">
+        SELECT COUNT(*) > 0
+        FROM information_schema.columns
+        WHERE table_schema = DATABASE()
+        AND table_name = #{tableName}
+        AND column_name = #{columnName}
+    </select>
+
+    <!-- 为表添加original_id列 -->
+    <update id="addOriginalIdColumn">
+        ALTER TABLE ${tableName}
+        ADD COLUMN original_id VARCHAR(64) COMMENT '原始数据ID'
+    </update>
+
+</mapper>