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

HornetQ嵌入模式下如何发送并转换POJO对象

HornetQ嵌入模式下发送POJO并支持属性过滤的实现方案

嘿,你的思路完全没问题——把POJO转成Map后将属性存入ClientMessage的properties中,既能实现消息发送,又能让消费者通过消息选择器按属性过滤,结合Jackson的ObjectMapper来做序列化/反序列化是非常高效的方案。下面给你梳理完整的实现步骤和代码示例:

1. 准备你的POJO类

首先确保你的POJO实现Serializable接口,同时完善构造方法、getter/setter(Jackson序列化需要无参构造):

import java.io.Serializable;

public class Pojo implements Serializable {
    private Integer id;
    private String name;
    private String phone;

    // 无参构造(Jackson序列化必备)
    public Pojo() {}

    // 全参构造
    public Pojo(Integer id, String name, String phone) {
        this.id = id;
        this.name = name;
        this.phone = phone;
    }

    // getter和setter方法
    public Integer getId() { return id; }
    public void setId(Integer id) { this.id = id; }
    public String getName() { return name; }
    public void setName(String name) { this.name = name; }
    public String getPhone() { return phone; }
    public void setPhone(String phone) { this.phone = phone; }
}

2. 发布端实现(转Map+设置消息属性)

这里我们先初始化HornetQ嵌入式服务器,将POJO转成Map后把每个属性存入消息的properties(支持过滤),同时把POJO转成JSON存入消息体(方便消费者直接还原对象):

import org.hornetq.api.core.TransportConfiguration;
import org.hornetq.api.core.client.*;
import org.hornetq.core.remoting.impl.invm.InVMConnectorFactory;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.util.Map;

public class PojoPublisher {
    public static void main(String[] args) throws Exception {
        // 1. 初始化嵌入式HornetQ连接
        ServerLocator serverLocator = HornetQClient.createServerLocatorWithoutHA(new TransportConfiguration(InVMConnectorFactory.class.getName()));
        ClientSessionFactory sessionFactory = serverLocator.createSessionFactory();
        ClientSession session = sessionFactory.createSession(false, true, true);

        // 2. 创建目标队列(不存在则创建)
        String queueName = "PojoQueue";
        if (!session.queueQuery(queueName).isExists()) {
            session.createQueue(queueName, queueName, true);
        }

        // 3. 创建消息生产者
        ClientProducer producer = session.createProducer(queueName);

        // 4. 准备要发送的POJO实例
        Pojo pojo = new Pojo(1, "张三", "13800138000");

        // 5. 用Jackson将POJO转成Map
        ObjectMapper mapper = new ObjectMapper();
        Map<String, Object> pojoMap = mapper.convertValue(pojo, Map.class);

        // 6. 创建ClientMessage并设置可过滤的属性
        ClientMessage message = session.createMessage(true);
        // 遍历Map,将属性存入消息properties(注意匹配HornetQ支持的属性类型)
        for (Map.Entry<String, Object> entry : pojoMap.entrySet()) {
            String key = entry.getKey();
            Object value = entry.getValue();
            if (value instanceof Integer) {
                message.putIntProperty(key, (Integer) value);
            } else if (value instanceof String) {
                message.putStringProperty(key, (String) value);
            }
            // 可根据POJO属性类型扩展Long、Boolean等其他类型
        }

        // 7. 将POJO转成JSON存入消息体(可选,方便消费者直接还原对象)
        String pojoJson = mapper.writeValueAsString(pojo);
        message.getBodyBuffer().writeString(pojoJson);

        // 8. 发送消息
        producer.send(message);
        System.out.println("POJO消息发送成功:" + pojoJson);

        // 9. 关闭资源
        producer.close();
        session.close();
        sessionFactory.close();
        serverLocator.close();
    }
}

3. 消费者实现(按属性过滤+还原POJO)

消费者这边可以通过消息选择器指定过滤条件(比如只接收name为"张三"的消息),既可以从properties读取属性,也可以从消息体直接还原POJO:

import org.hornetq.api.core.TransportConfiguration;
import org.hornetq.api.core.client.*;
import org.hornetq.core.remoting.impl.invm.InVMConnectorFactory;
import com.fasterxml.jackson.databind.ObjectMapper;

public class PojoConsumer {
    public static void main(String[] args) throws Exception {
        // 1. 初始化客户端连接
        ServerLocator serverLocator = HornetQClient.createServerLocatorWithoutHA(new TransportConfiguration(InVMConnectorFactory.class.getName()));
        ClientSessionFactory sessionFactory = serverLocator.createSessionFactory();
        ClientSession session = sessionFactory.createSession(false, true, true);

        // 2. 创建消费者,设置消息选择器(只接收name为"张三"的消息)
        String queueName = "PojoQueue";
        String selector = "name = '张三'";
        ClientConsumer consumer = session.createConsumer(queueName, selector);

        // 3. 启动会话开始接收消息
        session.start();

        // 4. 接收消息(5秒超时)
        ClientMessage message = consumer.receive(5000);
        if (message != null) {
            // 方式1:从消息properties读取属性
            Integer id = message.getIntProperty("id");
            String name = message.getStringProperty("name");
            String phone = message.getStringProperty("phone");
            System.out.println("从消息属性读取内容:id=" + id + ", name=" + name + ", phone=" + phone);

            // 方式2:从消息体还原POJO对象
            String pojoJson = message.getBodyBuffer().readString();
            ObjectMapper mapper = new ObjectMapper();
            Pojo pojo = mapper.readValue(pojoJson, Pojo.class);
            System.out.println("还原后的POJO对象:" + pojo.getName() + " - " + pojo.getPhone());

            // 确认消息已处理
            message.acknowledge();
        } else {
            System.out.println("超时未收到符合条件的消息");
        }

        // 5. 关闭资源
        consumer.close();
        session.close();
        sessionFactory.close();
        serverLocator.close();
    }
}

几个关键注意点

  • 属性类型匹配:HornetQ的消息properties只支持特定类型(Int、String、Long等),转Map时要注意属性类型对应,否则会导致设置失败或过滤不生效。
  • 消息选择器语法:选择器语法类似SQL,比如id > 5、phone LIKE '138%',字符串需用单引号包裹。
  • Jackson依赖:确保项目引入jackson-databind依赖,否则无法完成POJO与Map/JSON的转换。
  • 嵌入式服务器配置:如果需要持久化消息,可以在初始化服务器时添加存储目录等配置参数。

内容的提问来源于stack exchange,提问作者gekoramy bean

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:31:27