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

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自动回复的类型绑定限制:

  1. 从入站消息中提取JMSReplyTo指定的目的地
  2. 用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&amp;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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 16:45:41