Apache Camel:将拆分元素聚合为列表的实现问题
解决Camel自定义聚合器聚合位置列表的问题
我来帮你搞定这个聚合逻辑的问题,毕竟把单个处理后的结果聚合成列表是Camel里很常见的场景,咱们一步步来:
首先明确核心:你的PositionsAggregation需要正确实现Camel的AggregationStrategy接口,这个接口的aggregate方法是聚合逻辑的核心——它要负责把之前聚合的结果和新进来的单个位置合并成一个列表。
1. 正确实现自定义聚合器
这里给你一个完整的实现示例,包含常见的边界情况处理:
import org.apache.camel.Exchange; import org.apache.camel.processor.aggregate.AggregationStrategy; import java.util.ArrayList; import java.util.List; public class PositionsAggregation implements AggregationStrategy { @Override public Exchange aggregate(Exchange oldExchange, Exchange newExchange) { // 从新消息里提取刚处理好的位置(这里假设是String类型,换成你的实际类型就行) String newPosition = newExchange.getIn().getBody(String.class); // 跳过空的位置,避免列表里出现无效数据,可根据业务需求调整 if (newPosition == null || newPosition.isBlank()) { return oldExchange; } List<String> positionsList; if (oldExchange == null) { // 这是第一条进来的消息,初始化一个新列表 positionsList = new ArrayList<>(); positionsList.add(newPosition); // 把初始化后的列表放回新消息的body,作为初始聚合结果 newExchange.getIn().setBody(positionsList); return newExchange; } else { // 从之前的聚合结果里拿到已有的列表 positionsList = oldExchange.getIn().getBody(List.class); // 添加新位置到列表 positionsList.add(newPosition); // 更新旧消息的body,作为最新的聚合结果 oldExchange.getIn().setBody(positionsList); return oldExchange; } } }
2. 路由里的关键修正
你的原路由里有个容易踩的坑:aggregate(body(), new PositionsAggregation())这里用body()作为聚合关联键,意味着每个不同的位置值会被分到不同的聚合组,最后你会得到多个小列表而不是一个完整的列表!
要把所有位置聚合到同一个列表,你需要用一个固定的常量作为关联键,修改后的路由如下:
from("direct:cages-to-positions") .process(new CageToPositionProcessor()) // 用固定常量作为聚合键,确保所有消息都进入同一个聚合组 .aggregate(constant("ALL_POSITIONS"), new PositionsAggregation()) .completionTimeout(1000) // 1秒内无新消息就完成聚合 .to("mock:test");
3. 额外优化建议
- 如果你的位置是自定义对象,把代码里的
String换成对应的实体类即可,逻辑完全通用。 - 如果路由涉及并发处理,建议用线程安全的列表比如
CopyOnWriteArrayList替换ArrayList,避免并发修改问题。 - 可以在聚合完成后加一个处理器,对最终的列表做收尾(比如去重、排序、过滤空值等),根据你的业务需求调整。
测试验证
用Camel的测试框架发送多个笼子消息,然后校验mock:test收到的消息body是不是包含所有位置的完整列表,就能确认聚合逻辑是否正常工作了。
内容的提问来源于stack exchange,提问作者Neok
相关产品推荐
相关产品推荐

