这次我们来看一个 Java 开发中非常实际的问题:如何用 MyBatis 的流式查询,解决一次性加载海量数据导致的内存溢出(OOM)。很多开发者都遇到过,一个看似简单的SELECT * FROM large_table,在数据量达到百万级时,如果直接返回List,JVM 堆内存瞬间就会被撑爆,服务直接崩溃。MyBatis 提供的流式查询(Streaming Query)就是为这种场景而生的,它允许你像处理流一样,逐条从数据库读取和处理数据,而不是一次性把所有数据都加载到内存里。
这篇文章的重点不是讲复杂的理论,而是直接告诉你:流式查询能不能用?怎么用?用了之后效果如何?我们会从核心概念、代码实现、性能对比到生产环境的最佳实践,一步步拆解。如果你正在处理报表导出、大数据分析、数据迁移等需要处理大量数据的任务,这篇文章可以直接收藏备用。
1. 核心能力速览
在深入细节之前,我们先快速了解 MyBatis 流式查询的核心特性和适用边界。
| 能力项 | 说明 |
|---|---|
| 核心功能 | 以“流”的方式逐条处理数据库查询结果,避免一次性加载全部数据到 JVM 内存。 |
| 解决痛点 | 防止因查询百万级以上数据导致的内存溢出(OOM)和长时间 Full GC。 |
| 技术实现 | 基于 JDBC 的ResultSet游标,通过fetchSize和resultSetType等参数控制。 |
| 适用场景 | 大数据量报表生成、数据导出、ETL 数据迁移、日志分析等需要遍历海量结果集的场景。 |
| 不适用场景 | 需要随机访问结果集、频繁进行聚合计算或需要将全部数据在内存中多次遍历的场景。 |
| 性能影响 | 网络 I/O 和数据库连接占用时间可能变长,但内存占用极低,适合内存敏感型应用。 |
| 启动/使用方式 | 通过 MyBatis 的Cursor<T>接口或自定义ResultHandler实现,无需额外服务部署。 |
简单来说,流式查询就是把“一口吃成胖子”的查询,变成了“细嚼慢咽”的处理过程。
2. 适用场景与使用边界
2.1 谁需要流式查询?
- 后端开发者:需要从数据库导出大量数据到 CSV/Excel 文件。
- 数据平台工程师:负责将数据从在线库迁移到离线分析库。
- 报表系统维护者:系统需要生成包含数十万甚至百万行数据的复杂报表。
- 任何面临
OutOfMemoryError: Java heap space的 Java 服务开发者。
2.2 它能解决什么问题?
核心是解决内存瓶颈。传统查询List<User> users = userMapper.selectAll();会在内存中构造一个包含所有User对象的列表。假设一条记录 1KB,100 万条就是 1GB,这很容易超过 JVM 堆内存限制,触发 OOM。流式查询则是在遍历Cursor时,MyBatis 和 JDBC 驱动配合,每次只从数据库网络连接中获取少量记录(例如几百条)到内存,处理完后再获取下一批,内存中始终只保持一个很小的数据窗口。
2.3 不适合什么场景?
- 需要全量数据在内存中运算:例如需要对整个结果集进行排序、分组、多次遍历。流式查询是单向遍历,无法回头。
- 事务时间非常长的场景:流式查询通常需要保持数据库连接和事务直到处理完毕,长时间占用连接可能影响连接池。
- 网络环境极差:因为需要多次网络往返获取数据,高延迟网络下总耗时可能比一次性获取更长。
- 对数据库压力敏感:流式查询可能会对数据库造成更长时间的读压力(持有游标),需要评估数据库端的承受能力。
2.4 合规与边界
- 数据库兼容性:并非所有数据库和 JDBC 驱动都完美支持流式查询,需测试目标数据库(如 MySQL, PostgreSQL)的对应驱动。
- 连接池配置:使用如 HikariCP、Druid 等连接池时,需注意连接归还机制,避免流未关闭导致连接泄漏。
- 资源及时释放:必须确保
Cursor或ResultSet被正确关闭,否则会导致数据库游标泄漏和内存泄漏。
3. 环境准备与前置条件
在开始编码前,请确保你的环境满足以下条件。这不是一个需要 GPU 或特定硬件的 AI 模型,但对 Java 环境和数据库有一定要求。
Java 开发环境
- JDK 8 或更高版本(推荐 JDK 11 或 17)。
- Maven 或 Gradle 构建工具。
- 一个 IDE,如 IntelliJ IDEA 或 Eclipse。
MyBatis 依赖
- 在
pom.xml中引入 MyBatis 核心依赖。本文示例基于 MyBatis 3.5.x 以上版本。
<dependency> <groupId>org.mybatis</groupId> <artifactId>mybatis</artifactId> <version>3.5.10</version> <!-- 请使用最新稳定版 --> </dependency>- 如果使用 Spring Boot,可以直接使用
mybatis-spring-boot-starter。
<dependency> <groupId>org.mybatis.spring.boot</groupId> <artifactId>mybatis-spring-boot-starter</artifactId> <version>2.3.0</version> </dependency>- 在
数据库与驱动
- 一个测试数据库(如 MySQL、PostgreSQL)。
- 对应的 JDBC 驱动依赖。
<!-- MySQL 示例 --> <dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>8.0.33</version> <scope>runtime</scope> </dependency>测试数据
- 准备一张数据量较大的表用于测试。可以编写脚本插入 50 万到 100 万条记录。表结构可以简单如下:
CREATE TABLE `large_user` ( `id` bigint(20) NOT NULL AUTO_INCREMENT, `name` varchar(255) DEFAULT NULL, `email` varchar(255) DEFAULT NULL, `created_at` datetime DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (`id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
4. 两种流式查询实现方式详解
MyBatis 主要提供了两种方式实现流式查询:使用Cursor<T>接口和使用ResultHandler。我们分别来看。
4.1 方式一:使用Cursor<T>接口(推荐)
Cursor提供了迭代器风格的 API,使用起来最直观,类似于遍历一个List,但背后是流式读取。
第一步:Mapper 接口定义在 Mapper 接口中,将返回类型定义为Cursor<YourEntity>。
import org.apache.ibatis.cursor.Cursor; public interface UserMapper { Cursor<User> selectAllUsersStreaming(); }第二步:XML 映射文件在对应的 XML 映射文件中,编写 SQL。这里不需要特殊的配置,MyBatis 会根据返回类型自动处理。
<!-- UserMapper.xml --> <mapper namespace="com.example.mapper.UserMapper"> <select id="selectAllUsersStreaming" resultType="com.example.entity.User"> SELECT id, name, email, created_at FROM large_user <!-- 可以添加 WHERE 条件,但注意流式查询最好有排序保证顺序 --> ORDER BY id ASC </select> </mapper>第三步:Service 层调用与遍历这是关键步骤,必须在一个数据库事务中完成遍历,并且确保最终关闭Cursor。
import org.apache.ibatis.cursor.Cursor; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import java.io.IOException; @Service public class UserService { private final UserMapper userMapper; public UserService(UserMapper userMapper) { this.userMapper = userMapper; } /** * 使用 @Transactional 确保整个遍历过程在一个事务内 * 事务保证了数据库连接和游标的一致性 */ @Transactional public void processAllUsersWithCursor() { // 获取 Cursor try (Cursor<User> cursor = userMapper.selectAllUsersStreaming()) { // 遍历 Cursor,每次迭代才会真正从数据库取数据 for (User user : cursor) { // 在这里处理每一条用户数据 // 例如:写入文件、发送消息、进行计算 System.out.println("Processing user: " + user.getName()); // 模拟处理耗时 // Thread.sleep(1); // 谨慎使用,会拖慢整体速度 } } catch (IOException e) { // Cursor 的 close 方法会抛出 IOException throw new RuntimeException("Failed to close cursor", e); } // try-with-resources 会自动调用 cursor.close(),释放资源 } }核心要点:
@Transactional注解至关重要。流式查询需要在整个遍历期间保持同一个数据库连接和事务上下文。- 使用
try-with-resources语句确保Cursor被自动关闭,释放底层ResultSet和数据库游标。 - 遍历
Cursor时,才真正从数据库获取数据。userMapper.selectAllUsersStreaming()本身执行很快,它只是返回了一个Cursor对象,并没有立即加载数据。
4.2 方式二:使用ResultHandler
ResultHandler是一个回调接口,MyBatis 在读取每一条结果时都会调用它。这种方式控制粒度更细,但代码稍显繁琐。
第一步:定义 ResultHandler 实现类
import org.apache.ibatis.session.ResultHandler; import com.example.entity.User; public class UserResultHandler implements ResultHandler<User> { private int count = 0; @Override public void handleResult(ResultContext<? extends User> resultContext) { // 获取当前结果对象 User user = resultContext.getResultObject(); // 处理当前对象 System.out.println("Handling user[" + (++count) + "]: " + user.getName()); // 可以通过 resultContext 停止处理 // if (count > 10000) { // resultContext.stop(); // } } public int getCount() { return count; } }第二步:Mapper 接口定义Mapper 方法返回类型为void,需要传入ResultHandler参数。
public interface UserMapper { void selectAllUsersWithHandler(ResultHandler<User> handler); }第三步:XML 映射文件与普通查询一致。
<mapper namespace="com.example.mapper.UserMapper"> <select id="selectAllUsersWithHandler" resultType="com.example.entity.User"> SELECT id, name, email, created_at FROM large_user ORDER BY id ASC </select> </mapper>第四步:Service 层调用同样需要在事务中执行。
@Service public class UserService { private final UserMapper userMapper; public UserService(UserMapper userMapper) { this.userMapper = userMapper; } @Transactional public void processAllUsersWithHandler() { UserResultHandler handler = new UserResultHandler(); // 执行查询,结果会通过回调给 handler userMapper.selectAllUsersWithHandler(handler); System.out.println("Total processed: " + handler.getCount()); } }两种方式对比:
Cursor<T>:更现代,代码更简洁,类似迭代器模式,推荐大多数场景使用。ResultHandler:更底层,可以在处理每条数据时获得ResultContext,有能力提前停止处理 (resultContext.stop()),适合需要复杂控制流的场景。
5. 功能测试与效果验证
理论讲完了,我们来实际测试一下,看看流式查询到底如何解决 OOM 问题。
5.1 测试准备:制造内存危机
首先,我们写一个会 OOM 的传统查询方法作为对比。
// UserMapper.java List<User> selectAllUsers(); // 传统方法 // UserService.java public void processAllUsersTraditional() { List<User> allUsers = userMapper.selectAllUsers(); // 一次性加载所有数据到内存! for (User user : allUsers) { System.out.println("Processing user: " + user.getName()); } System.out.println("Total users: " + allUsers.size()); }在application.yml中,我们为测试限制 JVM 堆内存,让问题更容易暴露。
# Spring Boot 测试配置 spring: datasource: url: jdbc:mysql://localhost:3306/test_db?useSSL=false&serverTimezone=UTC username: root password: yourpassword # 通过JVM参数限制堆内存更直接,例如 -Xmx256m5.2 测试一:传统查询 vs 流式查询(内存占用)
- 启动服务:使用
-Xmx256m(最大堆内存 256MB)启动你的 Spring Boot 应用。 - 执行传统方法:调用
processAllUsersTraditional()处理一个 50 万条记录的表。极大概率你会看到控制台输出java.lang.OutOfMemoryError: Java heap space,程序崩溃。 - 执行流式方法:重启服务,调用
processAllUsersWithCursor()。观察控制台,它会开始一条条输出,进程不会崩溃。使用 JConsole、VisualVM 或jcmd <pid> GC.heap_info工具观察堆内存使用情况,你会发现内存使用是一条平稳的曲线,峰值远低于 256MB,而不会出现一个陡峭的峰值。
判断成功:流式查询方法在有限内存下能完成全部数据处理,且内存占用平稳;传统方法抛出 OOM 异常。
5.3 测试二:结合业务逻辑(导出 CSV)
流式查询最常见的用途是数据导出。我们来模拟将百万用户导出为 CSV 文件。
import org.apache.commons.csv.CSVFormat; import org.apache.commons.csv.CSVPrinter; import org.springframework.transaction.annotation.Transactional; import java.io.FileWriter; import java.io.IOException; import java.io.Writer; @Service public class DataExportService { private final UserMapper userMapper; public DataExportService(UserMapper userMapper) { this.userMapper = userMapper; } @Transactional // 关键:保持事务 public void exportUsersToCsv(String filePath) throws IOException { // 使用 try-with-resources 管理文件写入器和 CSV 打印机 try (Writer writer = new FileWriter(filePath); CSVPrinter csvPrinter = new CSVPrinter(writer, CSVFormat.DEFAULT.withHeader("ID", "Name", "Email", "CreatedAt")); Cursor<User> cursor = userMapper.selectAllUsersStreaming()) { // 关键:获取流式 Cursor for (User user : cursor) { // 将每一条记录写入 CSV 行 csvPrinter.printRecord(user.getId(), user.getName(), user.getEmail(), user.getCreatedAt()); // 每处理10000条,刷新一下缓冲区,避免内存中积累太多字符串 // if (cursor.getCurrentIndex() % 10000 == 0) { // csvPrinter.flush(); // } } // 最终刷新并关闭 csvPrinter.flush(); } // 自动关闭 cursor, csvPrinter, writer System.out.println("CSV export completed to: " + filePath); } }效果验证:运行此方法,即使数据量很大,也不会导致 OOM。同时,由于是边读边写,生成的文件会逐渐变大,而不是等所有数据在内存中组装好才一次性写入。
5.4 测试三:性能与资源观察
- 内存:使用
jstat -gc <pid> 1000观察 GC 情况。流式查询下,Young GC 可能更频繁(因为不断创建和回收 User 对象),但几乎不会发生 Full GC。传统查询则会因为分配超大数组触发多次 Full GC 甚至直接 OOM。 - 数据库连接:通过数据库的
SHOW PROCESSLIST;命令,可以看到流式查询执行期间,会有一个连接长时间处于Sending data状态,直到遍历结束。这意味着连接被占用,强调了事务管理和及时关闭游标的重要性。 - 耗时对比:流式查询的总耗时可能略高于传统查询,因为多了多次网络往返的开销。但对于内存敏感的场景,用稍长的时间换取服务的稳定性是绝对值得的。
6. 高级配置与调优
要让流式查询工作得更好,可能需要一些额外的配置。
6.1 配置fetchSize
fetchSize是 JDBC 驱动的一个提示,告诉数据库每次网络往返返回多少条记录。设置一个合理的值可以优化性能。
<!-- 在 MyBatis 的 <select> 标签中配置 --> <select id="selectAllUsersStreaming" resultType="com.example.entity.User" fetchSize="1000"> SELECT id, name, email, created_at FROM large_user ORDER BY id ASC </select>- MySQL:需要连接参数
useCursorFetch=true,并且驱动版本要支持。在 JDBC URL 中添加:jdbc:mysql://...?useCursorFetch=true。然后fetchSize才会生效。 - PostgreSQL:原生支持,配置
fetchSize即可。 - Oracle:也支持,但可能有自己的语法要求。
6.2 配置resultSetType
设置结果集类型为FORWARD_ONLY(默认)和READ_ONLY,这是流式查询的典型配置。
<select id="selectAllUsersStreaming" resultType="com.example.entity.User" fetchSize="1000" resultSetType="FORWARD_ONLY"> SELECT id, name, email, created_at FROM large_user ORDER BY id ASC </select>6.3 连接池注意事项
以 HikariCP 为例,需要确保连接不会因为查询时间过长而被回收或检测为“泄漏”。
spring: datasource: hikari: connection-timeout: 60000 # 连接超时时间,流式查询可能较长,适当调大 max-lifetime: 1800000 # 连接最大生命周期,调大 leak-detection-threshold: 60000 # 泄漏检测阈值,如果流式查询处理时间可能超过此值,需要调大或关闭最佳实践:对于已知的长时间流式查询任务,可以考虑使用独立的、配置更宽松的数据源。
7. 常见问题与排查方法
在实际使用流式查询时,你可能会遇到以下问题。
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
抛出InvalidResultSetAccessException或游标相关错误 | 1. 遍历Cursor时不在事务中。2. 遍历 Cursor的代码和获取Cursor的代码不在同一个事务/连接中。 | 检查方法是否被@Transactional注解,且传播级别正确(通常是REQUIRED)。 | 确保整个遍历过程在一个事务内。将获取和遍历Cursor的代码放在同一个@Transactional方法中。 |
| 数据库连接耗尽或连接泄漏 | Cursor没有正确关闭,导致数据库连接和游标未被释放。 | 检查代码是否使用了try-with-resources或finally块确保cursor.close()被调用。查看数据库连接池监控。 | 强制使用try (Cursor c = mapper.query())语法。确保 Service 方法不被异常中断。 |
| 流式查询速度比普通查询慢很多 | 1.fetchSize设置过小,网络往返次数过多。2. 数据库端排序或没有索引,导致流式传输本身慢。 3. 每行数据处理逻辑太重。 | 1. 检查并调整fetchSize。2. 对 ORDER BY字段加索引。3. 分析处理逻辑的耗时。 | 1. 适当增大fetchSize(如 1000-5000)。2. 优化 SQL 和索引。 3. 考虑异步或批量处理单行数据。 |
MySQL 下fetchSize不生效 | JDBC URL 缺少useCursorFetch=true参数。 | 检查 JDBC 连接字符串。 | 在 MySQL JDBC URL 后添加&useCursorFetch=true。注意,这可能会在某些版本和场景下有副作用,需测试。 |
| 处理过程中程序中断,游标未关闭 | 程序崩溃、被杀死或发生未捕获异常。 | 查看应用日志和数据库的SHOW PROCESSLIST,可能看到大量Sleep状态的连接。 | 1. 加强代码健壮性,做好异常处理,在catch块中关闭资源。2. 设置合理的数据库 wait_timeout和interactive_timeout,让数据库端自动关闭超时空闲连接。 |
在 Web 请求中使用流式查询返回Cursor | 试图在 HTTP 请求响应中直接返回Cursor。 | Cursor的生命周期与数据库连接绑定,HTTP 请求结束后连接可能已关闭,导致后续遍历失败。 | 绝对禁止在 Controller 层返回Cursor。流式查询应在 Service 层完成所有数据处理,将最终结果(如文件路径、处理状态)返回给前端。 |
8. 最佳实践与使用建议
- 事务是必须的:这是流式查询能工作的前提,务必牢记。
- 资源必须关闭:使用
try-with-resources是关闭Cursor最安全、最简洁的方式。 - 先测试,后上线:在生产环境使用前,务必在测试环境用同等数据量进行充分测试,验证内存、性能和稳定性。
- 监控与告警:对使用流式查询的服务,加强对其数据库连接数、长时间查询的监控。
- 区分使用场景:明确你的场景是“需要遍历大量数据并处理每一条”,而不是“需要将所有数据加载到内存进行复杂计算”。前者用流式,后者可能需要考虑分页查询或移到大数据平台处理。
- SQL 优化:即使流式查询,SQL 本身也要高效。确保
WHERE条件和ORDER BY的字段有索引,避免全表扫描的流式查询,那将是灾难。 - 批处理思想:在
ResultHandler或遍历Cursor时,可以考虑每积累 N 条记录进行一次批量操作(如批量插入到另一个库),而不是每条都操作,这能显著提升吞吐量。 - 超时控制:对于可能执行时间非常长的流式查询任务,要有超时中断机制,防止一个任务永远占用连接。
9. 总结与下一步
MyBatis 流式查询是一个强大的工具,它能将你从“大数据量查询必 OOM”的魔咒中解救出来。它的核心价值在于用可控的内存开销,换取处理海量数据的能力。
最值得尝试的点:如果你的系统中存在任何将数据库大量数据导出为文件、同步到其他系统或进行逐条清洗转换的任务,并且这些任务曾引发过内存告警,那么流式查询应该是你的首选优化方案。
最先应该验证的功能:从一个简单的Cursor遍历开始,搭配@Transactional和try-with-resources,先确保基础流程能在你的开发环境跑通。
最容易踩的坑:
- 忘记加
@Transactional:这是新手最容易犯的错误,会导致游标无效。 - 在 Web 层返回或泄漏
Cursor:记住,流式处理必须在 Service 层完成闭环。 - 不关闭
Cursor:导致连接和游标泄漏,这是生产事故的隐患。
下一步可以探索的方向:
- 结合 Spring Batch:对于更复杂、需要容错、重启能力的大批量数据处理任务,可以将 MyBatis 流式查询作为 Spring Batch 的
ItemReader,构建健壮的批处理作业。 - 异步处理:在流式遍历过程中,将获取到的数据放入队列,由消费者异步处理,可以进一步提高吞吐量,避免流式查询本身成为瓶颈。
- 多数据源流式查询:从一个数据库流式读取,同时写入另一个数据库,实现高效的数据迁移。
通过本文,你应该已经掌握了 MyBatis 流式查询从原理到实战的全部要点。它不是一个复杂的技术,但却是构建稳健后端服务不可或缺的技能。建议你将文中的示例代码在本地运行一遍,亲眼见证它如何优雅地处理百万数据而不崩,这种体验远比阅读文字来得深刻。