如何使用Apache Camel向同一接收方发送多条Kafka消息?
解决Apache Camel向同一接收方发送多条消息的问题
你的问题出在默认路由的处理逻辑上:顺序调用两个bean方法时,第二个方法会覆盖第一个方法设置的Exchange Out Body,最终recipientList只会发送最后一次设置的那条消息。下面提供两种可行的解决方案:
方案一:使用Multicast组件复制Exchange发送多条消息
Multicast会为每个后续处理器复制原始Exchange,每个处理器的输出都会独立发送到目标端点,避免Body被覆盖:
from(sourceEndpoint) .routeId(sourceEndpoint) .unmarshal().json(JsonLibrary.Jackson) .choice() .when(criteria) .multicast() .bean(xyz.class, "processAndSendMessage") .bean(xyz.class, "processAndSendAnotherMessage") .to(destinationEndpoint);
这里Multicast会分别执行两个bean方法,每个方法处理后的Exchange都会被发送到destinationEndpoint,实现一次触发发送两条消息的效果。
方案二:生成消息列表后用Splitter拆分发送
修改bean逻辑,让某个方法返回包含多条消息的集合,再通过Splitter拆分集合为单独的Exchange发送:
- 新增或修改bean方法,返回消息列表:
public List<Object> generateMessages(Exchange exchange) { List<Object> messages = new ArrayList<>(); // 调用原有逻辑生成第一条消息 Object msg1 = processAndSendMessage(exchange); messages.add(msg1); // 调用原有逻辑生成第二条消息 Object msg2 = processAndSendAnotherMessage(exchange); messages.add(msg2); return messages; }
- 调整路由配置,用Splitter拆分列表后发送:
from(sourceEndpoint) .routeId(sourceEndpoint) .unmarshal().json(JsonLibrary.Jackson) .choice() .when(criteria) .bean(xyz.class, "generateMessages") .split(body()) .recipientList(simple(destinationEndpoint));
Splitter会将列表中的每个元素拆分为独立的Exchange,每个Exchange都会通过recipientList发送到目标端点。
内容的提问来源于stack exchange,提问作者ankitkpd
相关产品推荐
相关产品推荐

