
本文介绍如何通过流式处理(Stream API)和异步编程优化多数据源批量查询,替代传统 List 全量收集方式,显著降低内存占用与响应延迟。
本文介绍如何通过流式处理(stream api)和异步编程优化多数据源批量查询,替代传统 `list
在处理 30+ 数据源并发查询时,原始代码使用 forEach 同步遍历并累积 List<map object>></map>,极易引发 OutOfMemoryError(如截图所示),根本原因在于:
- 所有查询结果一次性加载至堆内存;
-
Map<string object></string>中嵌套的List(如queryForList返回的每行记录)进一步放大内存压力; - 日志框架(如 Logback)尝试序列化打印该大对象,加剧 GC 压力与线程阻塞。
✅ 推荐解决方案:流式 + 异步 + 懒加载
1. 使用 Stream 替代 forEach,实现惰性求值
避免预分配大集合,改用 Stream 管道处理每个数据源结果,并立即消费(如写入文件、发送消息或分页返回):
// 将数据源映射为 Stream,逐个处理,不缓存全部结果
multiDataSourceProperties.getDatasources().entrySet().stream()
.map(entry -> {
jdbcTemplate.setDataSource((DataSource) entry.getValue());
String key = (String) entry.getKey();
List<Map<String, Object>> rows = namedJdbcTemplate.queryForList(rsConfig.getQuery(), queryParams);
return Map.entry(key, rows); // 返回键值对,避免嵌套 Map
})
.forEach(resultEntry -> {
// ✅ 立即处理单个数据源结果:例如写入 CSV 文件
writeToCsv(resultEntry.getKey(), resultEntry.getValue());
// 或发送至消息队列:kafkaTemplate.send("query-results", resultEntry);
});2. 引入异步非阻塞执行(关键优化)
使用 CompletableFuture 并行查询,避免线程阻塞,同时控制并发度防止数据库过载:
ExecutorService executor = Executors.newFixedThreadPool(10); // 限制并发数
List<CompletableFuture<Void>> futures = multiDataSourceProperties.getDatasources().entrySet().stream()
.map(entry -> CompletableFuture.runAsync(() -> {
jdbcTemplate.setDataSource((DataSource) entry.getValue());
String key = (String) entry.getKey();
List<Map<String, Object>> rows = namedJdbcTemplate.queryForList(rsConfig.getQuery(), queryParams);
writeToCsv(key, rows); // 异步写入,不阻塞主线程
}, executor))
.collect(Collectors.toList());
// 等待全部完成(可选),或直接返回响应,后台继续处理
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();3. 进阶:采用 Reactive 编程(Spring WebFlux + R2DBC)
若系统支持响应式栈,可彻底消除阻塞 I/O:
- 使用
R2dbcEntityTemplate替代JdbcTemplate; - 数据源查询返回
Flux<list object>>></list>; - 通过
flatMap并行处理,配合writeWith直接流式响应客户端(如 SSE 或 Chunked Transfer)。
⚠️ 重要注意事项
-
禁止日志全量打印:禁用
log.debug("Result: {}", resultSetList)类语句,改用log.debug("Processed {} data sources", count); -
资源清理:确保
ExecutorService在应用关闭时shutdown(); -
下游限流:若写入文件或消息队列,需添加背压(Backpressure)机制(如
onBackpressureBuffer); - 监控告警:对 JVM 堆内存、GC 时间、线程池队列长度设置 Prometheus + Grafana 监控。
综上,问题本质不是“如何用 Stream 发送 List
立即学习“Java免费学习笔记(深入)”;


















