NOTE

6.17 fork_join

What fork/join is, why it exists, usage, and work stealing.

JavaCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is It?

A framework for parallel execution. It splits a large task into multiple small tasks, calculates the result of each small task, and finally aggregates the results of all small tasks to obtain the result of the large task.

1.1. Why Did It Appear?

When implementing fork/join with a thread pool directly, you need to consider having the current thread participate in the work instead of becoming only a supervisor.

public class CustomForkJoin
{
    private static class CountTask implements Callable<Integer>
    {
        private static final int THRESHOLD = 1000;
        private int start;
        private int end;
        private ExecutorService executorService;

        public CountTask(int start, int end, ExecutorService executorService)
        {
            this.start = start;
            this.end = end;
            this.executorService = executorService;
        }

        @Override
        public Integer call() throws Exception
        {
            System.out.println(Thread.currentThread().getName() + " working for start:" + start + " end:" + end);
            int sum = 0;
            // If the threshold has been reached, calculate directly without splitting further.
            boolean canCompute = (end - start) <= THRESHOLD;
            if (canCompute)
            {
                for (int i = start; i <= end; i++)
                {
                    sum += i;
                }
            }
            else
            {
                int middle = (start + end) / 2;
                CountTask leftTask = new CountTask(start, middle, executorService);
                CountTask rightTask = new CountTask(middle + 1, end, executorService);

                Future<Integer> leftResult = executorService.submit(leftTask);
                Future<Integer> rightResult = executorService.submit(rightTask);

                return leftResult.get() + rightResult.get();
            }
            return sum;
        }
    }

    public static void main(String[] args) throws Exception
    {
//        ExecutorService executorService = Executors.newFixedThreadPool(10);//Not enough threads; blocks.
        ExecutorService executorService = Executors.newCachedThreadPool();

        CountTask countTask = new CountTask(1, 1000000, executorService);
        Integer val = executorService.submit(countTask).get();
        executorService.shutdown();
        System.out.println(val);
    }
}

2. Usage Scenarios

CPU-intensive tasks.

3. How to Use It

public class ForkJoinTest
{
    private static class CountTask extends RecursiveTask<Integer>//A task with a return value; override compute.
    {
        private static final int THRESHOLD = 2;
        private int start;
        private int end;

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

        @Override
        protected Integer compute()
        {
            int sum = 0;
            // If the threshold has been reached, calculate directly without splitting further.
            boolean canCompute = (end - start) <= THRESHOLD;
            if (canCompute)
            {
                for (int i = start; i <= end; i++)
                {
                    sum += i;
                }
            }
            else
            {
                int middle = (start + end) / 2;
                CountTask leftTask = new CountTask(start, middle);
                CountTask rightTask = new CountTask(middle + 1, end);

                invokeAll(leftTask, rightTask);//Let the current thread participate instead of becoming only a supervisor.
                // Split into two tasks for execution.
                leftTask.fork();
                rightTask.fork();

                // Aggregate the results.
                int leftResult = leftTask.join();
                int rightResult = rightTask.join();

                sum = leftResult + rightResult;
            }
            return sum;
        }
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException
    {
        ForkJoinPool forkJoinPool = new ForkJoinPool();
        // Calculate 1+2+3+4.
        CountTask task = new CountTask(1,4);
        Future<Integer> result = forkJoinPool.submit(task);
        System.out.println(result.get());
    }
}

4. Implementation Analysis

4.1. Work Stealing

Each thread has its own work queue. If a thread has already finished all of its tasks, it can steal tasks from another thread’s work queue for execution.

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub