拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

MyBatis 流式查询处理千万级数据:TaoToken 统一 Key 下的 Cursor 与 ResultHandler 实战

MyBatis 流式查询处理千万级数据:TaoToken 统一 Key 下的 Cursor 与 ResultHandler 实战

1. 千万级数据导出为什么会 OOM:从一次真实事故说起

先说结论:MyBatis 流式查询(Streaming Query)指的是查询成功后不返回List,而是返回一个迭代器或通过回调逐条消费结果,应用每次从结果集取一条记录。它能做什么?把千万级数据导出、批处理、对账这类场景的堆内存占用从「随数据量线性增长」压到「基本恒定」。适合谁?正在用 MyBatis / MyBatis-Plus 做数据导出、跑批、跨库同步,并且已经被OutOfMemoryError: Java heap space折磨过的后端同学。

我见过最典型的事故是这样:一个对账任务要导出 1200 万行订单明细,开发同学写了个selectList(wrapper),本地 10 万行跑得飞快,上线后堆内存 4G,跑到第 300 万行左右直接 OOM,服务重启,任务重跑,又 OOM,循环往复。根因不复杂——普通查询会把整个结果集一次性映射成 Java 对象塞进ArrayList,这些对象在方法返回前都是强引用,GC 根本回收不掉,堆内存被撑爆只是时间问题。

有人会说,那我分页查不就行了?分页当然可以,但深分页(limit 10000000, 1000)在 MySQL 上要先扫描并丢弃前面一千万行,效率取决于表设计和索引,设计不好就是全表扫,越翻越慢。流式查询绕开了这个问题:它保持一个数据库连接打开,服务端游标逐批吐数据,客户端逐条消费,用完即弃,内存曲线是平的。

这里有个关键前提必须提前说清楚:流式查询期间数据库连接是保持打开的,框架不负责帮你关,需要应用在取完数据后自己关闭。这也是后面A Cursor is already closed报错的根源。理解了这一点,Cursor 和 ResultHandler 两条路线的差异就很好懂了。

本文会交付三样东西:一是 Cursor 与 ResultHandler 的资源占用对比和适用边界;二是可直接复制的fetchSize、事务边界、连接池配置片段;三是压测验证步骤。同时,因为现在很多团队在批处理链路里会顺带调用大模型做数据清洗、摘要、分类,我会演示怎么用 TaoToken 的统一 Key 管理多模型调用,避免在代码里散落一堆厂商 Key。TaoToken 官网是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end ,API 入口是 https://taotoken.net/api ,后面配置里会用到。

先给一个直观的对照,帮你决定选哪条路:

维度CursorResultHandler
返回形态迭代器,业务侧主动forEach回调,框架推给你
连接持有需事务或 SqlSession 包裹同样需事务包裹
内存占用低,逐条消费低,逐条消费
代码侵入中,要处理 try-with-resources低,Mapper 方法无返回值
适用边界需要自己控制节奏、可中断纯消费、逻辑内聚在回调里
典型坑Cursor is already closed回调里抛异常导致事务回滚

这张表先放这儿,下面逐段拆开讲。

2. TaoToken 统一 Key 前置准备:批处理链路里的多模型调用怎么管

在讲配置之前,先把 TaoToken 这一层说清楚,因为后面的示例里会用到它。TaoToken 是一个统一的大模型 API 接入层,能做什么?你用一套 Key、一个 Base URL,就能调用多家模型,不用为每个厂商单独维护 Key、单独改 SDK 初始化代码。适合谁?批处理 / 数据管道里需要调用模型做清洗、打标、摘要,又不想把七八个厂商 Key 硬编码进配置文件的团队。

为什么批处理场景特别需要它?想象一下你的千万级数据导出任务,导出后要对每条记录做一次意图分类。如果直接对接某一家厂商,Key 泄露风险、限流、单点故障都压在你身上;如果对接多家做降级,代码里就会出现一堆if provider == "a" ... else if provider == "b"。TaoToken 把这层抽象掉了,你只面对一个 OpenAI 兼容的接口。

前置准备分三步。第一步,拿到统一 Key。访问 https://taotoken.net/api-keys ,登录后在控制台创建 API Key,复制出来形如sk-xxxxxxxx。这个 Key 就是你在所有模型调用里唯一需要配置的凭证。

第二步,确认 Base URL。TaoToken 的 API 入口是 https://taotoken.net/api ,注意这里不带任何查询参数,SDK 里配置的base_url就填这个。如果你用的是 OpenAI 官方 SDK,它会自动拼接/v1/chat/completions这类路径。

