【发布时间】:2020-06-12 11:31:40
【问题描述】:
我想使用 JMSXGroupId 属性对消息进行分组。但是消费者在没有分组的情况下检索队列的所有消息
我有一个本地 ActiveMQ 实例(版本 5.15.9、Windows 10、Java 8)。
我编写了 DoSend java 类来发送短信和 DoRead java 类来检索短信
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
public class DoSend {
public static void main(String[] args) throws JMSException {
ConnectionFactory factory = new ActiveMQConnectionFactory("tcp://localhost:61616");
Connection connection = null;
try {
connection = factory.createConnection("admin", "admin");
connection.start();
Session session = null;
try {
session = connection.createSession(true, Session.SESSION_TRANSACTED);
MessageProducer producer = null;
try {
Destination destination = session.createQueue("TEST");
producer = session.createProducer(destination);
int index = 0;
while (true) {
index += 1;
TextMessage message = session.createTextMessage(String.valueOf(index));
message.setStringProperty("JMSXGroupID", String.valueOf(index % 3));
producer.send(destination, message);
session.commit();
Thread.sleep(1000);
}
} finally {
if (producer != null) {
producer.close();
}
}
} catch (Throwable th) {
th.printStackTrace();
if (session != null) {
session.rollback();
}
throw new RuntimeException(th);
} finally {
if (session != null) {
session.close();
}
}
} finally {
if (connection != null) {
connection.close();
}
}
}
}
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
public class DoRead {
public static void main(String[] args) throws JMSException {
final ConnectionFactory factory = new ActiveMQConnectionFactory("tcp://localhost:61616");
System.out.println("START " + Thread.currentThread());
try {
Connection connection = null;
try {
connection = factory.createConnection("admin", "admin");
connection.start();
Session session = null;
try {
session = connection.createSession(true, Session.SESSION_TRANSACTED);
Destination destination = session.createQueue("TEST");
MessageConsumer consumer = null;
try {
consumer = session.createConsumer(destination);
TextMessage message;
while ((message = (TextMessage) consumer.receive(10000)) != null) {
System.out.println(Thread.currentThread() + ": message.index = " + message.getText());
System.out.println(Thread.currentThread() + ": message.JMSXGroupID = " + message.getStringProperty("JMSXGroupID"));
session.commit();
}
} finally {
if (consumer != null) {
consumer.close();
}
}
} finally {
if (session != null) {
session.close();
}
}
} finally {
if (connection != null) {
connection.close();
}
}
} catch (JMSException e) {
e.printStackTrace();
} finally {
System.out.println("FINISH " + Thread.currentThread());
}
}
}
我希望任何消费者都能检索具有相同 JMSXGroupId 的消息,但现在第一个消费者检索所有消息,而其他消费者什么也不检索。
【问题讨论】: