ExecutorService's surprising performance break-even point --- rules of thumb?

Viewed 23020

I'm trying to figure out how to correctly use Java's Executors. I realize submitting tasks to an ExecutorService has its own overhead. However, I'm surprised to see it is as high as it is.

My program needs to process huge amount of data (stock market data) with as low latency as possible. Most of the calculations are fairly simple arithmetic operations.

I tried to test something very simple: "Math.random() * Math.random()"

The simplest test runs this computation in a simple loop. The second test does the same computation inside a anonymous Runnable (this is supposed to measure the cost of creating new objects). The third test passes the Runnable to an ExecutorService (this measures the cost of introducing executors).

I ran the tests on my dinky laptop (2 cpus, 1.5 gig ram):

(in milliseconds)
simpleCompuation:47
computationWithObjCreation:62
computationWithObjCreationAndExecutors:422

(about once out of four runs, the first two numbers end up being equal)

Notice that executors take far, far more time than executing on a single thread. The numbers were about the same for thread pool sizes between 1 and 8.

Question: Am I missing something obvious or are these results expected? These results tell me that any task I pass in to an executor must do some non-trivial computation. If I am processing millions of messages, and I need to perform very simple (and cheap) transformations on each message, I still may not be able to use executors...trying to spread computations across multiple CPUs might end up being costlier than just doing them in a single thread. The design decision becomes much more complex than I had originally thought. Any thoughts?


import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class ExecServicePerformance {

 private static int count = 100000;

 public static void main(String[] args) throws InterruptedException {

  //warmup
  simpleCompuation();
  computationWithObjCreation();
  computationWithObjCreationAndExecutors();

  long start = System.currentTimeMillis();
  simpleCompuation();
  long stop = System.currentTimeMillis();
  System.out.println("simpleCompuation:"+(stop-start));

  start = System.currentTimeMillis();
  computationWithObjCreation();
  stop = System.currentTimeMillis();
  System.out.println("computationWithObjCreation:"+(stop-start));

  start = System.currentTimeMillis();
  computationWithObjCreationAndExecutors();
  stop = System.currentTimeMillis();
  System.out.println("computationWithObjCreationAndExecutors:"+(stop-start));


 }

 private static void computationWithObjCreation() {
  for(int i=0;i<count;i++){
   new Runnable(){

    @Override
    public void run() {
     double x = Math.random()*Math.random();
    }

   }.run();
  }

 }

 private static void simpleCompuation() {
  for(int i=0;i<count;i++){
   double x = Math.random()*Math.random();
  }

 }

 private static void computationWithObjCreationAndExecutors()
   throws InterruptedException {

  ExecutorService es = Executors.newFixedThreadPool(1);
  for(int i=0;i<count;i++){
   es.submit(new Runnable() {
    @Override
    public void run() {
     double x = Math.random()*Math.random();     
    }
   });
  }
  es.shutdown();
  es.awaitTermination(10, TimeUnit.SECONDS);
 }
}
11 Answers

The Fixed ThreadPool's ultimate porpose is to reuse already created threads. So the performance gains are seen in the lack of the need to recreate a new thread every time a task is submitted. Hence the stop time must be taken inside the submitted task. Just with in the last statement of the run method.

In case it is useful to others, here are test results with a realistic scenario - use ExecutorService repeatedly until the end of all tasks - on a Samsung Android device.

 Simple computation (MS): 102
 Use threads (MS): 31049
 Use ExecutorService (MS): 257

Code:

   ExecutorService executorService = Executors.newFixedThreadPool(1);
        int count = 100000;

        //Simple computation
        Instant instant = Instant.now();
        for (int i = 0; i < count; i++) {
            double x = Math.random() * Math.random();
        }
        Duration duration = Duration.between(instant, Instant.now());
        Log.d("ExecutorPerformanceTest", "Simple computation (MS): " + duration.toMillis());


        //Use threads
        instant = Instant.now();
        for (int i = 0; i < count; i++) {
            new Thread(() -> {
                double x = Math.random() * Math.random();
            }
            ).start();
        }
        duration = Duration.between(instant, Instant.now());
        Log.d("ExecutorPerformanceTest", "Use threads (MS): " + duration.toMillis());


        //Use ExecutorService
        instant = Instant.now();
        for (int i = 0; i < count; i++) {
            executorService.execute(() -> {
                        double x = Math.random() * Math.random();
                    }
            );
        }
        duration = Duration.between(instant, Instant.now());
        Log.d("ExecutorPerformanceTest", "Use ExecutorService (MS): " + duration.toMillis());

I've faced a similar problem, but Math.random() was not the issue.
The problem is having many small tasks that take just a few milliseconds to complete. It is not much but a lot of small tasks in series ends up being a lot of time and I needed to parallelize.

So, the solution I found, and it might work for those of you facing this same problem: do not use any of the executor services. Instead create your own long living Threads and feed them tasks.

Here is an example, just as an idea don't try to copy paste it cause it probably won't work as I am using Kotlin and translating to Java in my head. The concept is what's important:

First, the Thread, a Thread that can execute a task and then continue there waiting for the next one:

public class Worker extends Thread {
  private Callable task;
  private Semaphore semaphore;
  private CountDownLatch latch;

  public Worker(Semaphore semaphore) {
    this.semaphore = semaphore;
  }

  public void run() {
    while (true) {
      semaphore.acquire(); // this will block, the while(true) won't go crazy
      if (task == null) continue;

      task.run();
      if (latch != null) latch.countDown();

      task = null;
    }
  }

  public void setTask(Callable task) {
    this.task = task;
  }

  public void setCountDownLatch(CountDownLatch latch) {
    this.latch = latch;
  }
}

There is two things here that need explanation:

  • the Semaphore: gives you control over how many tasks and when they are executed by this thread
  • the CountDownLatch: is the way to notify someone else that a task was completed

So this is how you would use this Worker, first just a simple example:

  Semaphore semaphore = new Semaphore(0); // initially the semaphore is closed
  Worker worker = new Worker(semaphore);
  worker.start();

  worker.setTask( .. your callable task .. );
  semaphore.release(); // this will allow one task to be processed by the worker

Now a more complicated example, with two Threads and waiting for both to complete using the CountDownLatch:

  Semaphore semaphore1 = new Semaphore(0);
  Worker worker1 = new Worker(semaphore1);
  worker1.start();

  Semaphore semaphore2 = new Semaphore(0);
  Worker worker2 = new Worker(semaphore2);
  worker2.start();

  // same countdown latch for both workers, with a counter of 2
  CountDownLatch countDownLatch = new CountDownLatch(2);
  worker1.setCountDownLatch(countDownLatch);
  worker2.setCountDownLatch(countDownLatch);

  worker1.setTask( .. your callable task .. );
  worker2.setTask( .. your callable task .. );
  semaphore1.release();
  semaphore2.release();

  countDownLatch.await(); // this will block until 2 tasks have been completed

And after that code runs you could just add more tasks to the same threads and reuse them. That's the whole point of this, reusing the threads instead of creating new ones.

It is unpolished as f*** but hopefully this gives you an idea. For me this was an improvement compared to no multi threading. And it was much much better than any executor service with any number of threads in the pool by far.

Related