拓冰建站拓冰建站
首页 / 资讯中心 / 正文

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_sourcetaotoken_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 里提交到仓库。这一点在批处理任务里尤其重要因为跑批机器往往不止一台。如果你更习惯用命令行验证可以直接 curlcurl 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 Planhttps://taotoken.net/coding-plan 它更适合持续性的编码和 Agent 工作负载。本文的示例用按量调用就够了。3. 可复制配置fetchSize、事务边界与连接池三件套这一节是全文最核心的部分直接给可复制的配置。先说fetchSize它是流式查询的命门。fetchSize控制每次从数据库拉取多少条记录到客户端。设小了网络往返次数多性能差设大了单批内存占用高失去流式的意义。经验值MySQL 场景下 100 到 1000 之间比较稳PostgreSQL 必须配合autoCommitfalse才生效否则驱动会一次性拉全量。这一点很多人踩坑——在 PG 上设了fetchSize100却还是 OOM就是因为没关自动提交。XML 配置方式注意resultSetTypeFORWARD_ONLY是必须的只有单向游标才能流式select idselectFetchSize fetchSize500 resultSetTypeFORWARD_ONLY resultTypecom.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) QueryWrapperEntityDemo wrapper, ResultHandlerEntityDemo handler);Cursor 路线的 Mapper 方法则返回CursorTMapper public interface EntityDemoMapper extends BaseMapperEntityDemo { Select(select * from entity_demo limit #{limit}) CursorEntityDemo 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 (CursorEntityDemo 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: 60000leak-detection-threshold设成 60 秒一旦流式查询忘记关连接日志里会打出泄漏警告比等到连接池耗尽再排查强得多。另外跑批任务建议用独立的连接池或独立的 DataSource别和在线业务抢连接。ResultHandler 路线的完整写法Mapper 无返回值Service 里传回调Override public void streamGain() { QueryWrapperEntityDemo 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 MapString, Object testCursor(RequestParam int limit) { long start System.currentTimeMillis(); AtomicLong count new AtomicLong(); transactionTemplate.execute(status - { try (CursorEntityDemo cursor entityDemoMapper.scan(limit)) { cursor.forEach(item - { count.incrementAndGet(); // 模拟业务处理比如调用模型 }); } catch (IOException e) { throw new RuntimeException(e); } return null; }); MapString, Object result new HashMap(); result.put(count, count.get()); result.put(costMs, System.currentTimeMillis() - start); return result; }第三步启动时加上 JVM 参数方便观察 GCjava -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就说明链路通了可以放心把模型调用接进你的流式批处理管道。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门