‏הצגת רשומות עם תוויות activemq. הצג את כל הרשומות
‏הצגת רשומות עם תוויות activemq. הצג את כל הרשומות

יום שישי, 12 ביולי 2013

Get notification from activemq using apache camel

In order to listen to notification from the active MQ I add the following route declaration to me app
 

  1:   public static RouteBuilder CreateRouteFromActiveMQ ()
  2:   {
  3:     return new RouteBuilder() {
  4:       public void configure() throws Exception {
  5:     from("test-jms:queue:test.queue")
  6:     .unmarshal().jaxb("org.apache.camel.example.console").  
  7:     setHeader(Exchange.FILE_NAME, constant("Received.Csv"))
  8:         .to("file:target/TestmyFiles").       
  9:       to("stream:out");
 10:     }
 11:     };
 12:   
 13:   }

Register the routes


  1: CamelContext context = new DefaultCamelContext();  
  2:     
  3:          ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
  4:             // Note we can explicit name the component
  5:          context.addComponent("test-jms", JmsComponent.jmsComponentAutoAcknowledge(connectionFactory));
  6:   
  7:       context.addRoutes(createRoute());
  8:       
  9:       context.addRoutes(CreateRouteFromActiveMQ());
 10:       
 11:       context.start();

יום חמישי, 27 ביוני 2013

Send message to ActiveMQ using camel

The following code sends message to ActiveMQ

  1: public class TestCodeForBlog {
  2:   public static RouteBuilder createRoute() {
  3:     return new RouteBuilder() {
  4:       public void configure() throws Exception {
  5:       from("direct:toActiveMQ")
  6:       .marshal().jaxb("org.apache.camel.example.console").
  7:       setHeader(Exchange.FILE_NAME, constant("test.ActiveMQ"))
  8:       .to("file:target/TestmyFiles")
  9:       .to("test-jms:queue:test.queue").
 10:       to("stream:out");
 11:     }
 12:     };
 13:     }
 14:   
 15:   public static void main(String[] args) {
 16:   
 17: 
 18:     try {
 19:       CamelContext context = new DefaultCamelContext();  
 20:     
 21:          ConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
 22:          // Explicit name the jms active mq component
 23:          context.addComponent("test-jms", JmsComponent.jmsComponentAutoAcknowledge(connectionFactory));
 24: 
 25:       context.addRoutes(createRoute());
 26:       context.start();
 27:     
 28:       PersoneData thePersoneData = new PersoneData();
 29:       
 30:       thePersoneData.Name = "Zvika";
 31:       
 32:       thePersoneData.Age = 42;
 33:       
 34:       thePersoneData.Email = "Loveme42@hotmail.com";
 35:       
 36:       thePersoneData.Address ="Bilo 9b Ness-Ziona";
 37:       
 38:       Thread.sleep(1000);  
 39:       
 40:       ProducerTemplate template = context.createProducerTemplate();
 41:       
 42:       template.sendBody("direct:toActiveMQ", thePersoneData);
 43:     
 44:       Thread.sleep(1000);
 45:       System.out.print("I'm out");
 46:        
 47:     } catch (Exception e) {
 48:       e.printStackTrace();
 49:     }
 50:   }
 51: }

The following dependencies should be inserted into the pom file :


  1: <dependency>
  2:       <groupId>org.apache.activemq</groupId>
  3:       <artifactId>activemq-core</artifactId>
  4:     </dependency>
  5:     <dependency>
  6:       <groupId>org.apache.activemq</groupId>
  7:       <artifactId>activemq-camel</artifactId>
  8:       <exclusions>
  9:         <exclusion>
 10:           <groupId>org.apache.camel</groupId>
 11:           <artifactId>camel-web</artifactId>
 12:         </exclusion>
 13:       </exclusions>
 14:     </dependency>
 15:     <dependency>
 16:       <groupId>org.apache.xbean</groupId>
 17:       <artifactId>xbean-spring</artifactId>
 18:       <exclusions>
 19:         <exclusion>
 20:           <groupId>org.springframework</groupId>
 21:           <artifactId>spring</artifactId>
 22:         </exclusion>
 23:       </exclusions>
 24:     </dependency>

