如何通过JMX协议向ActiveMQ Artemis队列发送带指定Headers的消息?
问题描述
我需要通过JMX协议向ActiveMQ Artemis队列发送消息,目前不清楚如何正确传递消息Headers(如消息生命周期、消息类型等,区别于Properties)。现有代码似乎没有传递Headers的对应字段,同时我还需要发送字节类型的消息,设置消息生命周期为300秒。
现有发送消息的代码、JMXMessage、JMXHeaders类定义以及JMX库的sendMessage方法签名如下:
现有发送代码
public static void sendMessage(QueueControl queueControl, JMXMessage jmxMessage, String user, String password) throws Exception { boolean flag = false; queueControl.sendMessage(jmxMessage.getJmxProperties(), Integer.parseInt(jmxMessage.getJmxHeaders().getJmxType()),jmxMessage.getJmxBody(),flag,user, password); }
JMXMessage类
public class JMXMessage implements Serializable { private JMXHeaders jmxHeaders; private Map<String, String> jmxProperties; private String jmxBody; public JMXMessage(JMXHeaders jmxHeaders,String jmxBody) { this.jmxHeaders = jmxHeaders; this.jmxBody = jmxBody; } public JMXHeaders getJmxHeaders() { return jmxHeaders; } public void setJmxHeaders(JMXHeaders jmxHeaders) { this.jmxHeaders = jmxHeaders; } public Map<String, String> getJmxProperties() { return jmxProperties; } public void setJmxProperties(Map<String, String> jmxProperties) { this.jmxProperties = jmxProperties; } public String getJmxBody() { return jmxBody; } public void setJmxBody(String jmxBody) { this.jmxBody = jmxBody; } @Override public String toString() { return "Headers" + "\n" + "{" + "\n" + " MessageId=" + jmxHeaders.getJmxMessageId() + ";" + "\n" + " Priority=" + jmxHeaders.getJmxPriority() + ";" + "\n" + " Timestamp=" + jmxHeaders.getJmxTimestamp() + ";" + "\n" + " Expires=" + jmxHeaders.getJmxExpires() + ";" + "\n" + " Type=" + jmxHeaders.getJmxType() + ";" + "\n" + " DeliveryMode=" + jmxHeaders.getJmxDeliveryMode() + ";" + "\n" + " CorrelationId=" + jmxHeaders.getJmxCorrelationId() + ";" + "\n" + " ReplyTo=" + jmxHeaders.getJmxReplyTo() + ";" + "\n" + "}" + "\n" + "\n" + "Properties" + "\n" + "{" + "\n" + " Properties=" + jmxProperties + ";" + "\n" + "}" + "\n" + "\n" + "Body" + "\n" + "{" + "\n" + jmxBody + ";" + "\n" + "}" + "\n" + "\n"; } }
JMXHeaders类
public class JMXHeaders implements Serializable { private String jmxType; private String jmxExpires; private String jmxDeliveryMode; private String jmxReplyTo; private String jmxMessageId; private String jmxTimestamp; private String jmxCorrelationId; private String jmxPriority; public JMXHeaders(String jmxType, String jmxExpires, String jmxDeliveryMode, String jmxReplyTo, String jmxMessageId, String jmxTimestamp, String jmxCorrelationId, String jmxPriority) { this.jmxType = jmxType; this.jmxExpires = jmxExpires; this.jmxDeliveryMode = jmxDeliveryMode; this.jmxReplyTo = jmxReplyTo; this.jmxMessageId = jmxMessageId; this.jmxTimestamp = jmxTimestamp; this.jmxCorrelationId = jmxCorrelationId; this.jmxPriority = jmxPriority; } public JMXHeaders(String jmxType, String jmxExpires, String jmxDeliveryMode, String jmxReplyTo) { this.jmxType = jmxType; this.jmxExpires = jmxExpires; this.jmxDeliveryMode = jmxDeliveryMode; this.jmxReplyTo = jmxReplyTo; } public String getJmxExpires() { return jmxExpires; } public void setJmxExpires(String jmxExpires) { this.jmxExpires = jmxExpires; } public String getJmxMessageId() { return jmxMessageId; } public void setJmxMessageId(String jmxMessageId) { this.jmxMessageId = jmxMessageId; } public String getJmxTimestamp() { return jmxTimestamp; } public void setJmxTimestamp(String jmxTimestamp) { this.jmxTimestamp = jmxTimestamp; } public String getJmxCorrelationId() { return jmxCorrelationId; } public void setJmxCorrelationId(String jmxCorrelationID) { this.jmxCorrelationId = jmxCorrelationID; } public String getJmxPriority() { return jmxPriority; } public void setJmxPriority(String jmxPriority) { this.jmxPriority = jmxPriority; } public String getJmxType() { return jmxType; } public void setJmxType(String jmxType) { this.jmxType = jmxType; } public String getJmxDeliveryMode() { return jmxDeliveryMode; } public void setJmxDeliveryMode(String jmxDeliveryMode) { this.jmxDeliveryMode = jmxDeliveryMode; } public String getJmxReplyTo() { return jmxReplyTo; } public void setJmxReplyTo(String jmxReplyTo) { this.jmxReplyTo = jmxReplyTo; } @Override public String toString() { return "JMXHeaders{" + "jmxType='" + jmxType + '\'' + ", jmxExpires='" + jmxExpires + '\'' + ", jmxDeliveryMode='" + jmxDeliveryMode + '\'' + ", jmxReplyTo='" + jmxReplyTo + '\'' + ", jmxMessageId='" + jmxMessageId + '\'' + ", jmxTimestamp='" + jmxTimestamp + '\'' + ", jmxCorrelationId='" + jmxCorrelationId + '\'' + ", jmxPriority='" + jmxPriority + '\'' + '}'; } }
JMX库sendMessage方法签名
String sendMessage(@Parameter(name = "headers",desc = "The headers to add to the message") Map<String, String> var1, @Parameter(name = "type",desc = "A type for the message") int var2, @Parameter(name = "body",desc = "The body (byte[]) of the message encoded as a string using Base64") String var3, @Parameter(name = "durable",desc = "Whether the message is durable") boolean var4, @Parameter(name = "user",desc = "The user to authenticate with") String var5, @Parameter(name = "password",desc = "The users password to authenticate with") String var6) throws Exception;
解决方案
针对你的需求,需要从消息头构造、字节消息体处理、sendMessage调用修正三个方面调整代码:
1. 正确构造消息头Map
sendMessage的第一个参数headers是对应JMS标准消息头的Map,你需要将JMXHeaders中的字段映射到这个Map里,重点处理消息生命周期(expires):
expires:消息过期时间,需计算当前时间戳 + 300秒(300*1000毫秒),作为字符串传入priority:消息优先级,对应JMXHeaders的jmxPrioritycorrelationID:关联ID,对应jmxCorrelationIdreplyTo:回复地址,对应jmxReplyTomessageID:消息ID,对应jmxMessageId
2. 处理字节类型消息体
sendMessage的body参数要求是Base64编码的字符串,所以需要将字节数组转换为Base64格式。可以用Java自带的Base64.getEncoder().encodeToString(byte[])方法处理。
3. 修正sendMessage调用逻辑
现有代码错误地将jmxProperties传给了headers参数,需要替换为构造好的消息头Map;同时durable参数对应JMXHeaders中的jmxDeliveryMode(持久化模式对应true,非持久化对应false)。
修改后的sendMessage方法
import java.util.Base64; import java.util.HashMap; import java.util.Map; public static void sendMessage(QueueControl queueControl, JMXMessage jmxMessage, String user, String password) throws Exception { JMXHeaders headers = jmxMessage.getJmxHeaders(); // 构造消息头Map Map<String, String> messageHeaders = new HashMap<>(); // 设置消息生命周期:当前时间+300秒(毫秒) long expiresTime = System.currentTimeMillis() + 300 * 1000; messageHeaders.put("expires", String.valueOf(expiresTime)); // 其他消息头字段映射 if (headers.getJmxPriority() != null) { messageHeaders.put("priority", headers.getJmxPriority()); } if (headers.getJmxCorrelationId() != null) { messageHeaders.put("correlationID", headers.getJmxCorrelationId()); } if (headers.getJmxReplyTo() != null) { messageHeaders.put("replyTo", headers.getJmxReplyTo()); } if (headers.getJmxMessageId() != null) { messageHeaders.put("messageID", headers.getJmxMessageId()); } // 处理字节消息体:将字节数组转为Base64字符串 String base64Body; if (jmxMessage.getJmxBody() instanceof byte[]) { base64Body = Base64.getEncoder().encodeToString((byte[]) jmxMessage.getJmxBody()); } else { // 如果传入的是字符串,直接用(如果原本就是Base64编码的话) base64Body = jmxMessage.getJmxBody(); } // 处理durable参数:对应DeliveryMode,持久化=1对应true,非持久化=2对应false boolean durable = "1".equals(headers.getJmxDeliveryMode()); // 调用sendMessage,传入正确的headers、类型、body、durable等参数 queueControl.sendMessage( messageHeaders, Integer.parseInt(headers.getJmxType()), base64Body, durable, user, password ); }
构造JMXMessage示例(发送字节消息)
// 字节消息体 byte[] byteBody = "Hello Artemis".getBytes(); // 构造JMXHeaders,设置消息类型、投递模式等 JMXHeaders headers = new JMXHeaders( "2", // 消息类型:BytesMessage对应2,根据需求调整 null, // expires在sendMessage方法中计算,此处传null "1", // 投递模式:1=持久化,2=非持久化 null // 回复地址,不需要则传null ); // 将字节体转为Base64字符串,传入JMXMessage String base64Body = Base64.getEncoder().encodeToString(byteBody); JMXMessage jmxMessage = new JMXMessage(headers, base64Body); // 调用sendMessage发送 sendMessage(queueControl, jmxMessage, "username", "password");
内容的提问来源于stack exchange,提问作者user21859885
相关产品推荐
相关产品推荐

