使用Qpid JMS向ActiveMQ Artemis发送大消息失败求助
我尝试从磁盘读取大文件并发送至ActiveMQ Artemis的队列,但每次操作都会抛出相同异常:
Caused by: jakarta.jms.MessageFormatException: Only objectified primitive objects and String types are allowed but was: java.io.BufferedInputStream@3a8489ac type: class java.io.BufferedInputStream at org.apache.qpid.jms.message.JmsMessagePropertySupport.checkValidObject(JmsMessagePropertySupport.java:121) at org.apache.qpid.jms.message.JmsMessagePropertyIntercepter.setProperty(JmsMessagePropertyIntercepter.java:725) at org.apache.qpid.jms.message.JmsMessage.setObjectProperty(JmsMessage.java:348) at org.apache.camel.component.jms.JmsBinding.createJmsMessageForType(JmsBinding.java:661) at org.apache.camel.component.jms.JmsBinding.createJmsMessage(JmsBinding.java:568) at org.apache.camel.component.jms.JmsBinding.createJmsMessage(JmsBinding.java:525) at org.apache.camel.component.jms.JmsBinding.makeJmsMessage(JmsBinding.java:348) at org.apache.camel.component.jms.JmsProducer$2.createMessage(JmsProducer.java:325) at org.apache.camel.component.jms.JmsConfiguration$CamelJmsTemplate.doSendToDestination(JmsConfiguration.java:616) at org.apache.camel.component.jms.JmsConfiguration$CamelJmsTemplate.lambda$send$0(JmsConfiguration.java:574) at org.springframework.jms.core.JmsTemplate.execute(JmsTemplate.java:507) ... 39 more
原因分析
经调试发现,Camel的JmsBinding组件在启用Artemis流处理(artemisStreamingEnabled)时,会执行message.setObjectProperty("JMS_AMQ_InputStream", is);代码,将InputStream对象设为JMS消息属性。但Qpid JMS客户端的checkValidObject方法仅允许基本类型(Boolean、Byte、Short、Integer、Long、Float、Double、Character)和String类型作为消息属性,InputStream类型不符合要求,因此抛出该异常。
相关核心代码逻辑:
- Camel的流式处理逻辑:
if (endpoint.isArtemisStreamingEnabled()) { LOG.trace("Optimised for Artemis: Streaming payload in BytesMessage"); InputStream is = context.getTypeConverter().mandatoryConvertTo(InputStream.class, exchange, body); message.setObjectProperty("JMS_AMQ_InputStream", is); LOG.trace("Optimised for Artemis: Finished streaming payload in BytesMessage"); } else { // 常规消息处理逻辑 }
- Qpid JMS的属性校验逻辑:
public static void checkValidObject(Object value) throws MessageFormatException { boolean valid = value instanceof Boolean || value instanceof Byte || value instanceof Short || value instanceof Integer || value instanceof Long || value instanceof Float || value instanceof Double || value instanceof Character || value instanceof String || value == null; if (!valid) { throw new MessageFormatException("Only objectified primitive objects and String types are allowed but was: " + value + " type: " + value.getClass()); } }
解决方案
方案1:关闭Artemis流处理特性
在Camel的JMS端点配置中,将artemisStreamingEnabled设置为false。此时Camel会将文件内容作为常规BytesMessage的负载发送,而非尝试通过消息属性传递InputStream。
Java DSL示例:
from("file:/path/to/your/files") .to("jms:queue:yourTargetQueue?artemisStreamingEnabled=false");
XML配置示例:
<route> <from uri="file:/path/to/your/files"/> <to uri="jms:queue:yourTargetQueue?artemisStreamingEnabled=false"/> </route>
方案2:替换为Artemis原生JMS客户端
JMS_AMQ_InputStream是ActiveMQ Artemis原生客户端支持的扩展属性,专门用于实现消息内容的流式传输。如果必须使用流式处理来优化大文件发送性能,可以将Qpid JMS客户端替换为Artemis原生的artemis-jms-client依赖,它支持将InputStream作为该属性传递,不会触发属性类型校验异常。
方案3:自定义Camel JmsBinding实现流式负载传输
继承JmsBinding类,重写createJmsMessageForType方法,修改流式处理逻辑,将InputStream直接写入BytesMessage的负载,而非设置为消息属性。
自定义Binding示例:
public class CustomArtemisJmsBinding extends JmsBinding { public CustomArtemisJmsBinding() { super(); } @Override protected Message createJmsMessageForType(Message message, Object body, Exchange exchange, Session session) throws JMSException { JmsEndpoint endpoint = (JmsEndpoint) exchange.getFromEndpoint(); if (endpoint != null && endpoint.isArtemisStreamingEnabled()) { InputStream is = exchange.getContext().getTypeConverter().mandatoryConvertTo(InputStream.class, exchange, body); BytesMessage bytesMessage = session.createBytesMessage(); byte[] buffer = new byte[4096]; int readBytes; try { while ((readBytes = is.read(buffer)) != -1) { bytesMessage.writeBytes(buffer, 0, readBytes); } } catch (IOException e) { throw new JMSException("Failed to stream file content to BytesMessage: " + e.getMessage()); } return bytesMessage; } else { return super.createJmsMessageForType(message, body, exchange, session); } } }
注册自定义Binding到Camel:
JmsComponent jmsComponent = new JmsComponent(); jmsComponent.setConnectionFactory(yourConnectionFactory); jmsComponent.setBinding(new CustomArtemisJmsBinding()); camelContext.addComponent("jms", jmsComponent);
内容的提问来源于stack exchange,提问作者ave4496

