利用MyBatis ResultHandler与游标驱动,构建高效大数据流式处理管道
·
1. 为什么需要流式处理大数据?
想象一下你要从水库里取水浇灌农田。如果一次性把水库抽干,不仅需要巨大的储水容器,还可能压垮运输管道。更合理的做法是用水管持续引流,按需取用——这就是流式处理的核心逻辑。
在数据处理领域,传统的一次性加载方式(如List<User> users = userMapper.selectAll())在处理百万级数据时,相当于把整个水库的水瞬间灌进内存。我曾在生产环境遇到过因此导致的内存溢出(OOM),JVM堆内存直接飙到8GB上限,服务崩溃的惨痛教训。
流式处理的优势在于:
- 内存可控:每次只处理单条或小批量数据
- 实时响应:数据边读取边处理,无需等待全量加载
- 资源节约:减少数据库连接占用时间
2. MyBatis流式处理双剑客:ResultHandler与游标
2.1 ResultHandler的工作原理
ResultHandler是MyBatis提供的回调接口,其核心方法是:
void handleResult(ResultContext<?> context);
当数据库返回结果集时,MyBatis会逐行调用这个方法。我常用它来做这些事:
- 实时写入文件(如CSV导出)
- 分批插入到其他数据库
- 内存计算(如实时统计)
这里有个实战技巧:通过ResultContext对象可以获取当前处理状态:
// 获取当前行数据
User user = (User)context.getResultObject();
// 判断是否停止处理
if(shouldStop()) context.stop();
2.2 数据库游标的关键配置
单纯使用ResultHandler还不够,必须配合正确的游标类型。重点配置这两个参数:
@Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 1000)
void streamData(ResultHandler<User> handler);
FORWARD_ONLY:声明只进游标,禁止随机访问(节省内存)fetchSize:控制每次从数据库拉取的行数(实测1000-5000是较优区间)
踩坑提醒:MySQL驱动默认会缓存所有结果集到内存,必须同时在JDBC连接字符串添加useCursorFetch=true参数才能真正启用流式:
jdbc:mysql://host/db?useCursorFetch=true
3. 构建完整流式处理管道
3.1 基础实现方案
完整代码示例:
// 1. 自定义处理器
public class UserExportHandler implements ResultHandler<User> {
private final CSVWriter writer;
@Override
public void handleResult(ResultContext<? extends User> ctx) {
User user = ctx.getResultObject();
writer.writeRow(userToArray(user));
if(ctx.getResultCount() % 1000 == 0) {
writer.flush(); // 每1000行刷盘
}
}
}
// 2. Mapper接口配置
@Options(resultSetType = FORWARD_ONLY, fetchSize = 1000)
void streamUsers(@Param("minId") long minId, ResultHandler<User> handler);
// 3. 调用示例
try(CSVWriter writer = new CSVWriter(new FileWriter("output.csv"))) {
UserExportHandler handler = new UserExportHandler(writer);
userMapper.streamUsers(1000L, handler);
}
3.2 高并发优化方案
当需要并行处理时,推荐采用生产者-消费者模式:
// 创建有界队列防止内存膨胀
BlockingQueue<User> queue = new ArrayBlockingQueue<>(1000);
// ResultHandler作为生产者
public class QueueProducer implements ResultHandler<User> {
@Override
public void handleResult(ResultContext<? extends User> ctx) {
queue.put(ctx.getResultObject()); // 阻塞式插入
}
}
// 启动消费者线程池
ExecutorService pool = Executors.newFixedThreadPool(4);
for(int i=0; i<4; i++) {
pool.submit(() -> {
while(true) {
User user = queue.take();
processUser(user);
}
});
}
性能数据:在16核服务器上测试处理1000万条数据:
- 单线程耗时:182秒
- 4消费者线程:63秒
- 内存占用始终稳定在500MB以内
4. 生产环境实战经验
4.1 事务管理要点
流式处理中常见的事务误区:
// 错误示例:整个流式过程在事务中
@Transactional
public void exportData() {
mapper.streamData(handler); // 可能导致连接长时间占用
}
正确做法应该是:
- 分批次开启事务(每1000条提交一次)
- 使用
TransactionTemplate编程式事务 - 对于只读操作,设置
@Transactional(readOnly=true)
4.2 资源释放技巧
必须确保流关闭后释放资源。推荐try-with-resources模式:
try(SqlSession session = sqlSessionFactory.openSession()) {
UserMapper mapper = session.getMapper(UserMapper.class);
mapper.streamData(new MyHandler());
} // 自动关闭session和结果集
4.3 异常处理方案
流式处理中异常需要特殊处理:
public class SafeResultHandler implements ResultHandler<User> {
@Override
public void handleResult(ResultContext<? extends User> ctx) {
try {
businessLogic(ctx.getResultObject());
} catch (Exception e) {
ctx.stop(); // 终止流
throw new RuntimeException("处理失败", e);
}
}
}
5. 性能对比与调优
5.1 不同方案内存占用对比
| 处理方式 | 100万条数据内存峰值 | 耗时 |
|---|---|---|
| 传统List查询 | 1.2GB | 28s |
| 基础流式处理 | 50MB | 32s |
| 并行流式处理 | 80MB | 11s |
5.2 关键参数调优
-
fetchSize:
- MySQL建议1000-5000
- Oracle建议100-200
- PostgreSQL建议500-2000
-
JVM参数:
-XX:+UseG1GC -Xmx2g -Xms2g使用G1垃圾回收器能更好处理大量短期对象
-
批处理大小: 写入数据库时,每500-1000条提交一次事务效率最高
6. 典型应用场景
6.1 数据迁移工具
最近用这套方案实现了跨库数据迁移:
public class DataMigrator implements ResultHandler<SourceData> {
private final TargetMapper targetMapper;
private int batchCount = 0;
@Override
public void handleResult(ResultContext<? extends SourceData> ctx) {
targetMapper.insert(convert(ctx.getResultObject()));
if(++batchCount % 1000 == 0) {
sqlSession.commit(); // 分批提交
}
}
}
6.2 实时报表生成
金融行业日终报表处理示例:
public class ReportGenerator {
public void generateDailyReport(LocalDate date) {
ReportHandler handler = new ReportHandler();
// 流式处理交易数据
transactionMapper.streamByDate(date, handler);
// 实时生成统计结果
handler.saveReport();
}
}
6.3 大数据预处理
配合Spark等大数据框架时,可以这样对接:
JavaRDD<User> rdd = sparkContext.emptyRDD();
ResultHandler handler = context -> {
rdd.union(sparkContext.parallelize(
Arrays.asList(context.getResultObject())));
};
userMapper.streamAll(handler);
更多推荐
所有评论(0)