1. 从一次定时任务超时说起:Java MongoDB Aggregation 复杂查询到底难在哪
如果你正在用 Java + Spring Data MongoDB 做定时任务,每隔一两秒扫一次集合,筛选出「已过期」「已失效」「该被回收」的文档,那你大概率踩过和我一样的坑:用@Query或者find()加一堆Criteria拼条件,代码越写越长,最后还得在 JVM 里再遍历一遍做时间比较和标记判断。数据量一上来,服务端 CPU 直接飙红。
这个场景的核心矛盾在于:筛选逻辑本应该在数据库侧完成,却被搬到了应用侧。MongoDB 的 Aggregation Pipeline 就是为这种「多阶段、带计算、带关联」的查询而生的。它允许你在$project阶段动态计算字段(比如用createdTime + childTaskAlivePeriod * 1000算出过期时间),在$match阶段用$or/$and组合多个条件,甚至用$lookup做跨集合关联、$unwind拆数组、$group做聚合统计。
但问题也随之而来:管道 JSON 嵌套层级深,手写容易出错;用原生com.mongodb驱动写BasicDBObject层层嵌套,几十行代码看得头皮发麻;换成 Spring Data 的Aggregation.newAggregation()虽然优雅,但表达式语法、Criteria组合、OutputMode配置又各有各的坑。更现实的是,当你想在开发阶段用模型辅助生成或校验这些管道时,还需要一个稳定的 API 通道来调用模型能力。
这篇文章就围绕「Java + Spring Data MongoDB 的 Aggregation 多阶段管道落地」展开,给出可复制的管道 JSON、Spring Data 代码片段,以及通过 TaoToken 统一 Key 完成模型侧调用配置的完整步骤。适合正在做定时任务、数据清洗、跨集合关联查询的后端同学。
2. TaoToken 统一 Key 前置准备:为 Aggregation 开发配一条模型调用通道
在正式写管道之前,先说清楚为什么这个场景需要 TaoToken。你在开发 Aggregation 时,经常会遇到几类需要模型辅助的时刻:把一段复杂的原生管道 JSON 翻译成 Spring Data 的Aggregation对象;检查$lookup的localField/foreignField是否写反;让模型帮你把$group的_id分组键和累加器解释清楚。这些操作如果每次都去网页端复制粘贴,效率很低,而通过统一的 API 通道接入,就能在 IDE 或脚本里直接调用。
TaoToken 在这里扮演的是「统一 Key / API 通道」的角色。你只需要一个 Key,就能通过兼容的接口协议调用模型对话能力,用于生成、校验、解释 Aggregation 管道。它的官网入口是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 基址是 https://taotoken.net/api (注意这个地址不加 UTM 参数)。
具体操作上,你需要先拿到 Key。进入控制台创建 API Key,路径是 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite ,在 API Keys 页面生成一个新的 Key 并保存好。这个 Key 就是你后续所有模型调用的凭证。
拿到 Key 之后,建议先做一次连通性验证,确认通道可用。你可以用 curl 发一个最简单的对话请求:
curl -X POST https://taotoken.net/api/v1/chat/completions \ -H "Content-Type: application/json" \ -H "Authorization: Bearer 你的_API_KEY" \ -d '{ "model": "claude-sonnet-4-20250514", "messages": [ {"role": "user", "content": "用一句话解释 MongoDB 的 $lookup 阶段作用"} ] }'如果返回的 JSON 里有choices字段且内容正常,说明 Key 和通道都没问题。这一步很关键,因为后面你在写 Aggregation 时如果让模型帮忙生成管道,走的就是这个通道。模型对话的入口在 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite ,你可以在那里先手动试几条 prompt,确认模型对 MongoDB 聚合管道的理解符合预期。
对于长期做编码和 Agent 场景的同学,如果调用频率高,可以考虑 Coding Plan,入口是 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite ,里面有完整的接口说明和参数列表。
这里要强调一点:TaoToken 是模型调用的统一通道,不是数据库代理,也不替代你的 MongoDB 或 Spring Data。它的价值在于让你在开发 Aggregation 的过程中,有一个稳定的模型侧能力可以随时调用,用来生成管道、排查语法、解释阶段行为。
3. 可复制配置:Spring Data Aggregation 管道 JSON 与 Java 代码片段
这一节是全文的核心,直接给你能复制运行的配置和代码。我们以一个典型的「定时任务扫描过期文档」场景为例,集合名为RefreshTask,需要筛选出满足以下任一条件的文档:expireFlag为 true;或者collectState为 1 且nextCollectTime小于当前时间;或者动态计算出的expireTime(等于createdTime + childTaskAlivePeriod * 1000)小于当前时间。
先看原生的 Aggregation 管道 JSON,这是你在 MongoDB Shell 或 Compass 里可以直接跑的版本:
[ { "$project": { "_id": 1, "uri": 1, "scope": 1, "desc": 1, "pushState": 1, "collectState": 1, "collectInterval": 1, "lastCollectTime": 1, "refreshRequestHeaders": 1, "effectivePercent": 1, "hasDataVersion": 1, "childTaskAlivePeriod": 1, "parentId": 1, "createdTime": 1, "lastModifiedTime": 1, "lastCollectValue": 1, "lastPushValue": 1, "expireFlag": 1, "nextCollectTime": 1, "expireTime": { "$add": [ "$createdTime", { "$multiply": ["$childTaskAlivePeriod", 1000] } ] } } }, { "$match": { "$or": [ { "expireFlag": { "$eq": true } }, { "$and": [ { "collectState": { "$eq": 1 } }, { "nextCollectTime": { "$lt": "$$NOW" } } ] }, { "expireTime": { "$lt": "$$NOW" } } ] } } ]注意$$NOW是 MongoDB 4.2+ 支持的聚合变量,如果你用的是 3.6 版本,需要替换成具体的ISODate()值。这个管道分两个阶段:$project负责保留字段并计算expireTime,$match负责用$or组合三个筛选条件。
接下来是 Spring Data MongoDB 的等价写法。你需要引入spring-boot-starter-data-mongodb,然后在 Service 里注入MongoTemplate:
import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.data.mongodb.core.aggregation.Aggregation; import org.springframework.data.mongodb.core.aggregation.AggregationResults; import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.stereotype.Service; import java.util.Date; import java.util.List; @Service public class RefreshTaskService { private final MongoTemplate mongoTemplate; public RefreshTaskService(MongoTemplate mongoTemplate) { this.mongoTemplate = mongoTemplate; } public List<RefreshTaskDO> findExpiredTasks() { Date now = new Date(); Aggregation aggregation = Aggregation.newAggregation( Aggregation.project( "_id", "uri", "scope", "desc", "pushState", "collectState", "collectInterval", "lastCollectTime", "refreshRequestHeaders", "effectivePercent", "hasDataVersion", "childTaskAlivePeriod", "parentId", "createdTime", "lastModifiedTime", "lastCollectValue", "lastPushValue", "expireFlag", "nextCollectTime" ) .andExpression("createdTime + childTaskAlivePeriod * 1000") .as("expireTime"), Aggregation.match( new Criteria().orOperator( Criteria.where("expireFlag").is(true), new Criteria().andOperator( Criteria.where("collectState").is(1), Criteria.where("nextCollectTime").lt(now) ), Criteria.where("expireTime").lt(now) ) ) ); AggregationResults<RefreshTaskDO> results = mongoTemplate.aggregate(aggregation, "RefreshTask", RefreshTaskDO.class); return results.getMappedResults(); } }这段代码的关键点有三个。第一,Aggregation.project()里用.andExpression()做字段计算,语法是 SpEL 风格的表达式,createdTime + childTaskAlivePeriod * 1000会被翻译成$add和$multiply。第二,Aggregation.match()里用new Criteria().orOperator()组合多个条件,嵌套的andOperator对应管道里的$and。第三,mongoTemplate.aggregate()直接返回映射好的实体列表,不需要手动遍历DBObject再 set 字段。
如果你需要把这段配置写成外部化的 JSON 文件(比如放在resources/aggregation/refresh-task-pipeline.json),可以用Aggregation.parse()或者直接读文件后转成Document列表再构造Aggregation。不过对于大多数场景,上面的 Java 代码已经足够清晰。
另外,如果你的管道里涉及$lookup、$unwind、$group,Spring Data 也都有对应的方法:
Aggregation.newAggregation( Aggregation.lookup("order", "orderId", "_id", "orderInfo"), Aggregation.unwind("orderInfo"), Aggregation.group("scope").count().as("total").sum("effectivePercent").as("avgPercent"), Aggregation.sort(Sort.Direction.DESC, "total"), Aggregation.limit(100) );$lookup的四个参数依次是:从集合名、本地字段、外部字段、输出数组字段名。$unwind把数组拆成多行。$group的_id是分组键,后面跟累加器。这套组合在跨集合统计场景里非常常用。
4. 验证请求与结果比对:确认管道真的按预期执行
写完代码不代表管道就对了。你需要一套验证方法,确认生成的 Mongo 语句和预期一致,并且返回结果正确。
第一步,打印生成的管道。Spring Data 的Aggregation对象可以转成Document查看:
Document pipeline = aggregation.toPipeline(Aggregation.DEFAULT_CONTEXT); System.out.println(pipeline.toJson());或者更直接地,在MongoTemplate上开启 DEBUG 日志,Spring Data 会把最终发给 MongoDB 的命令打印出来。你会在日志里看到类似这样的内容:
{ "aggregate": "RefreshTask", "pipeline": [ { "$project": { ... "expireTime": { "$add": ["$createdTime", { "$multiply": ["$childTaskAlivePeriod", 1000] }] } } }, { "$match": { "$or": [ ... ] } } ], "cursor": { "batchSize": 2147483647 } }把这个 JSON 复制到 MongoDB Shell 里直接执行,对比返回的文档数量和字段值。如果 Shell 里返回 5 条,Java 代码里getMappedResults().size()也应该是 5。如果对不上,说明映射或条件有问题。
第二步,做边界值比对。构造几条测试数据:一条expireFlag=true,一条collectState=1且nextCollectTime是过去时间,一条expireTime计算后小于当前时间,再构造一条都不满足的。跑一遍管道,确认只返回前三条。这种「正例 + 反例」的比对能快速暴露$or/$and的括号嵌套错误。
第三步,用模型辅助校验。把生成的管道 JSON 贴给模型,让它检查$lookup的字段是否匹配、$group的_id是否遗漏、$unwind是否会导致数据膨胀。通过 TaoToken 的模型对话入口 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 可以直接发起这类校验请求。实测下来,模型对$lookup的localField/foreignField写反、$group累加器用错这类问题识别得比较准。
第四步,性能比对。在数据量较大的集合上,分别跑「应用侧过滤」和「Aggregation 管道过滤」,用explain()看执行计划。管道里的$match如果放在$project之前,能更早地减少文档数量;但如果$match依赖$project计算出的字段,就必须放在后面。这个顺序直接影响性能,值得多试几次。
5. 本篇常见错误排查:从 401 到 cursor 报错逐个击破
这一节整理我在这个场景里真实踩过的坑,以及对应的报错和解决方式。
报错一:com.mongodb.MongoCommandException: Command failed with error 9: 'The 'cursor' option is required, except for aggregate with the explain argument'
这个报错在 MongoDB 3.6+ 配合旧版驱动时非常常见。原因是AggregationOptions.builder()默认的OutputMode是INLINE,而 3.6 之后的版本要求聚合必须返回游标。解决方式是把OutputMode改成CURSOR:
AggregationOptions options = AggregationOptions.builder() .outputMode(AggregationOptions.OutputMode.CURSOR) .build();如果你用的是 Spring Data,MongoTemplate.aggregate()内部已经处理了这个问题,一般不会遇到。但如果你混用了原生DBCollection.aggregate(),就要注意这个配置。
报错二:401 Unauthorized或local proxy failed
这类报错通常出现在调用模型 API 时。401 说明 Key 无效或没带上,检查Authorization: Bearer头里的 Key 是否和控制台里生成的一致。local proxy failed一般是本地网络配置问题,确认你的请求地址是https://taotoken.net/api而不是其他地址。如果是在 IDE 插件里配置,注意 Base URL 要填完整路径,不要漏掉/api。
报错三:reading choices相关解析错误
当你用脚本解析模型返回时,如果报reading choices或choices is undefined,说明返回的 JSON 结构和你预期的不一样。先打印完整响应体,确认choices字段存在。常见原因是请求体里model参数写错,或者messages格式不对。正确的请求体结构是:
{ "model": "claude-sonnet-4-20250514", "messages": [{"role": "user", "content": "..."}] }报错四:OAuth相关认证失败
如果你在 Claude Code 或类似工具里配置,遇到 OAuth 报错,检查是不是把 API Key 和 OAuth Token 混用了。TaoToken 的 API Key 走的是 Bearer 认证,不需要额外的 OAuth 流程。Claude Code 的接入配置里,Base URL 填https://taotoken.net/api,Key 填控制台生成的 API Key,Model ID 填你实际使用的模型名。这三件套缺一不可。
报错五:$lookup返回空数组
管道语法没错,但$lookup结果全是空数组。检查localField和foreignField的类型是否一致。MongoDB 不会做隐式类型转换,如果本地字段是String而外部字段是ObjectId,关联就会失败。用$toString或$toObjectId在$project阶段做转换。
报错六:$group后字段丢失
$group阶段只会保留_id和累加器字段,其他字段都会丢。如果你需要保留原始字段,要么在$group前用$first取,要么在$group后用$lookup回查。这是聚合管道的设计逻辑,不是 bug。
6. 语义一致 CTA:把 Key 配好,让 Aggregation 开发更顺
回到整个流程。你现在已经有了可复制的管道 JSON、Spring Data 的 Java 代码、验证方法和排错清单。剩下的就是把模型调用通道配好,让后续的管道生成和校验更顺手。
如果你主要是在排障和接入阶段,建议先去 API Keys 页面生成 Key,路径是 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite ,然后对照接入文档 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 完成配置。文档里有完整的 Base URL、认证方式和请求示例。
如果你需要先验证模型对 MongoDB 聚合管道的理解是否符合预期,可以直接在模型对话页面 https://taotoken.net/models?utm_source=taotoken_aicg_blog_end&utm_content=models&utm_campaign=rewrite 发几条测试 prompt,比如「把这段 $lookup 管道翻译成 Spring Data 代码」,确认输出质量后再接入到自己的开发流程里。
如果你长期做编码和 Agent 场景,调用频率高,Coding Plan 会更合适,入口是 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。它针对编码场景做了优化,适合把模型能力集成到日常开发工作流里。
最后提醒一点:Aggregation 管道的调试,最有效的方法永远是「打印生成的 Mongo 语句 + 在 Shell 里直接跑 + 对比结果」。模型可以帮你生成和校验,但最终的验证还是要靠真实数据。把 Key 配好,把管道跑通,剩下的就是根据业务需求调整$match条件和$project字段了。