1. 从一次 OOM 事故说起:100w 行数据到底该怎么读
线上有个对账任务,每天凌晨从 MySQL 拉一张流水表做汇总,表里大概 100w 行。最早写得很朴素:selectList一把梭,结果某天数据涨到 80w 之后,服务开始隔三差五抛java.lang.OutOfMemoryError: Java heap space,堆 dump 下来一看,一半内存全是BigDataSearchEntity对象。这就是典型的「大结果集读取」翻车现场。
先把概念说清楚:流式查询(Streaming Query)指的是查询执行后不把整个结果集塞进客户端内存,而是返回一个游标/迭代器,应用每次从里面取一条或一批;游标查询(Cursor Query)本质是流式的一种实现,通过fetchSize控制每次从服务端拉多少行到客户端。它适合谁?适合做数据迁移、数据导出、批量对账、跨库结果集合并这类「结果集很大但单条处理逻辑很轻」的场景。
为什么不能无脑分页?LIMIT 1000000, 20这种深分页,MySQL 要先扫描并丢弃前 100w 行,越翻越慢,翻到后面基本等于全表扫。而一次性selectList的问题更直接:MyBatis 会把每一行映射成实体对象,100w 个对象加上 ResultSet 缓冲,堆内存直接爆炸。所以面试官问「100w 数据怎么处理」,真正想听的是你对内存模型和连接生命周期的理解,而不是背一句「用流式查询」。
这篇我会按「复现 OOM → 配置流式 → 验证内存 → 排错」的顺序走一遍,所有配置都能直接抄。中间涉及模型调用和长任务编排时,我会用 TaoToken 做统一入口,把 API Key、Base URL、Model ID 三件套讲清楚,避免你在多个平台之间来回切。
2. 前置准备:TaoToken 接入与依赖配置
在动手写流式查询之前,先把工程环境和调用入口理顺。TaoToken 在这里扮演的是「统一模型网关」的角色:你写代码时只需要认一个 Base URL 和一把 Key,模型切换、额度查看、调用日志都在一个控制台里完成。官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 地址是 https://taotoken.net/api (这个不加 UTM)。
先拿 Key:进控制台 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite ,在 API Keys 页面 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite 新建一把,复制出来存到环境变量,别硬编码进代码。
export TAOTOKEN_API_KEY="sk-你的key" export TAOTOKEN_BASE_URL="https://taotoken.net/api"Maven 依赖这块,流式查询只需要 MyBatis 和 MySQL 驱动,版本别太老:
<dependency> <groupId>org.mybatis.spring.boot</groupId> <artifactId>mybatis-spring-boot-starter</artifactId> <version>3.0.3</version> </dependency> <dependency> <groupId>com.mysql</groupId> <artifactId>mysql-connector-j</artifactId> <version>8.3.0</version> </dependency>数据源配置里有个关键点:MySQL 驱动要允许服务端游标。连接串加上useCursorFetch=true,否则fetchSize在某些版本下不生效,驱动会偷偷把整个结果集读回来。
spring: datasource: url: jdbc:mysql://127.0.0.1:3306/demo?useCursorFetch=true&useSSL=false&serverTimezone=Asia/Shanghai username: root password: your_password hikari: maximum-pool-size: 10如果你用 Cline MCP 或 Claude Code 这类工具辅助写代码,配置里同样要写全三件套:Base URL 填https://taotoken.net/api,Key 填上面那把,Model ID 按你选的模型填。Coding Plan 适合长期跑这类数据任务的场景,入口在 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。模型对话调试可以用 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 先验证连通性。
3. 可复制配置:JDBC fetchSize 与 MyBatis Cursor 写法
这一节是核心,直接给能跑的代码。先看错误示范,也就是会 OOM 的那种:
// 危险:一次性加载 100w 行 List<BigDataSearchEntity> list = bigDataSearchMapper.selectList(null); for (BigDataSearchEntity e : list) { process(e); }正确姿势有三种,按内存占用从低到高排。
方式一:MyBatis Cursor 逐行处理(推荐)
Mapper 接口返回Cursor<T>,注意方法上不加@Options的 fetchSize 也行,但建议显式声明:
@Mapper public interface BigDataSearchMapper { @Select("SELECT id, order_no, amount, create_time FROM big_data_search") @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = Integer.MIN_VALUE) Cursor<BigDataSearchEntity> streamAll(); }fetchSize = Integer.MIN_VALUE是 MySQL 驱动的约定,表示「逐行流式读取」,等价于服务端游标模式。Service 层必须用 try-with-resources 包住,否则连接不释放:
@Service public class BigDataService { @Autowired private BigDataSearchMapper mapper; public void handleAll() { try (Cursor<BigDataSearchEntity> cursor = mapper.streamAll()) { cursor.forEach(this::process); } catch (IOException e) { throw new RuntimeException("cursor close failed", e); } } private void process(BigDataSearchEntity e) { // 单条处理,别在这里再查库 } }方式二:ResultHandler 回调(一次查询,逐行回调)
@Select("SELECT id, order_no, amount FROM big_data_search") @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 1000) @ResultType(BigDataSearchEntity.class) void listData(ResultHandler<BigDataSearchEntity> handler);调用时传一个 lambda,返回类型必须是void:
mapper.listData(ctx -> { BigDataSearchEntity e = ctx.getResultObject(); process(e); });方式三:原生 JDBC 控制 fetchSize(最直观)
try (Connection conn = dataSource.getConnection(); PreparedStatement ps = conn.prepareStatement( "SELECT id, order_no, amount FROM big_data_search", ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY)) { ps.setFetchSize(Integer.MIN_VALUE); // MySQL 流式 try (ResultSet rs = ps.executeQuery()) { while (rs.next()) { process(rs.getLong("id"), rs.getString("order_no")); } } }三种方式的内存差异用一张表对照更清楚:
| 方式 | 内存增长 | 连接占用 | 适用场景 |
|---|---|---|---|
| selectList 一次性 | 随行数线性增长 | 查询完即释放 | 小结果集(<1w) |
| Cursor 逐行 | 基本恒定 | 全程独占 | 百万级导出/迁移 |
| ResultHandler | 基本恒定 | 全程独占 | 需要复用 MyBatis 映射 |
| 分页 LIMIT | 每页恒定 | 每页释放 | 能走索引的浅分页 |
注意:流式查询期间这条数据库连接是独占的,处理逻辑里千万不要再去查同一数据源,否则会死等连接池。要么把处理结果攒批后异步写,要么单独配一个小连接池给流式任务。
4. 验证请求:内存对比与成功结果
配置写完,得用数据证明它真的省内存。我一般用两段代码做对照实验,JVM 参数固定-Xmx256m,表里灌 100w 行。
对照组:一次性加载
long start = System.currentTimeMillis(); List<BigDataSearchEntity> list = mapper.selectList(null); System.out.println("size=" + list.size() + ", cost=" + (System.currentTimeMillis() - start) + "ms");跑起来大概率在size打印之前就抛 OOM,堆日志里能看到java.lang.OutOfMemoryError: Java heap space,这就是复现成功。
实验组:Cursor 流式
long start = System.currentTimeMillis(); AtomicLong count = new AtomicLong(); try (Cursor<BigDataSearchEntity> cursor = mapper.streamAll()) { cursor.forEach(e -> { count.incrementAndGet(); if (count.get() % 100000 == 0) { Runtime rt = Runtime.getRuntime(); long used = (rt.totalMemory() - rt.freeMemory()) / 1024 / 1024; System.out.println("已处理 " + count.get() + " 行, 堆占用 " + used + "MB"); } }); } System.out.println("total=" + count.get() + ", cost=" + (System.currentTimeMillis() - start) + "ms");实测下来,堆占用会稳定在 40–80MB 之间波动,不会随行数上涨,100w 行跑完大概几十秒到几分钟,取决于单条处理逻辑。控制台输出类似:
已处理 100000 行, 堆占用 52MB 已处理 200000 行, 堆占用 55MB ... total=1000000, cost=48213ms如果你还想验证模型侧调用是否正常,可以用模型对话页面发一条测试请求,确认 Base URL 和 Key 没问题:https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 。接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite ,里面有各语言的示例。
5. 常见报错排查:401、local proxy failed 与 reading choices
流式任务跑起来后,报错往往不在 SQL 本身,而在调用链路上。下面几个是我踩过的坑,对照着查。
报错一:401 Unauthorized
{"error":{"message":"Invalid API key","type":"invalid_request_error"}}原因通常是 Key 没读到环境变量,或者复制时带了空格。检查echo $TAOTOKEN_API_KEY,确认以sk-开头且无换行。如果是在 Cline MCP 或 Claude Code 里配置,三件套要写全:Base URLhttps://taotoken.net/api、Key、Model ID,缺一个都会 401 或 404。
报错二:local proxy failed
Error: local proxy failed: dial tcp 127.0.0.1:7890: connect: connection refused这是本地代理端口没起来,但配置里还写着代理地址。把HTTP_PROXY/HTTPS_PROXY环境变量清掉,或者确认代理服务在跑。注意别把这类代理和数据库连接混在一起排查,两者没关系。
报错三:reading choices 相关
Error: reading choices: unexpected end of JSON input这种一般是流式响应被中途截断,常见于网络抖动或超时设置过短。把客户端超时调大,比如 120s,并确认响应体是完整的 SSE 流。如果是 Codex 的auth.json配置问题,检查里面的base_url和api_key字段是否和上面三件套一致。
报错四:Cursor 用完没关导致连接池耗尽
HikariPool-1 - Connection is not available, request timed out after 30000ms九成是Cursor没放进 try-with-resources,或者forEach里抛异常后没走到 close。养成习惯:只要返回类型是Cursor,外面必须包 try-with-resources。
报错五:fetchSize 不生效,还是 OOM
检查连接串有没有useCursorFetch=true,以及resultSetType是不是FORWARD_ONLY。SCROLL_INSENSITIVE会让驱动把结果集缓存到客户端,等于白配。
提示:排障时优先看接入文档的 FAQ 部分,大部分 401 和超时问题都有现成答案:https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite
6. 把流式查询接进你的日常工具链
最后说下怎么把这套东西固化下来。如果你用 Claude Code 做辅助开发,配置入口在 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_content=claude-code-anthropic&utm_campaign=rewrite ,把 Base URL、Key、Model ID 填进去,之后让它帮你生成 Mapper 和 Cursor 处理模板会快很多。长期跑数据迁移、对账这类 Agent 任务,Coding Plan 的额度模型更划算,入口前面给过了。
几个实用技巧收尾:流式任务尽量放在独立线程池,别和 Web 请求抢连接;处理逻辑里攒批写库,比如每 1000 条 flush 一次,比逐条 insert 快一个数量级;如果结果集要跨库合并,先在各库用 Cursor 拉,再在内存里做归并,别把两个大 List 同时驻留。至于 100w 行到底选 Cursor 还是分页,判断标准很简单:能走索引的浅分页用分页,深分页或全表扫描一律上 Cursor。