problems with priorities
Kulvir Bhogal <[email protected]> Thu, 6 May 2004 08:33:36 -0700 (PDT)
| Newsgroups | gmane.comp.java.openjms.devel |
|---|---|
| Message-ID | <[email protected]> |
I am having problems with the JMS priority feature of OpenJMS. I have created a simple topicpublisher that publishes a few messages of different JMSpriority. My subscriber is designed to pick up only messages of JMSPriority>5. For some reason, when I put in this SQL92 clause, my subscriber does not receive any messages. I have attached the files. Can someone please help me? Thanks, Kulvir
TopicPublish.java
(text/plain, 3 KB)
import java.util.Hashtable;
import javax.jms.JMSException;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.jms.Topic;
import javax.jms.TopicConnection;
import javax.jms.TopicConnectionFactory;
import javax.jms.TopicPublisher;
import javax.jms.TopicSession;
import javax.naming.Context;
import javax.naming.InitialContext;
import javax.naming.NamingException;
public class TopicPublish
{
public static void main(String[] args)
{
try
{
Hashtable properties = new Hashtable();
properties.put(
Context.INITIAL_CONTEXT_FACTORY,
"org.exolab.jms.jndi.InitialContextFactory");
properties.put(Context.PROVIDER_URL, "rmi://localhost:1099/");
Context context = new InitialContext(properties);
// retrieve topic connection factory
TopicConnectionFactory factory =
(TopicConnectionFactory)context.lookup(
"JmsTopicConnectionFactory");
// create a topic connection using factory
TopicConnection topicConnection =
factory.createTopicConnection();
topicConnection.start();
// create a topic session
// set transactions to false and set auto
// acknowledgement of receipt of messages
TopicSession topicSession =
topicConnection.createTopicSession(
false,Session.AUTO_ACKNOWLEDGE);
// lookup the topic, topic1
Topic topic = (Topic)
context.lookup("topic1");
// create a topic publisher and associate
// to the retrieved topic
TopicPublisher topicPublisher =
topicSession.createPublisher(topic);
// publish a message to the topic
System.out.println(
"Sample application publishing " +
"a message to a topic.");
TextMessage messageOne = topicSession.createTextMessage();
System.out.println("Message #1 - Priority 9 - Sport: Basketball");
messageOne.setText("Message #1 - Priority 9 - Sport: Basketball");
messageOne.setJMSPriority(9);
messageOne.setStringProperty("Sport","Basketball");
topicPublisher.publish(messageOne);
TextMessage messageTwo = topicSession.createTextMessage();
System.out.println("Message #2 - Priority 5 - Sport: Baseball");
messageTwo.setText("Message #2 - Priority 5 - Sport: Baseball");
messageTwo.setJMSPriority(5);
messageTwo.setStringProperty("Sport","Baseball");
topicPublisher.publish(messageTwo);
// clean up
topicPublisher.close();
topicSession.close();
topicConnection.close();
}
catch (NamingException e)
{
e.printStackTrace();
}
catch (JMSException e)
{
e.printStackTrace();
}
}
}
TopicSubscribeAsynchronous.java
(text/plain, 2.6 KB)
import java.util.Hashtable;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageListener;
import javax.jms.Queue;
import javax.jms.QueueConnectionFactory;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.jms.Topic;
import javax.jms.TopicConnection;
import javax.jms.TopicConnectionFactory;
import javax.jms.TopicSession;
import javax.jms.TopicSubscriber;
import javax.naming.Context;
import javax.naming.InitialContext;
import javax.naming.NamingException;
public class TopicSubscribeAsynchronous implements MessageListener
{
private TopicConnection topicConnection;
private TopicSession topicSession;
private Topic topic;
private TopicSubscriber topicSubscriber;
TopicSubscribeAsynchronous()
{
try
{
// create a JNDI context
Hashtable properties = new Hashtable();
properties.put(Context.INITIAL_CONTEXT_FACTORY,"org.exolab.jms.jndi.InitialContextFactory");
properties.put(Context.PROVIDER_URL,"rmi://localhost:1099/");
Context context = new InitialContext(properties);
// retrieve topic connection factory
TopicConnectionFactory topicConnectionFactory =
(TopicConnectionFactory)context.lookup("JmsTopicConnectionFactory");
// create a topic connection
topicConnection = topicConnectionFactory.createTopicConnection();
// create a topic session
// set transactions to false and set auto acknowledgement of receipt of messages
topicSession = topicConnection.createTopicSession(false,Session.AUTO_ACKNOWLEDGE);
// retrieve topic
topic = (Topic)context.lookup("topic1");
// create a topic subscriber and associate to the retrieved topic
String filterQuery = "JMSPriority > 5";
topicSubscriber = topicSession.createSubscriber(topic,filterQuery,false);
// associate message listener
topicSubscriber.setMessageListener(this);
// start delivery of incoming messages
topicConnection.start();
}
catch (NamingException e)
{
e.printStackTrace();
}
catch (JMSException e)
{
e.printStackTrace();
}
}
public static void main(String[] args)
{
System.out.println("Example to asynchronously listen to a topic.");
try
{
TopicSubscribeAsynchronous listener = new TopicSubscribeAsynchronous();
Thread.currentThread().sleep(2000);
}
catch (InterruptedException e)
{
e.printStackTrace();
}
}
// process incoming topic messages
public void onMessage(Message message)
{
try
{
String messageText = null;
if (message instanceof TextMessage)
messageText = ((TextMessage)message).getText();
System.out.println(messageText);
}
catch (JMSException e)
{
e.printStackTrace();
}
}
}