第三步,选模型。你可以在 https://taotoken.net/models 查看当前可用的模型列表,也可以直接在 https://taotoken.net/chat 里对话验证某个模型是否可用、效果是否符合预期。批处理里常用的做法是:便宜模型做粗筛,贵模型做精修,两者都通过同一个 Key 调用。

这里给一个 Spring Boot 里配置 TaoToken 客户端的片段,用application.yml管理,避免硬编码:

taotoken: base-url: https://taotoken.net/api api-key: ${TAOTOKEN_API_KEY} default-model: your-cheap-model-id refine-model: your-strong-model-id timeout-seconds: 60 max-retries: 3

对应的配置类:

@Configuration @ConfigurationProperties(prefix = "taotoken") @Data public class TaoTokenProperties { private String baseUrl; private String apiKey; private String defaultModel; private String refineModel; private int timeoutSeconds = 60; private int maxRetries = 3; }

注意api-key用环境变量注入,不要写死在 yml 里提交到仓库。这一点在批处理任务里尤其重要,因为跑批机器往往不止一台。

如果你更习惯用命令行验证,可以直接 curl:

curl https://taotoken.net/api/v1/chat/completions \ -H "Authorization: Bearer $TAOTOKEN_API_KEY" \ -H "Content-Type: application/json" \ -d '{ "model": "your-model-id", "messages": [{"role": "user", "content": "ping"}] }'

返回里能看到choices[0].message.content,说明 Key 和 Base URL 都通了。这一步先做,别等到流式查询跑起来才发现模型调用 401,那时候排查成本翻倍。

关于长期跑批和 Agent 场景,如果你打算把模型调用做成常驻服务而不是一次性脚本,可以了解下 Coding Plan:https://taotoken.net/coding-plan ,它更适合持续性的编码和 Agent 工作负载。本文的示例用按量调用就够了。

3. 可复制配置:fetchSize、事务边界与连接池三件套

这一节是全文最核心的部分,直接给可复制的配置。先说fetchSize,它是流式查询的命门。

fetchSize控制每次从数据库拉取多少条记录到客户端。设小了,网络往返次数多,性能差;设大了,单批内存占用高,失去流式的意义。经验值:MySQL 场景下 100 到 1000 之间比较稳,PostgreSQL 必须配合autoCommit=false才生效,否则驱动会一次性拉全量。这一点很多人踩坑——在 PG 上设了fetchSize=100却还是 OOM,就是因为没关自动提交。

XML 配置方式,注意resultSetType="FORWARD_ONLY"是必须的,只有单向游标才能流式:

<select id="selectFetchSize" fetchSize="500" resultSetType="FORWARD_ONLY" resultType="com.example.poi.entity.EntityDemo"> select * from entity_demo </select>

注解方式,注意 Mapper 方法必须没有返回值,这是 ResultHandler 路线的硬性要求:

@Select("select * from entity_demo t ${ew.customSqlSegment}") @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize = 500) @ResultType(EntityDemo.class) void selectFetchSize(@Param(Constants.WRAPPER) QueryWrapper<EntityDemo> wrapper, ResultHandler<EntityDemo> handler);

Cursor 路线的 Mapper 方法则返回Cursor<T>:

@Mapper public interface EntityDemoMapper extends BaseMapper<EntityDemo> { @Select("select * from entity_demo limit #{limit}") Cursor<EntityDemo> scan(@Param("limit") int limit); }

接下来是事务边界,这是最容易出错的地方。Cursor 必须在事务内消费,否则 Mapper 方法一返回连接就关了,游标跟着失效。三种方案,我推荐TransactionTemplate,因为它对「只在外部调用时生效」这个注解坑免疫:

