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

Apache Camel路由问题:转换后的值未被转发,仍发送原始消息

问题原因与解决方案

问题根源

你用到的.multicast()是Apache Camel的Multicast EIP,它默认会**复制原始Exchange(包含初始接收的消息体)**发送给每个子处理器分支。这就导致你的.bean(PricingLifeCyclebService.class, "paraMap")处理后的消息体仅在该子分支内有效,后续的.to(...)步骤仍然使用原始消息体,而非转换后的值。


解决方案

方案1:移除不必要的Multicast

如果你的路由不需要并行处理多个分支(当前代码里只有一条处理链路),直接删除.multicast()和.parallelProcessing()即可:

from(azureServicebus(AZvalue)
        .connectionString(connectionString)
        .receiverAsyncClient(serviceBusReceiverAsyncClient)
        .serviceBusReceiveMode(ServiceBusReceiveMode.RECEIVE_AND_DELETE)
        .serviceBusType(ServiceBusType.topic)
        .prefetchCount(100)
        .consumerOperation(ServiceBusConsumerOperationDefinition.receiveMessages)
        //.maxAutoLockRenewDuration(Duration.ofMinutes(10))
)
.messageHistory()
.routeId(Endpoints.SEDA_PROCESS_SB_MESSAGE_ENDPOINT)
.bean(PricingLifeCyclebService.class, "paraMap") 
.log("Final :- ${body}")
.to(Endpoints.SEDA_SEND_PRICING_LIFE_CYCLE_MESSAGE);

方案2:保留Multicast并共享上下文(若确实需要)

如果后续要扩展多分支并行处理,需要给Multicast添加.shareUnitOfWork()配置,让所有子分支共享同一个Exchange上下文,这样分支内对消息体的修改会同步到主Exchange:

from(azureServicebus(AZvalue)
        .connectionString(connectionString)
        .receiverAsyncClient(serviceBusReceiverAsyncClient)
        .serviceBusReceiveMode(ServiceBusReceiveMode.RECEIVE_AND_DELETE)
        .serviceBusType(ServiceBusType.topic)
        .prefetchCount(100)
        .consumerOperation(ServiceBusConsumerOperationDefinition.receiveMessages)
        //.maxAutoLockRenewDuration(Duration.ofMinutes(10))
)
.messageHistory()
.routeId(Endpoints.SEDA_PROCESS_SB_MESSAGE_ENDPOINT)
.multicast().shareUnitOfWork() // 共享工作单元,同步消息体修改
.parallelProcessing()
.bean(PricingLifeCyclebService.class, "paraMap") 
.log("Final :- ${body}")
.to(Endpoints.SEDA_SEND_PRICING_LIFE_CYCLE_MESSAGE)
.end();

额外提示

如果后续SEDA端点接收消息时出现序列化问题,可以在paraMap方法里手动将LifeCycleTopicC对象转为JSON字符串:

public String paraMap(String object) throws JsonProcessingException { 
    final var paraMapper = LifeCycleUtility.getObjectMapper();
    final var paraDTO = paraMapper.readValue(object, ParaDTO.class);
    
    LifecycleDTO lifecycleDTO = paraToLifecleDTOMapper.mapToLifecycleDTO(paraDTO);
    LifeCycleTopicC obj = lifecycleTopicmapper.toLifeCycleTopicC(lifecycleDTO);
    
    log.info("Event={}, Body={}", "ConvertParatoTopic", obj);
    return paraMapper.writeValueAsString(obj); // 返回JSON字符串
}

内容的提问来源于stack exchange,提问作者John

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 19:00:26