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
相关产品推荐
相关产品推荐