יום שני, 10 ביוני 2013

ActiveMQ MapMessage

ActiveMQ propose a message of type MapMessage that its is a Wrapper of HashMap collection.

Example for Setting the Message Values:

  1: package com.mycompany.client;
  2: 
  3: import javax.jms.*;
  4: import javax.naming.InitialContext;
  5: import javax.naming.NamingException;
  6: import java.util.Properties;
  7: import java.util.logging.Level;
  8: import java.util.logging.Logger;
  9:  
 10: public class MessageProducer {
 11:  private String topicName = "myTopic.Programming";
 12:  
 13:  private String initialContextFactory = "org.apache.activemq"
 14: +".jndi.ActiveMQInitialContextFactory";
 15:  private String connectionString = "tcp://localhost:61616";
 16:   
 17: 
 18:  public void publishWithTopicLookup() {
 19:      try {
 20:          Properties properties = new Properties();
 21:          TopicConnection topicConnection = null;
 22:          properties.put("java.naming.factory.initial", initialContextFactory);
 23:          properties.put("connectionfactory.QueueConnectionFactory",
 24:            connectionString);
 25:          properties.put("topic." + topicName, topicName);
 26:         
 27:          
 28:           // initialize
 29:           // the required connection factories
 30:           InitialContext ctx = new InitialContext(properties);
 31:           TopicConnectionFactory topicConnectionFactory = (TopicConnectionFactory) ctx
 32:                            .lookup("QueueConnectionFactory");
 33:           topicConnection = topicConnectionFactory.createTopicConnection();
 34:          
 35:           TopicSession topicSession = topicConnection.createTopicSession(
 36:              false, Session.AUTO_ACKNOWLEDGE);
 37:            Topic topic = (Topic) ctx.lookup(topicName);
 38: 
 39:            javax.jms.TopicPublisher topicPublisher = topicSession
 40:                        .createPublisher(topic);
 41:         
 42:            MapMessage message = topicSession.createMapMessage();
 43:             
 44:            message.setString("Name", "Zvika");
 45:            message.setDouble("Age", 42.6);
 46:          
 47:            topicPublisher.publish(message);
 48:            System.out.println("Publishing message ");
 49:            topicPublisher.close();
 50:            topicSession.close(); 
 51:            topicConnection.close();
 52:      } catch (NamingException ex) {
 53:          Logger.getLogger(MessageProducer.class.getName()).log(Level.SEVERE, null, ex);
 54:      }
 55:          catch (JMSException e) {
 56:    throw new RuntimeException("Error in initial context lookup", e);
 57:   }
 58:  }
 59: }

