Sunday, February 8, 2015

Distributed Computing Using Hazelcast

I was reading about in memory data grid solutions that are available in market when I came across Hazelcast. One of the feature that most of the in-memory data grid solutions offer is distributed computing.It means executing your code in multiple machines instead of just one to maximize throughput.

I found this very interesting.So I have done some experiments with it using Hazelcast.What i will do is to make the computer sum from 1 to 1000.But not on one node, but on two or three nodes and see how distributed system computing performs.

For the time being, I will do this on one machine(on my PC) with multiple JVMs ( multiple nodes) and in some other post I will try to move this experiment to AWS.

Here is the configuration of my machine:

version:Intel(R) Core(TM) i7-4600U CPU @ 2.10GHz
cores = 2
enabledcores=2
threads=4
memory=8GB

So there are 2 cores with 8 gb of memory.

I will create a task called SumCalculatorTask to which i give "start" and "end" numbers to calculate the sum.Also I have added some delay to get some decent numbers...

public class SumCalculatorTask implements Callable<Integer>, HazelcastInstanceAware, Serializable {

    private transient HazelcastInstance hazelcastInstance;

    private int start;
    private int end;

    public SumCalculatorTask(int start, int end) {
        this.start = start;
        this.end = end;
    }

    @Override    public Integer call() throws Exception {
        System.out.println("Calculating sum from " + start + " till " + end);
        Integer sum = 0;
        long startTime = System.currentTimeMillis();
        for (int i = start; i <= end; i++) {
            sum += i;
            Thread.sleep(10);
        }
        long endTime = System.currentTimeMillis();
        System.out.println("Total time taken to calculate sum from " + start + " till " + end + " is " + (endTime - startTime) + " millis");
        return sum;
    }

    @Override    public void setHazelcastInstance(HazelcastInstance hazelcastInstance) {
        this.hazelcastInstance = hazelcastInstance;
    }
}

Hazelcast has its own ExecutorService which executes task on members of cluster.So you need to implement HazelcastInstanceAware so that it can be submitted to IExecutorService of Hazelcast.Also Serializable needs to be implements,because this task would go through wire to different machines.
Next step is to create a Hazelcast node.

 public class Node {
    public static void main(String[] args) throws FileNotFoundException {
        Config config = new ClasspathXmlConfig("hazelcast-cluster-config.xml");
        config.setInstanceName("hazelgrid");
        HazelcastInstance instance = Hazelcast.getOrCreateHazelcastInstance(config);
    }
}

If you run this class one time it will create one node,if you run it two times it will create two nodes..and so on..After running the node you need to submit this task to node(s) from a client.This is done by this code...

public class Client {
    public static void main(String[] args) throws IOException, ExecutionException, InterruptedException {
        InputStream inputStream = Client.class.getClassLoader().getResourceAsStream("hazelcast-client-config.xml");
        ClientConfig clientConfig = new XmlClientConfigBuilder(inputStream).build();
        HazelcastInstance client = HazelcastClient.newHazelcastClient(clientConfig);

        calculateSum(client);
    }

    private static void calculateSum(HazelcastInstance client) throws InterruptedException, ExecutionException {
        IExecutorService executorService = client.getExecutorService("default");
        Integer finalSum = 0;
        List<Future<Integer>> futures = new ArrayList<>();
        long startTime = System.currentTimeMillis();
        for (int i = 0; i < 10; i++) {
            futures.add(executorService.submit(new SumCalculatorTask(100 * i + 1, 100 * i + 100)));
        }

        for (Future<Integer> future : futures) {
            finalSum += future.get();
        }
        long endTime = System.currentTimeMillis();
        System.out.println("Final Sum: " + finalSum + " calculated in " + (endTime-startTime) + " millis.");
    }

}

Now, I ran different scenarios on my machine and collected the data which shows the distributed computing is amazing.If you can divide your process into independent tasks, then you can use the computational power of multiple machines and collect the result.

 Further you can scale - up and scale-out to speed up the process.For scale-up you can add more CPU,RAM to your machine and to scale-out you can add more such machines.Also Hazelcast offers "thread-pool-size" which signfies how many threads would be run in a machine.(Will be useful when you scale-up).

Here is the data which I collected on my machine.Can vary on yours...

