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 :)