Example for consuming a MapMessage:


  1: package com.mycompany.client;
  2: 
  3: import java.util.Enumeration;
  4: import java.util.Properties;
  5: import java.util.logging.Level;
  6: import java.util.logging.Logger;
  7: import javax.jms.JMSException;
  8: import javax.jms.MapMessage;
  9: import javax.jms.Message;
 10: import javax.jms.MessageListener;
 11: import javax.jms.Session;
 12: import javax.jms.Topic;
 13: import javax.jms.TopicConnection;
 14: import javax.jms.TopicConnectionFactory;
 15: import javax.jms.TopicSession;
 16: import javax.jms.TopicSubscriber;
 17: import javax.naming.InitialContext;
 18: import javax.naming.NamingException;
 19: 
 20: public class MessageConsumer {
 21:   
 22:     private class TestMessageListener implements MessageListener {
 23:  public void onMessage(Message message) {
 24:   try {
 25:         System.out.println("Got the Message  TimeStamp: "
 26:      +  message.getJMSTimestamp());
 27:         System.out.println("Got the Message JMS ID : "
 28:      +  message.getJMSMessageID() );
 29:         
 30:         MapMessage theMsg = (MapMessage)message;
 31:             
 32:         Enumeration theEnumeration = theMsg.getMapNames();
 33:         
 34:         while ( theEnumeration.hasMoreElements() )
 35:         {
 36:             String nextKey = theEnumeration.nextElement().toString();
 37:         
 38:             System.out.println("The key:" + nextKey + " the value:" + theMsg.getObjectProperty(nextKey).toString());
 39:         }
 40:   } catch (JMSException e) {
 41:    e.printStackTrace();
 42:   }
 43:  }
 44: } 
 45:     
 46:  private String topicName = "myTopic.Programming";
 47:  
 48:  private String initialContextFactory = "org.apache.activemq"
 49: +".jndi.ActiveMQInitialContextFactory";
 50:  private String connectionString = "tcp://localhost:61616";
 51:   
 52: 
 53:  public void ListenWithTopicLookup() {
 54:      try {
 55:          Properties properties = new Properties();
 56:          TopicConnection topicConnection = null;
 57:          properties.put("java.naming.factory.initial", initialContextFactory);
 58:          properties.put("connectionfactory.QueueConnectionFactory",
 59:            connectionString);
 60:          properties.put("topic." + topicName, topicName);
 61:         
 62:          
 63:           // initialize
 64:           // the required connection factories
 65:           InitialContext ctx = new InitialContext(properties);
 66:           TopicConnectionFactory topicConnectionFactory = (TopicConnectionFactory) ctx
 67:                            .lookup("QueueConnectionFactory");
 68:           topicConnection = topicConnectionFactory.createTopicConnection();
 69:          
 70:           TopicSession topicSession = topicConnection.createTopicSession(
 71:              false, Session.AUTO_ACKNOWLEDGE);
 72:            Topic topic = (Topic) ctx.lookup(topicName);
 73: 
 74:            TopicSubscriber theTopicSubscriber = topicSession.createSubscriber(topic);
 75:            
 76:            theTopicSubscriber.setMessageListener(new TestMessageListener());
 77:            
 78:            topicConnection.start();
 79:     
 80:      } catch (NamingException ex) {
 81:          Logger.getLogger(MessageProducer.class.getName()).log(Level.SEVERE, null, ex);
 82:      }
 83:          catch (JMSException e) {
 84:    throw new RuntimeException("Error in initial context lookup", e);
 85:   }
 86:  }
 87: }

יום ראשון, 2 ביוני 2013

ActiveMQ simple listener

I wrote a simple listener to the publisher from previous ActiveMQ post: activemq-topic-sender

The MessageListener callback class that’s receives notification when a topic message has received.

  1: private class TextMessageListener implements MessageListener {
  2:  public void onMessage(Message message) {
  3:   try {
  4:         System.out.println("Got the Message  TimeStamp: "
  5:             +  message.getJMSTimestamp());
  6:         System.out.println("Got the Message JMS ID : "
  7:             +  message.getJMSMessageID() );
  8:         
  9:         TextMessage theMsg = (TextMessage)message;
 10:         System.out.println("The message:" + theMsg.getText());
 11:   
 12:   } catch (JMSException e) {
 13:    e.printStackTrace();
 14:   }
 15:  }
 16: } 

