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

如何通过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的jmxPriority
  • correlationID:关联ID,对应jmxCorrelationId
  • replyTo:回复地址,对应jmxReplyTo
  • messageID:消息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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 15:33:07