基于Spring Boot+Spring Batch的多表动态数据库迁移实现问询
动态数据库信息下的Spring Batch跨库迁移方案
一、抛弃静态Bean定义,运行时动态构建核心组件
Spring Batch的Job、Step、ItemReader/ItemWriter无需提前用@Bean声明,而是在API请求触发时,根据传入的数据库参数实时构建:
1. 动态生成数据源
基于请求体中的数据库信息,用DriverManagerDataSource快速创建源/目标数据源:
private DataSource buildDataSource(String url, String username, String password, String driverClass) { DriverManagerDataSource dataSource = new DriverManagerDataSource(); dataSource.setUrl(url); dataSource.setUsername(username); dataSource.setPassword(password); dataSource.setDriverClassName(driverClass); return dataSource; }
2. 适配未知表结构的Reader/Writer
针对源库多Schema、未知表结构的场景,用JDBC元数据+通用映射类处理:
// 动态构建读取器,适配任意表结构 private ItemReader<Map<String, Object>> buildDynamicReader(DataSource sourceDs, String schema, String table) { JdbcCursorItemReader<Map<String, Object>> reader = new JdbcCursorItemReader<>(); reader.setDataSource(sourceDs); // 自动查询表所有字段 reader.setSql(String.format("SELECT * FROM %s.%s", schema, table)); reader.setRowMapper(new ColumnMapRowMapper()); // 将结果映射为Map,适配未知字段 return reader; } // 动态构建写入器,通用插入逻辑 private ItemWriter<Map<String, Object>> buildDynamicWriter(DataSource targetDs, String targetTable) throws SQLException { JdbcBatchItemWriter<Map<String, Object>> writer = new JdbcBatchItemWriter<>(); writer.setDataSource(targetDs); writer.setItemSqlParameterSourceProvider(new MapSqlParameterSourceProvider()); // 读取目标表元数据,动态生成INSERT语句 DatabaseMetaData metaData = targetDs.getConnection().getMetaData(); ResultSet columns = metaData.getColumns(null, null, targetTable, "%"); StringBuilder columnsSb = new StringBuilder(); StringBuilder placeholdersSb = new StringBuilder(); while (columns.next()) { columnsSb.append(columns.getString("COLUMN_NAME")).append(","); placeholdersSb.append(":").append(columns.getString("COLUMN_NAME")).append(","); } // 移除末尾逗号 String columnStr = columnsSb.substring(0, columnsSb.length()-1); String placeholderStr = placeholdersSb.substring(0, placeholdersSb.length()-1); writer.setSql(String.format("INSERT INTO %s (%s) VALUES (%s)", targetTable, columnStr, placeholderStr)); return writer; }
3. API接口中动态构建并启动Job
在接口里接收迁移请求,遍历源库表,逐个构建Job并执行:
@RestController public class MigrationController { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final JobLauncher jobLauncher; public MigrationController(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, JobLauncher jobLauncher) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.jobLauncher = jobLauncher; } @PostMapping("/start-migration") public String startMigration(@RequestBody MigrationReq req) throws Exception { // 构建源、目标数据源 DataSource sourceDs = buildDataSource(req.getSourceUrl(), req.getSourceUsername(), req.getSourcePassword(), req.getSourceDriver()); DataSource targetDs = buildDataSource(req.getTargetUrl(), req.getTargetUsername(), req.getTargetPassword(), req.getTargetDriver()); // 获取源库指定Schema下的所有表 DatabaseMetaData sourceMeta = sourceDs.getConnection().getMetaData(); ResultSet tables = sourceMeta.getTables(null, req.getSourceSchema(), "%", new String[]{"TABLE"}); while (tables.next()) { String tableName = tables.getString("TABLE_NAME"); // 构建Reader、Writer ItemReader<Map<String, Object>> reader = buildDynamicReader(sourceDs, req.getSourceSchema(), tableName); ItemWriter<Map<String, Object>> writer = buildDynamicWriter(targetDs, tableName); // 构建Step Step migrateStep = stepBuilderFactory.get("step-" + tableName + "-" + System.currentTimeMillis()) .<Map<String, Object>, Map<String, Object>>chunk(1000) .reader(reader) .writer(writer) .transactionManager(new DataSourceTransactionManager(targetDs)) .build(); // 构建Job Job migrateJob = jobBuilderFactory.get("job-" + tableName + "-" + System.currentTimeMillis()) .start(migrateStep) .build(); // 启动Job JobParameters params = new JobParametersBuilder() .addString("table", tableName) .addLong("timestamp", System.currentTimeMillis()) .toJobParameters(); jobLauncher.run(migrateJob, params); } return "迁移任务已启动"; } // 内部类:接收请求参数 static class MigrationReq { private String sourceUrl; private String sourceUsername; private String sourcePassword; private String sourceDriver; private String sourceSchema; private String targetUrl; private String targetUsername; private String targetPassword; private String targetDriver; // getter/setter省略 } }
二、处理多Schema与未知表结构的核心技巧
- 通过JDBC元数据获取库表信息:用
DatabaseMetaData可以遍历指定Schema下的所有表,以及表的字段结构,完全不需要提前知晓表定义:
// 获取指定Schema下的所有表 ResultSet tables = metaData.getTables(null, schemaName, "%", new String[]{"TABLE"}); // 获取指定表的所有字段 ResultSet columns = metaData.getColumns(null, schemaName, tableName, "%");
- 用Map作为通用数据载体:
ColumnMapRowMapper将查询结果自动转为Map<String, Object>,适配任意表结构的读取;写入时通过MapSqlParameterSourceProvider自动映射参数,配合动态生成的INSERT语句完成通用写入。
三、关键注意事项
- 保证Job/Step名称唯一:每次动态创建时加入表名+时间戳,避免Spring Batch因重复注册Bean报错。
- 事务隔离:针对目标数据源单独配置事务管理器,避免跨数据源事务问题。
- 大表优化:如果表数据量巨大,改用
JdbcPagingItemReader实现分页读取,避免内存溢出。 - 异常监控:给Job添加
JobExecutionListener,记录迁移失败的表名和原因,方便后续重试。
内容的提问来源于stack exchange,提问作者XPulse
相关产品推荐
相关产品推荐