The message consumer class :
Note:the connection must be started in order to receive notifications


  1: public class MessageConsumer {
  2: 
  3: private String topicName = "myTopic.Programming";
  4: private String initialContextFactory = "org.apache.activemq"
  5:   +".jndi.ActiveMQInitialContextFactory";
  6: private String connectionString = "tcp://localhost:61616";
  7:   
  8: public void ListenWithTopicLookup() {
  9:      try {
 10:          Properties properties = new Properties();
 11:          TopicConnection topicConnection = null;
 12:          properties.put("java.naming.factory.initial", initialContextFactory);
 13:          properties.put("connectionfactory.QueueConnectionFactory",
 14:            connectionString);
 15:          properties.put("topic." + topicName, topicName);
 16:         
 17:          
 18:           // initialize
 19:           // the required connection factories
 20:           InitialContext ctx = new InitialContext(properties);
 21:           TopicConnectionFactory topicConnectionFactory = (TopicConnectionFactory) ctx
 22:                            .lookup("QueueConnectionFactory");
 23:           topicConnection = topicConnectionFactory.createTopicConnection();
 24:          
 25:           TopicSession topicSession = topicConnection.createTopicSession(
 26:              false, Session.AUTO_ACKNOWLEDGE);
 27:            Topic topic = (Topic) ctx.lookup(topicName);
 28: 
 29:            TopicSubscriber theTopicSubscriber = topicSession.createSubscriber(topic);
 30:            
 31:            theTopicSubscriber.setMessageListener(new TextMessageListener());
 32:            
 33:            topicConnection.start();
 34:     
 35:      } catch (NamingException ex) {
 36:          Logger.getLogger(MessageProducer.class.getName()).log(Level.SEVERE, null, ex);
 37:      }
 38:          catch (JMSException e) {
 39:    throw new RuntimeException("Error in initial context lookup", e);
 40:   }
 41:  }
 42: }

