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
相关产品推荐
相关产品推荐