Cluster Executorservice configured for 1 thread per member in threadpool
Final Sum: 500500 calculated in 10191 millis. --> 1 node cluster
Final Sum: 500500 calculated in 6128 millis.  -->  2 node cluster
Final Sum: 500500 calculated in 4087 millis. -->   3 node cluster
 
Cluster Executorservice configured for 2 thread per member in threadpool
Final Sum: 500500 calculated in 5118 millis. --> 1 node cluster
Final Sum: 500500 calculated in 3060 millis.  -->  2 node cluster
Final Sum: 500500 calculated in 2062 millis. -->   3 node cluster

 
Cluster Executorservice configured for 3 thread per member in threadpool
Final Sum: 500500 calculated in 4087 millis. --> 1 node cluster
Final Sum: 500500 calculated in 3074 millis.  -->  2 node cluster
Final Sum: 500500 calculated in 2038 millis. -->   3 node cluster


So as you can see, once we start adding nodes and threads the processing time starts reducing which is just amazing.Just imagine when you have 100 of nodes with each machine with 8 -16 cores how fast can you compute.This concept is really scalable for most of the problems.

 
I have shared this code on
https://github.com/dheerajsharma1990/hazelgrid
You can look into it..Thanks :)


Thursday, January 15, 2015

ActiveMQ Broker and Spring JMS Integration

I was writing a test case in which I was sending an event through my test case and was allowing the application to listen to that event and then further do the end to end testing.

To start the ActiveMQ in my test cases i was using BrokerService with the url on which my application is listening.After the broker service got started i was sending event immediately.

The problem was that my application was not able to listen to that event.

Sequence was..

1. start spring context of my application ( which use service-activator of spring to listen to event).

2. start broker

3.publish event through my test case.

At that time,I didn't have time to see what's wrong so I put Thread.sleep(10*1000) between step 2 & 3 and my application started to listen to event.

Now today I spent some time finding out the issue.Initially I thought that I am sending message before the broker is even started so I created a main class , started the broker and send the message.But the broker was started before i published the message.

Then what went wrong?

I digged more and placed Thread.sleep above "start broker" and below "start broker" to find out what's really happening.

After going through the logs I found that the spring context listener is somewhat synchronizing with broker that i started which is taking 4-5 seconds..