I change the main class in order to start 2 new message consumers :


  1:  public static void main( String[] args )
  2:     {
  3:         System.out.println( "Hello World!" );
  4:         
  5:         MessageConsumer theMessageConsumer= new     MessageConsumer();
  6:         
  7:         theMessageConsumer.ListenWithTopicLookup();
  8:    
  9:         MessageConsumer OthertheMessageConsumer= new     MessageConsumer();
 10:         
 11:         OthertheMessageConsumer.ListenWithTopicLookup();
 12:         
 13:         MessageProducer publisher = new MessageProducer();
 14:   
 15:         publisher.publishWithTopicLookup();
 16:         try {
 17:             Thread.sleep( 1000);
 18:         } catch (InterruptedException ex) {
 19:             Logger.getLogger(App.class.getName()).log(Level.SEVERE, null, ex);
 20:         }

The Result :
[exec:exec]
Hello World!
log4j:WARN No appenders could be found for logger (org.apache.activemq.thread.TaskRunnerFactory).
log4j:WARN Please initialize the log4j system properly.
log4j:WARN See
http://logging.apache.org/log4j/1.2/faq.html#noconfig for more info.
Publishing message ActiveMQTextMessage {commandId = 0, responseRequired = false, messageId = ID:Zvika-PC-2639-1370199196674-5:1:1:1:1, originalDestination = null, originalTransactionId = null, producerId = null, destination = topic://myTopic.Programming, transactionId = null, expiration = 0, timestamp = 1370199196906, arrival = 0, brokerInTime = 0, brokerOutTime = 0, correlationId = null, replyTo = null, persistent = true, type = null, priority = 4, groupID = null, groupSequence = 0, targetConsumerId = null, compressed = false, userID = null, content = null, marshalledProperties = null, dataStructure = null, redeliveryCounter = 0, size = 0, properties = null, readOnlyProperties = false, readOnlyBody = false, droppable = false, text = This is a test message}
Got the Message  TimeStamp: 1370199196906
Got the Message JMS ID : ID:Zvika-PC-2639-1370199196674-5:1:1:1:1
The message:This is a test message
Got the Message  TimeStamp: 1370199196906
Got the Message JMS ID : ID:Zvika-PC-2639-1370199196674-5:1:1:1:1
The message:This is a test message

Each listener got notification about the message that was send .

יום שני, 27 במאי 2013

ActiveMQ topic sender

Starting ActiveMQ service
Capture19
Use the following maven pom for my simple client test

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>
  <groupId>com.mycompany</groupId>
  <artifactId>Client</artifactId>
  <version>1.0-SNAPSHOT</version>
  <packaging>jar</packaging>
  <name>Client</name>
  <url>http://maven.apache.org</url>
  <properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
  </properties>
<repositories> 
    <repository>
      <id>repository.jboss.org-public</id>
      <name>JBoss.org Maven repository</name>
      <url>https://repository.jboss.org/nexus/content/groups/public</url>
    </repository>  
</repositories>
  <dependencies>
    <dependency>
      <groupId>junit</groupId>
      <artifactId>junit</artifactId>
      <version>3.8.1</version>
      <scope>test</scope>
    </dependency>
   <dependency>
	<groupId>javax.jms</groupId>
	<artifactId>jms</artifactId>
	<version>1.1</version>
</dependency>
<dependency>
  <groupId>org.apache.activemq</groupId>
  <artifactId>activemq-all</artifactId>
  <version>5.8.0</version>
</dependency>
  </dependencies>
</project>

The ActiveMQ monitor can be run  from the following url http://localhost:8161


The code for the demo

  1: package com.mycompany.client;
  2: 
  3: import javax.jms.*;
  4: import javax.naming.InitialContext;
  5: import javax.naming.NamingException;
  6: import java.util.Properties;
  7: import java.util.logging.Level;
  8: import java.util.logging.Logger;
  9:  
 10: public class MessageProducer {
 11:  private String topicName = "myTopic.Programming";
 12:  
 13:  private String initialContextFactory = "org.apache.activemq"
 14: +".jndi.ActiveMQInitialContextFactory";
 15:  private String connectionString = "tcp://localhost:61616";
 16:   
 17: 
 18:  public void publishWithTopicLookup() {
 19:      try {
 20:          Properties properties = new Properties();
 21:          TopicConnection topicConnection = null;
 22:          properties.put("java.naming.factory.initial", initialContextFactory);
 23:          properties.put("connectionfactory.QueueConnectionFactory",
 24:            connectionString);
 25:          properties.put("topic." + topicName, topicName);
 26:         
 27:          
 28:           // initialize
 29:           // the required connection factories
 30:           InitialContext ctx = new InitialContext(properties);
 31:           TopicConnectionFactory topicConnectionFactory = (TopicConnectionFactory) ctx
 32:                            .lookup("QueueConnectionFactory");
 33:           topicConnection = topicConnectionFactory.createTopicConnection();
 34:          
 35:           TopicSession topicSession = topicConnection.createTopicSession(
 36:              false, Session.AUTO_ACKNOWLEDGE);
 37:            Topic topic = (Topic) ctx.lookup(topicName);
 38: 
 39:            javax.jms.TopicPublisher topicPublisher = topicSession
 40:                        .createPublisher(topic);
 41:         
 42:            String msg = "This is a test message";
 43:            TextMessage textMessage = topicSession.createTextMessage(msg);
 44:         
 45:            topicPublisher.publish(textMessage);
 46:            System.out.println("Publishing message " +textMessage);
 47:            topicPublisher.close();
 48:            topicSession.close(); 
 49:            topicConnection.close();
 50:      } catch (NamingException ex) {
 51:          Logger.getLogger(MessageProducer.class.getName()).log(Level.SEVERE, null, ex);
 52:      }
 53:          catch (JMSException e) {
 54:    throw new RuntimeException("Error in initial context lookup", e);
 55:   }
 56:  }
 57: }

The connections:
http://localhost:8161/admin/connections.jsp


Note I set a break point before the connection closing in order to inspect the connection in the connections page


Capture21


Capture20


The sample use JNDI in order to generate the ActiveMQ objects .
The context factory that  is used to generate the objects is  ActiveMQInitialContextFactory.
It is set in the following command:
properties.put("java.naming.factory.initial", initialContextFactory);


The created topic in the topics page (http://localhost:8161/admin/topics.jsp)Capture22