【发布时间】:2016-05-09 21:05:23
【问题描述】:
我想做的是通过 Apache Activemq 在 C# 应用程序和 Java 应用程序之间发送消息。
C#:
using (IConnection connection = factory.CreateConnection())
using (ISession session = connection.CreateSession())
{
IDestination destination = SessionUtil.GetDestination(session, "queue://ISI");
// Create a consumer and producer
using (IMessageProducer producer = session.CreateProducer(destination))
{
// Start the connection so that messages will be processed.
connection.Start();
ITextMessage request = session.CreateTextMessage(JsonConvert.SerializeObject(obj));
/*request.NMSCorrelationID = "abc";
request.Properties["NMSXGroupID"] = "cheese";
request.Properties["myHeader"] = "Cheddar";*/
producer.Send(request);
return request;
}
}
Java:
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(isiProperties.getMqUrl());
connection = connectionFactory.createConnection();
connection.start();
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Destination destination = session.createQueue("ISI");
MessageConsumer consumer = session.createConsumer(destination);
Message message = consumer.receive();
if(message instanceof TextMessage) {
try {
String text = ((TextMessage) message).getText();
ObjectMapper mapper = new ObjectMapper();
StatusChangeMessage obj = mapper.readValue(text, StatusChangeMessage.class);
if (obj instanceof StatusChangeMessage) {
StatusChangeMessage received = (StatusChangeMessage) obj;
Order order = orderRepository.findOne(received.getOrderId());
order.setStatus(received.getStatus());
orderRepository.saveAndFlush(order);
}
} catch(JMSException e) {
} catch(IOException e) {
}
}
C# 应用程序正确发送消息(它在 activemq 管理界面中可见)但没有活动订阅者(Java 应用程序应该这样做)。你看这里有什么不对吗?
基本上,if(message instanceof TextMessage) { 上的断点不会被执行。
【问题讨论】:
-
如果您使用的是调试器,那么
message的类型是什么?您可以在调试器中看到它,看看为什么它不是TextMessage的实例... -
调试器永远无法到达
if(message instanceof TextMessage)。它挂在Message message = consumer.receive(); -
您是否确认您的客户正在连接到同一个代理实例?
-
是的,当我杀死activemq时,Java应用程序会收到消息。