@Resource private EntityDemoMapper entityDemoMapper; @Resource private TransactionTemplate transactionTemplate; public void streamWithCursor(int limit) { transactionTemplate.execute(status -> { try (Cursor<EntityDemo> cursor = entityDemoMapper.scan(limit)) { cursor.forEach(item -> { // 逐条业务处理,用完即弃 process(item); }); } catch (IOException e) { throw new RuntimeException("cursor consume failed", e); } return null; }); }

如果你用@Transactional注解,务必确认调用方是另一个 Bean,同类内部调用注解不生效,游标照样关闭。这个坑我在两个项目里都见过。

连接池配置同样关键。流式查询会长时间占用一个连接,如果连接池最大连接数太小,跑批任务会把连接池占满,导致其他接口拿不到连接。HikariCP 的推荐配置:

spring: datasource: hikari: maximum-pool-size: 20 minimum-idle: 5 connection-timeout: 30000 max-lifetime: 1800000 idle-timeout: 600000 leak-detection-threshold: 60000

leak-detection-threshold设成 60 秒,一旦流式查询忘记关连接,日志里会打出泄漏警告,比等到连接池耗尽再排查强得多。另外,跑批任务建议用独立的连接池或独立的 DataSource,别和在线业务抢连接。

ResultHandler 路线的完整写法,Mapper 无返回值,Service 里传回调:

@Override public void streamGain() { QueryWrapper<EntityDemo> wrapper = new QueryWrapper<>(); entityDemoMapper.selectFetchSize(wrapper, resultContext -> { EntityDemo row = resultContext.getResultObject(); process(row); }); }

注意 ResultHandler 回调里如果抛异常,整个事务会回滚,已经处理的数据不会自动补偿。所以回调里要么做幂等,要么把异常吞掉记录到死信表,别让它冒泡。

最后补一个多模型调用的配置片段,把 TaoToken 的 Key 和模型 ID 一起管起来,避免散落:

{ "taotoken": { "baseUrl": "https://taotoken.net/api", "apiKeyEnv": "TAOTOKEN_API_KEY", "models": { "classify": "your-cheap-model-id", "summarize": "your-strong-model-id" } } }

三件套齐了:fetchSize控制批次,事务边界保证游标存活,连接池防止连接耗尽。下面验证。

4. 验证请求与成功结果:压测步骤和内存曲线怎么看

配置写完不算完,得验证。这一节给一套可复制的压测步骤,从 10 万行到 1000 万行逐级放大,观察内存和耗时。

第一步,准备测试数据。用存储过程或脚本灌 1000 万行到entity_demo表,字段别太少,至少 10 个字段,模拟真实宽度。灌数据时关掉 binlog 或调大innodb_flush_log_at_trx_commit,否则灌数据本身就要跑很久。

第二步,写一个验证接口,分别用 Cursor 和 ResultHandler 跑同一份数据,记录处理条数和耗时:

@GetMapping("/stream/cursor") public Map<String, Object> testCursor(@RequestParam int limit) { long start = System.currentTimeMillis(); AtomicLong count = new AtomicLong(); transactionTemplate.execute(status -> { try (Cursor<EntityDemo> cursor = entityDemoMapper.scan(limit)) { cursor.forEach(item -> { count.incrementAndGet(); // 模拟业务处理,比如调用模型 }); } catch (IOException e) { throw new RuntimeException(e); } return null; }); Map<String, Object> result = new HashMap<>(); result.put("count", count.get()); result.put("costMs", System.currentTimeMillis() - start); return result; }

第三步,启动时加上 JVM 参数,方便观察 GC:

java -Xms2g -Xmx2g \ -XX:+UseG1GC \ -XX:+PrintGCDetails \ -Xloggc:gc.log \ -jar your-app.jar

第四步,逐级压测。先跑 10 万行,确认功能通;再跑 100 万行,看内存是否平稳;最后跑 1000 万行,重点看三件事:堆内存是否稳定在某个水位不再上涨、GC 频率是否正常、任务总耗时是否可接受。

成功的结果长这样:1000 万行数据,堆内存稳定在 1.2G 左右不再增长,Young GC 每隔几秒一次,没有 Full GC,总耗时 8 到 15 分钟(取决于单条处理逻辑)。如果堆内存持续上涨直到 OOM,说明流式没生效,回去检查resultSetType和事务边界。

第五步,验证模型调用链路。在process方法里插入一次 TaoToken 调用,确认批处理过程中模型调用不报错:

private void process(EntityDemo row) { // 业务处理 String text = row.getContent(); if (text != null && text.length() > 10) { String label = taotokenClient.classify(text); row.setLabel(label); } }

跑 1000 行验证即可,别一上来就 1000 万行调模型,成本和耗时都不可控。确认链路通了再放大。

这里有个实测经验:流式查询的耗时瓶颈往往不在数据库,而在单条业务处理逻辑。如果process里有一次同步 HTTP 调用(比如调模型),1000 万行就是 1000 万次调用,哪怕每次 50ms,总耗时也是 138 小时。所以批处理里调模型一定要做批量聚合,比如攒 100 条一起调,或者用异步 + 限流。这一点比流式查询本身更容易被忽略。

5. 本篇常见错排查:401、Cursor is already closed、OOM 怎么定位

这一节按真实报错来,每个报错给现象、根因、修法。

报错一:java.lang.IllegalStateException: A Cursor is already closed.

现象:调用cursor.forEach时抛这个异常。根因:Mapper 方法执行完连接就关了,游标跟着失效。修法:用TransactionTemplate或SqlSessionFactory.openSession()把整个消费过程包在事务里,确保消费期间连接不释放。注意@Transactional同类内部调用不生效这个坑。

报错二:401 Unauthorized或invalid api key

现象:调用 TaoToken 接口返回 401。根因:Key 没配、配错、或者环境变量没注入。修法:先确认TAOTOKEN_API_KEY环境变量在当前进程可见,再确认 Base URL 是 https://taotoken.net/api 而不是别的路径。用 curl 单独验证一次,排除代码问题。如果 curl 通、代码不通,检查 SDK 是否自动拼接了/v1,导致路径变成/api/v1/v1/...。

报错三:OutOfMemoryError: Java heap space依旧出现

现象:明明用了流式查询还是 OOM。根因通常有三个:一是fetchSize没生效(PG 没关 autoCommit,或 MySQL 驱动版本太老);二是消费过程中把数据攒进了另一个 List,比如cursor.forEach(list::add),等于白流式;三是resultSetType没设成FORWARD_ONLY。修法:逐个排查,重点看消费逻辑里有没有隐式集合累积。

报错四:local proxy failed或连接超时

现象:批处理跑一半连接断开。根因:连接池max-lifetime小于数据库wait_timeout,或者网络中间层掐断长连接。修法:把max-lifetime设成小于数据库wait_timeout,并开启leak-detection-threshold观察是否有连接泄漏。流式查询持有连接时间长,这个配置比普通查询更敏感。

报错五:reading choices相关解析失败

现象:模型返回体解析报错,提示读不到choices字段。根因:Base URL 配错导致返回的不是标准 OpenAI 格式,或者模型 ID 写错返回了错误结构。修法:先用 curl 看原始返回,确认结构里有choices数组。TaoToken 是 OpenAI 兼容格式,正常返回一定有choices。如果返回的是错误对象,先解决错误码。

报错六:OAuth 或鉴权相关错误

现象:提示鉴权失败、token 无效。根因:Key 过期、复制时带了空格、或者用了错误的鉴权头。修法:重新在 https://taotoken.net/api-keys 生成 Key,注意Authorization: Bearer <key>格式,Bearer 后面有一个空格。复制 Key 时别把首尾空白带进去。

排查顺序建议固定下来:先 curl 验证 Key 和 Base URL,再验证单条模型调用,再验证 100 行流式查询,最后放大到千万级。逐级验证能把问题定位到具体环节,比一上来跑全量高效得多。

6. 把流式查询和统一 Key 固化成团队规范

走到这里,你已经有了完整的三件套配置、压测步骤和排错清单。最后说点工程化的东西,让这套方案在团队里能复用。

第一,把流式查询封装成模板方法。Cursor 的 try-with-resources 和事务包裹是重复代码,抽成一个StreamQueryTemplate,业务侧只传消费逻辑,避免每个新人都重新踩一遍Cursor is already closed。

第二,把 TaoToken 的 Key 和模型 ID 收敛到配置中心,不要散落在各个跑批脚本里。批处理任务往往由不同人维护,Key 散落意味着轮换时漏改、泄露时找不到源头。统一走环境变量或配置中心,配合 https://taotoken.net/api-keys 的 Key 管理,轮换成本最低。

第三,给流式查询加监控。记录每次流式任务的消费条数、耗时、峰值堆内存、模型调用次数和失败率。这些指标能帮你在 OOM 之前发现问题,而不是等告警。

第四,批处理里调模型一定要做批量聚合和限流。单条调用在千万级数据下不可行,攒批 + 异步 + 重试是标配。TaoToken 的统一接口让切换模型变得容易,但调用模式的设计仍然是你自己的责任。

如果你还在选型阶段,建议先用 10 万行数据把 Cursor 和 ResultHandler 各跑一遍,对比代码复杂度和资源占用,再决定用哪个。多数纯消费场景 ResultHandler 更简洁,需要精细控制节奏或中途中断的场景 Cursor 更灵活。选定了就固化成模板,别每次重新发明。

最后留一个可以直接跑的验证命令,确认你的 TaoToken 配置在批处理环境里可用:

curl -s https://taotoken.net/api/v1/chat/completions \ -H "Authorization: Bearer $TAOTOKEN_API_KEY" \ -H "Content-Type: application/json" \ -d '{"model":"your-model-id","messages":[{"role":"user","content":"ok"}]}' \ | head -c 300

返回里有choices就说明链路通了,可以放心把模型调用接进你的流式批处理管道。

返回列表