如何在JDBC Batch Item Writer中实现动态SQL完成MongoDB多集合同步到SQL Server
问题根源与修复方案
现有代码的核心错误
SqlString类的schema、table、document三个属性没有初始化赋值,Bean初始化阶段调用getObject()会直接抛出空指针异常- 强转
document.keySet()为CopyOnWriteArraySet类型不匹配,org.bson.Document的键集合是普通的Set实现,强转会抛出类转换异常 - 生成的INSERT语句占位符格式错误,JDBC命名参数需要加冒号前缀,你当前生成的
VALUES (a,b,c)会被识别为字段引用,直接执行会报错 - 依赖单个Document生成SQL逻辑不可靠:同一个Mongo集合的不同文档可能存在字段差异,用单个文档的键生成SQL会导致其他缺失对应字段的文档写入失败
- 你在Spring Bean初始化阶段生成SQL,此时还没有读取到Mongo的任何数据,根本拿不到文档字段列表
修复方案
1. 先预取目标集合的全量字段列表
在任务启动前,扫描目标Mongo集合的样本数据,合并所有出现过的字段,生成完整的字段列表,避免字段遗漏。
2. 使用StepScope动态生成SQL
使用Spring Batch的@StepScope注解,在Step启动阶段根据当前要处理的集合动态生成SQL,适配100+集合的批量处理需求。
3. 适配参数缺失场景
在参数提供者中补全缺失的字段,统一赋值为null,保证和SQL占位符数量匹配。
修正后代码
Step配置
// 用Job参数传递当前要处理的集合名、SQL Server schema、表名 @Bean @StepScope public MongoItemReader<Document> reader(@Value("#{jobParameters['mongoCollection']}") String collection) { MongoItemReader<Document> reader = new MongoItemReader<>(); reader.setTemplate(mongoTemplate); reader.setCollection(collection); reader.setTargetType(Document.class); reader.setQuery("{}"); reader.setSort(Collections.singletonMap("_id", Sort.Direction.ASC)); reader.setPageSize(1000); return reader; } @Bean @StepScope public JdbcBatchItemWriter<Document> writer( @Value("#{jobParameters['sqlSchema']}") String schema, @Value("#{jobParameters['sqlTable']}") String table, @Value("#{jobParameters['mongoCollection']}") String collection ) throws Exception { // 预取当前集合的全量字段 Set<String> allFields = getAllFieldsFromMongoCollection(collection); // 生成动态SQL String insertSql = generateInsertSql(schema, table, allFields); JdbcBatchItemWriter<Document> writer = new JdbcBatchItemWriter<>(); writer.setDataSource(dataSource); writer.setAssertUpdates(true); writer.setItemSqlParameterSourceProvider(new TableSourceProvider(allFields)); writer.setSql(insertSql); // 适配SQL Server的命名参数处理 writer.setJdbcTemplate(new NamedParameterJdbcTemplate(dataSource)); return writer; }
字段拉取与SQL生成工具方法
// 预拉取Mongo集合的所有字段,拉取前1000条样本合并字段,可根据实际情况调整样本量 private Set<String> getAllFieldsFromMongoCollection(String collectionName) { Set<String> fields = new HashSet<>(); mongoTemplate.executeQuery(new Query().limit(1000), collectionName, document -> { fields.addAll(document.keySet()); }); // 可选:把_id字段转成id字符串,根据业务需求调整 if(fields.contains("_id") && !fields.contains("id")) { fields.add("id"); } return fields; } // 生成符合SQL Server规范的INSERT语句 private String generateInsertSql(String schema, String table, Set<String> fields) { StringBuilder sql = new StringBuilder("INSERT INTO "); sql.append(schema).append(".").append(table).append(" "); // 字段部分加方括号避免关键字冲突 String columns = fields.stream().map(f -> "[" + f + "]").collect(Collectors.joining(",", "(", ")")); sql.append(columns); // 占位符加冒号前缀 String placeholders = fields.stream().map(f -> ":" + f).collect(Collectors.joining(",", " VALUES (", ")")); sql.append(placeholders); return sql.toString(); }
修正后的ItemSqlParameterSourceProvider
public class TableSourceProvider implements ItemSqlParameterSourceProvider<Document> { private final Set<String> allFields; public TableSourceProvider(Set<String> allFields) { this.allFields = allFields; } @Override public SqlParameterSource createSqlParameterSource(Document document) { Map<String, Object> paramMap = new HashMap<>(); // 先填充所有字段,默认值为null allFields.forEach(field -> paramMap.put(field, null)); // 用文档实际值覆盖 document.forEach((k, v) -> { if("_id".equals(k)) { // 把ObjectId转成字符串存储,可根据业务需求调整 paramMap.put(k, v.toString()); paramMap.put("id", v.toString()); } else if (v instanceof List || v instanceof org.bson.Document) { // 复杂类型转JSON字符串存储,可根据业务需求调整为拆分关联表 paramMap.put(k, com.mongodb.util.JSON.serialize(v)); } else { paramMap.put(k, v); } }); return new MapSqlParameterSource(paramMap); } }
注意事项
- 如果Mongo集合的文档字段差异极大,建议加大预拉取的样本量,避免遗漏字段
- 数值类型、日期类型要根据SQL Server的表结构做适配,避免类型转换错误
- 批量处理100+集合时,可以用循环生成Job实例,每个集合对应一个Job,传入不同的集合名、schema、表名参数即可,不需要重复写配置代码
内容的提问来源于stack exchange,提问作者Rio
相关产品推荐
相关产品推荐