11:54:07.029 [main] INFO  o.a.activemq.broker.BrokerService - For help or more information please see: http://activemq.apache.org
11:54:07.029 [main] INFO  c.s.c.p.activemq.ActiveMqRunner - Time taken to start broker: 903 millis
11:54:07.029 [main] INFO  c.s.c.p.activemq.ActiveMqRunner - Broker started !!
Jan 16, 2015 11:54:07 AM org.glassfish.grizzly.http.server.NetworkListener start
INFO: Started listener bound to [0.0.0.0:8092]
Jan 16, 2015 11:54:07 AM org.glassfish.grizzly.http.server.HttpServer start
INFO: [HttpServer] Started.
11:54:11.138 [org.springframework.jms.listener.DefaultMessageListenerContainer#0-1] DEBUG o.a.a.transport.WireFormatNegotiator - Sending: WireFormatInfo { version=9, properties={MaxFrameSize=9223372036854775807, CacheSize=1024, CacheEnabled=true, SizePrefixDisabled=false, MaxInactivityDurationInitalDelay=10000, TcpNoDelayEnabled=true, MaxInactivityDuration=0, TightEncodingEnabled=true, StackTraceEnabled=true}, magic=[A,c,t,i,v,e,M,Q]}

 This was the main problem.So i decided to start broker at first step so that all the application context can get in sync with my broker.

Final sequence was:

1. start broker

2. start spring context of my application ( which use service-activator of spring to listen to event).

3. publish event

And my problem got solved ..

Monday, January 5, 2015

SVN or GIT (Know your use case !!)

I am working on project that uses GIT as source code manager.Prior to this I worked on a project that was using SVN. Now everytime I make a commit on to GIT I just wonder why we are using GIT or what is wrong with SVN.In search of answer, I found some really good posts on google.

Ultimately, I came to the conclusion that its your use case that matters.They both are good. Now in which use case you should use what is a question. Here are some of my findings..


  • GIT is decentralized. Which essentially means you can commit to your local machine even when you are not connected to the internet.
Now the question comes, does your projects or team would be frequently in a situation where you dont have internet access!! Is your company is offering you special kind of holidays where you go to the mountains and you have to commit. How often are you going to do that?Just because once in a bluemoon you dont have internet access and you are desperate to commit ( and not wishing to keep your changes in changelist) you are using GIT.Even though you commit, no body would be able to see your changes because you are not connected to internet!! 
Using GIT..
  • GIT is the new cool thing that is trending in the developer fashion space.If your use case is to use version control that is trending..use GIT
We like to follow the latest trend..
  • GIT is well suited when multiple developers are working on different branches and not connected to main branch.(trunk/origin whatever..)
Now, does your project in your company work in this way?If yes then go for git, because git have sophisticated techniques of merging,rebasing,switching between branches. My project however work in a way most projects work.All developers make changes for the upcoming release.We have builds running to ensure nothing gets broken.Then why use GIT??

  • GIT is best suited for open source projects where multiple developers fork a branch and everyone adds a functionality and the owner chooses the best one from several branches and pulls the changes to the main..
In a nutshell, I would like to say if you and your team really needs some kind of offline version control , and you guys can take out time to learn this new system with entirely different terminology ( checkout /clone and commit/push) then go for GIT.And if you are interested in spending more time on delivering performant system and just a version control system that manages source code for you then go for SVN because its very intuitive and easy to use..


Sunday, October 19, 2014

InCompatibleClassChangeError

I am working on a new project that has a really messed up classpath jars.I was trying to run a class in local when i came accross this error InCompatibleClassChangeError.

The logs looked somewhat like this..

Exception in thread "main" org.springframework.beans.factory.BeanCreationException: 
Error creating bean with name 'org.springframework.dao.annotation.PersistenceExceptionTranslationPostProcessor#0' 
defined in class path resource [commonApplicationContextJPA.xml]: 
Initialization of bean failed; nested exception is org.springframework.beans.factory.BeanCreationException: Error creating bean with name 'entityManagerFactory' defined in class path resource [commonApplicationContextJPA.xml]: Invocation of init method failed; nested exception is java.lang.IncompatibleClassChangeError: Implementing class


at java.lang.ClassLoader.loadClass(ClassLoader.java:358) 
at org.hibernate.ejb.Ejb3Configuration.<clinit>(Ejb3Configuration.java:127) 
at org.hibernate.ejb.HibernatePersistence.createContainerEntityManagerFactory(HibernatePersistence.java:71) 
at org.springframework.orm.jpa.LocalContainerEntityManagerFactoryBean.createNativeEntityManagerFactory(LocalContainerEntityManagerFactoryBean.java:268) 
at org.springframework.orm.jpa.AbstractEntityManagerFactoryBean.afterPropertiesSet(AbstractEntityManagerFactoryBean.java:310) 
at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.invokeInitMethods(AbstractAutowireCapableBeanFactory.java:1514) 
at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.initializeBean(AbstractAutowireCapableBeanFactory.java:1452)



Baffled by it, I started to look what exactly is going on.So I searched google and found what exactly InCompatibleClassChangeError is and I found this..



This is a complete list of changes in Java library API that may cause clients built with an old version of the library to throw java.lang.IncompatibleClassChangeError if they run on a new one (i.e. breaking BC):
  1. Non-final field become static,
  2. Non-constant field become non-static,
  3. Class become interface,
  4. Interface become class,
  5. if you add a new field to class/interface (or add new super-class/super-interface) then a static field from a super-interface of a client class C may hide an added field (with the same name) inherited from the super-class of C (very rare case).
By this time, I was sure that definitely there is some mismatch between the version of some jar.I started looking for mismatch in the spring versions, tried everything i could to make the spring version of entire project to one version but no luck.
Then i googled more and found that, such error comes for hibernate when tow different hibernate versions are present.Finally, i found i built maven dependency tree using "mvn dependency:tree" and searched for two different versions of hibernate.And voila,there were two of them.
Finally,i excluded the older version using this and the exception went away.I am happy now.I have learnt to approach problems in a more practical way.


<Exclusion> 
    <GroupId> org.hibernate </ groupId> 
    <ArtifactId> Hibernate </ artifactId> 
</ Exclusion> 
<Exclusion> 
    <GroupId> org.hibernate </ groupId> 
    <ArtifactId> hibernate-commons-annotations </ artifactId>
</ Exclusion>