Camel使用JMSReplyTo时如何发送不同类型的响应POJO?
我有一条Camel路由,监听ActiveMQ队列并处理消息,希望将结果发送到JMSReplyTo头指定的队列作为响应。旧版本需要显式设置目的地,通过ProducerTemplate.sendBody()发送;但根据Camel文档,JMSReplyTo会被自动识别,因此尝试不设置目的地,却发现Camel要求响应体必须是入站消息的POJO类型。
我尝试在路由最后执行exchange.getIn().setBody(myString);或exchange.getMessage().setBody(myString);,但出现如下错误:
Caused by: org.apache.camel.NoTypeConversionAvailableException: No type converter available to convert from type: java.lang.String to the required type: com.example.IncomingMessage with value {"failedEndpointsAndCauses":{},"message":"Message does not contain some ID field","status":"ERROR"}
这个String正是我想作为响应的内容。IncomingMessage是从ActiveMQ接收的入站POJO,显然Camel期望响应POJO与入站类型一致。
getOut().setBody()已被废弃无法使用。参考相关问题提到设置disableReplyTo,但这不符合请求-响应EIP的预期,我需要发送不同类型的响应体。尝试返回自定义OutboundMessage POJO时,也出现类似的类型转换错误,完整堆栈信息如下:
2022-12-20 14:09:49,293 DEBUG [org.apa.cam.com.jms.DefaultJmsMessageListenerContainer] (Camel (camel-1) thread #6 - JmsConsumer[my.queue]) Initiating transaction rollback on application exception: org.apache.camel.CamelExecutionException: Exception occurred during execution on the exchange: Exchange[] at org.apache.camel.CamelExecutionException.wrapCamelExecutionException(CamelExecutionException.java:45) at org.apache.camel.support.builder.ExpressionBuilder$33.evaluate(ExpressionBuilder.java:1030) at org.apache.camel.support.ExpressionAdapter.evaluate(ExpressionAdapter.java:45) at org.apache.camel.component.bean.MethodInfo$ParameterExpression.evaluateParameterBinding(MethodInfo.java:738) at org.apache.camel.component.bean.MethodInfo$ParameterExpression.evaluateParameterExpressions(MethodInfo.java:624) at org.apache.camel.component.bean.MethodInfo$ParameterExpression.evaluate(MethodInfo.java:592) at org.apache.camel.component.bean.MethodInfo.initializeArguments(MethodInfo.java:263) at org.apache.camel.component.bean.MethodInfo.createMethodInvocation(MethodInfo.java:271) at org.apache.camel.component.bean.BeanInfo.createInvocation(BeanInfo.java:277) at org.apache.camel.component.bean.AbstractBeanProcessor.process(AbstractBeanProcessor.java:126) at org.apache.camel.component.bean.BeanProcessor.process(BeanProcessor.java:81) at org.apache.camel.processor.errorhandler.NoErrorHandler.process(NoErrorHandler.java:47) at org.apache.camel.impl.engine.CamelInternalProcessor.process(CamelInternalProcessor.java:398) at org.apache.camel.processor.Pipeline$PipelineTask.run(Pipeline.java:109) at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187) at org.apache.camel.impl.engine.DefaultReactiveExecutor.scheduleMain(DefaultReactiveExecutor.java:64) at org.apache.camel.processor.Pipeline.process(Pipeline.java:184) at org.apache.camel.impl.engine.CamelInternalProcessor.process(CamelInternalProcessor.java:398) at org.apache.camel.impl.engine.DefaultAsyncProcessorAwaitManager.process(DefaultAsyncProcessorAwaitManager.java:83) at org.apache.camel.support.AsyncProcessorSupport.process(AsyncProcessorSupport.java:41) at org.apache.camel.component.jms.EndpointMessageListener.onMessage(EndpointMessageListener.java:132) at org.springframework.jms.listener.AbstractMessageListenerContainer.doInvokeListener(AbstractMessageListenerContainer.java:736) at org.springframework.jms.listener.AbstractMessageListenerContainer.invokeListener(AbstractMessageListenerContainer.java:696) at org.springframework.jms.listener.AbstractMessageListenerContainer.doExecuteListener(AbstractMessageListenerContainer.java:674) at org.springframework.jms.listener.AbstractPollingMessageListenerContainer.doReceiveAndExecute(AbstractPollingMessageListenerContainer.java:318) at org.springframework.jms.listener.AbstractPollingMessageListenerContainer.receiveAndExecute(AbstractPollingMessageListenerContainer.java:245) at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.invokeListener(DefaultMessageListenerContainer.java:1237) at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.executeOngoingLoop(DefaultMessageListenerContainer.java:1227) at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.run(DefaultMessageListenerContainer.java:1120) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.camel.InvalidPayloadException: No body available of type: com.example.IncomingMessage but has value: OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={}) of type: com.example.OutboundMessage on: Message[ID:foo-40805-1671538168705-4:1:1:4:1]. Caused by: No type converter available to convert from type: com.example.OutboundMessage to the required type: com.example.IncomingMessage with value OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={}). Exchange[]. Caused by: [org.apache.camel.NoTypeConversionAvailableException - No type converter available to convert from type: com.example.OutboundMessage to the required type: com.example.IncomingMessage with value OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={})] at org.apache.camel.support.MessageSupport.getMandatoryBody(MessageSupport.java:125) at org.apache.camel.support.builder.ExpressionBuilder$33.evaluate(ExpressionBuilder.java:1028) ... 30 more Caused by: org.apache.camel.NoTypeConversionAvailableException: No type converter available to convert from type: com.example.OutboundMessage to the required type: com.example.IncomingMessage with value OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={}) at org.apache.camel.impl.converter.CoreTypeConverterRegistry.mandatoryConvertTo(CoreTypeConverterRegistry.java:275) at org.apache.camel.support.MessageSupport.getMandatoryBody(MessageSupport.java:123) ... 31 more 2022-12-20 14:09:49,298 DEBUG [org.apa.cam.com.jms.DefaultJmsMessageListenerContainer] (Camel (camel-1) thread #6 - JmsConsumer[my.queue]) Rolling back transaction because of listener exception thrown: org.apache.camel.CamelExecutionException: Exception occurred during execution on the exchange: Exchange[] 2022-12-20 14:09:49,299 WARN [org.apa.cam.com.jms.EndpointMessageListener] (Camel (camel-1) thread #6 - JmsConsumer[my.queue]) Execution of JMS message listener failed. Caused by: [org.apache.camel.CamelExecutionException - Exception occurred during execution on the exchange: Exchange[]]: org.apache.camel.CamelExecutionException: Exception occurred during execution on the exchange: Exchange[] at org.apache.camel.CamelExecutionException.wrapCamelExecutionException(CamelExecutionException.java:45) at org.apache.camel.support.builder.ExpressionBuilder$33.evaluate(ExpressionBuilder.java:1030) at org.apache.camel.support.ExpressionAdapter.evaluate(ExpressionAdapter.java:45) at org.apache.camel.component.bean.MethodInfo$ParameterExpression.evaluateParameterBinding(MethodInfo.java:738) at org.apache.camel.component.bean.MethodInfo$ParameterExpression.evaluateParameterExpressions(MethodInfo.java:624) at org.apache.camel.component.bean.MethodInfo$ParameterExpression.evaluate(MethodInfo.java:592) at org.apache.camel.component.bean.MethodInfo.initializeArguments(MethodInfo.java:263) at org.apache.camel.component.bean.MethodInfo.createMethodInvocation(MethodInfo.java:271) at org.apache.camel.component.bean.BeanInfo.createInvocation(BeanInfo.java:277) at org.apache.camel.component.bean.AbstractBeanProcessor.process(AbstractBeanProcessor.java:126) at org.apache.camel.component.bean.BeanProcessor.process(BeanProcessor.java:81) at org.apache.camel.processor.errorhandler.NoErrorHandler.process(NoErrorHandler.java:47) at org.apache.camel.impl.engine.CamelInternalProcessor.process(CamelInternalProcessor.java:398) at org.apache.camel.processor.Pipeline$PipelineTask.run(Pipeline.java:109) at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:187) at org.apache.camel.impl.engine.DefaultReactiveExecutor.scheduleMain(DefaultReactiveExecutor.java:64) at org.apache.camel.processor.Pipeline.process(Pipeline.java:184) at org.apache.camel.impl.engine.CamelInternalProcessor.process(CamelInternalProcessor.java:398) at org.apache.camel.impl.engine.DefaultAsyncProcessorAwaitManager.process(DefaultAsyncProcessorAwaitManager.java:83) at org.apache.camel.support.AsyncProcessorSupport.process(AsyncProcessorSupport.java:41) at org.apache.camel.component.jms.EndpointMessageListener.onMessage(EndpointMessageListener.java:132) at org.springframework.jms.listener.AbstractMessageListenerContainer.doInvokeListener(AbstractMessageListenerContainer.java:736) at org.springframework.jms.listener.AbstractMessageListenerContainer.invokeListener(AbstractMessageListenerContainer.java:696) at org.springframework.jms.listener.AbstractMessageListenerContainer.doExecuteListener(AbstractMessageListenerContainer.java:674) at org.springframework.jms.listener.AbstractPollingMessageListenerContainer.doReceiveAndExecute(AbstractPollingMessageListenerContainer.java:318) at org.springframework.jms.listener.AbstractPollingMessageListenerContainer.receiveAndExecute(AbstractPollingMessageListenerContainer.java:245) at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.invokeListener(DefaultMessageListenerContainer.java:1237) at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.executeOngoingLoop(DefaultMessageListenerContainer.java:1227) at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.run(DefaultMessageListenerContainer.java:1120) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.camel.InvalidPayloadException: No body available of type: com.example.IncomingMessage but has value: OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={}) of type: com.example.OutboundMessage on: Message[ID:foo-40805-1671538168705-4:1:1:4:1]. Caused by: No type converter available to convert from type: com.example.OutboundMessage to the required type: com.example.IncomingMessage with value OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={}). Exchange[]. Caused by: [org.apache.camel.NoTypeConversionAvailableException - No type converter available to convert from type: com.example.OutboundMessage to the required type: com.example.IncomingMessage with value OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={})] at org.apache.camel.support.MessageSupport.getMandatoryBody(MessageSupport.java:125) at org.apache.camel.support.builder.ExpressionBuilder$33.evaluate(ExpressionBuilder.java:1028) ... 30 more Caused by: org.apache.camel.NoTypeConversionAvailableException: No type converter available to convert from type: com.example.OutboundMessage to the required type: com.example.IncomingMessage with value OutboundMessage(jobResourceId=null, status=something, message=Message does not contain some ID, failedEndpointsAndCauses={}) at org.apache.camel.impl.converter.CoreTypeConverterRegistry.mandatoryConvertTo(CoreTypeConverterRegistry.java:275) at org.apache.camel.support.MessageSupport.getMandatoryBody(MessageSupport.java:123) ... 31 more 2022-12-20 14:09:49,873 SEVERE [org.ecl.yas.int.Unmarshaller] (awaitility-thread) Unexpected char 111 at (line no=1, column no=2, offset=1), expecting 'u'
方法1:显式使用ProducerTemplate发送响应
保留请求-响应模式的同时,绕过Camel自动回复的类型绑定限制:
- 从入站消息中提取
JMSReplyTo指定的目的地 - 用
ProducerTemplate直接发送自定义类型的响应到该目的地,同时传递JMSCorrelationID关联请求与响应
示例代码:
// 获取JMSReplyTo目的地 Destination replyTo = exchange.getIn().getHeader(JmsConstants.JMS_REPLY_TO, Destination.class); if (replyTo != null) { producerTemplate.send(replyTo, ex -> { // 设置自定义响应体(String或OutboundMessage都可) ex.getIn().setBody(myString); // 关联请求消息ID,确保响应可被匹配 ex.getIn().setHeader(JmsConstants.JMS_CORRELATION_ID, exchange.getIn().getHeader(JmsConstants.JMS_MESSAGE_ID)); return ex; }); }
方法2:修改JMS端点配置解除类型绑定
在Camel的JMS消费者端点添加参数,关闭强制类型转换并明确请求-响应模式:
forceTypeConversion=false:禁止Camel强制将响应体转换为入站消息类型replyToType=RequestReply:明确启用请求-响应模式但不限制响应类型
XML DSL示例
<from uri="activemq:queue:my.queue?forceTypeConversion=false&replyToType=RequestReply"/>
Java DSL示例
from("activemq:queue:my.queue?forceTypeConversion=false&replyToType=RequestReply") .process(exchange -> { // 直接设置自定义响应体 exchange.getMessage().setBody(new OutboundMessage("ERROR", "缺失ID字段")); });
方法3:自定义类型转换器(临时妥协方案)
若必须依赖Camel自动回复机制,可创建自定义转换器让Camel能将响应类型转换为入站类型,但此方法仅为满足框架要求,业务逻辑上可能无意义:
@Converter public class CustomTypeConverter { @Converter public static IncomingMessage toIncomingMessage(OutboundMessage outbound) { IncomingMessage incoming = new IncomingMessage(); incoming.setMessage(outbound.getMessage()); // 按需补充转换逻辑 return incoming; } @Converter public static IncomingMessage toIncomingMessage(String json) { ObjectMapper mapper = new ObjectMapper(); try { return mapper.readValue(json, IncomingMessage.class); } catch (JsonProcessingException e) { throw new RuntimeException(e); } } }
注册转换器到Camel上下文:
context.getTypeConverterRegistry().addTypeConverter(IncomingMessage.class, OutboundMessage.class, new CustomTypeConverter()); context.getTypeConverterRegistry().addTypeConverter(IncomingMessage.class, String.class, new CustomTypeConverter());
内容的提问来源于stack exchange,提问作者WesternGun

