You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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发送:

  1. 新增或修改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;
}
  1. 调整路由配置,用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.23 05:24:30