【问题标题】:Execute the method in parallel within the loop using Java使用 Java 在循环内并行执行方法
【发布时间】:2018-05-07 06:21:05
【问题描述】:

我有如下代码。在一个循环中,它正在执行方法“process”。它按顺序运行。我想并行运行这个方法,但它应该在循环内完成,以便我可以在下一行求和。即即使它是并行运行的,所有函数都应该在第二个 for 循环执行之前完成。

在Jdk1.7而不是JDK1.8版本如何解决这个问题?

public static void main(String s[]){
    int arrlen = 10;
    int arr[] = new int[arrlen] ;

    int t =0;
    for(int i=0;i<arrlen;i++){
        arr[i] = i;
        t = process(arr[i]);
        arr[i] = t;
    }

    int sum =0;
    for(int i=0;i<arrlen;i++){
        sum += arr[i];
    }
    System.out.println(sum);

}

public static int process(int arr){
    return arr*2;
}

【问题讨论】:

标签: java multithreading java-threads


【解决方案1】:

以下示例可能会对您有所帮助。我已经使用 fork/join 框架来做到这一点。

对于像您的示例这样的小数组,传统方法可能更快,我怀疑 fork/join 方式会花费稍长的时间。但对于较大的规模或进程,fork/join 框架是合适的。甚至 java 8 并行流也使用 fork/join 框架作为底层基础。

public class ForkMultiplier extends RecursiveAction {
        int[] array;
        int threshold = 3;
        int start;
        int end;

        public ForkMultiplier(int[] array,int start, int end) {
            this.array = array;
            this.start = start;
            this.end = end;
        }

        protected void compute() {
            if (end - start < threshold) {
                computeDirectly();
            } else {
                int middle = (end + start) / 2;
                ForkMultiplier f1= new ForkMultiplier(array, start, middle);
                ForkMultiplier f2= new ForkMultiplier(array, middle, end);
                invokeAll(f1, f2);
            }
        }

        protected void computeDirectly() {
            for (int i = start; i < end; i++) {
                array[i] = array[i] * 2;
            }
        }
    }

你的主班会喜欢下面这个

 public static void main(String s[]){

        int arrlen = 10;
        int arr[] = new int[arrlen] ;


        for(int i=0;i<arrlen;i++){
            arr[i] = i;
        }

        ForkJoinPool pool = new ForkJoinPool();
        pool.invoke(new ForkMultiplier(arr, 0, arr.length));

        int sum =0;
        for(int i=0;i<arrlen;i++){
            sum += arr[i];
        }

        System.out.println(sum);

    }

【讨论】:

    【解决方案2】:

    您基本上需要使用自 Java 1.5 以来存在的 Executors 和 Futures 组合(请参阅Java Documentation)。

    在以下示例中,我创建了一个主类,它使用另一个帮助类,该类的作用类似于您要并行化的处理器。

    主类分为3个步骤:

    1. 创建进程池并并行执行任务。
    2. 等待所有任务完成。
    3. 从任务中收集结果。

    出于教学原因,我放了一些日志,更重要的是,我在每个流程的业务逻辑中放了一个随机等待时间,模拟由 Process 类运行的耗时算法.

    每个进程的最大等待时间为2秒,这也是第2步的最高等待时间,即使你增加并行任务的数量(只需尝试更改以下代码的变量totalTasks来测试一下)。

    这里是主类:

    package com.example;
    
    import java.util.ArrayList;
    import java.util.concurrent.ExecutionException;
    import java.util.concurrent.ExecutorService;
    import java.util.concurrent.Executors;
    import java.util.concurrent.Future;
    
    public class Main
    {
        public static void main(String[] args) throws InterruptedException, ExecutionException
        {
            int totalTasks = 100;
    
            ExecutorService newFixedThreadPool = Executors.newFixedThreadPool(totalTasks);
    
            System.out.println("Step 1 - Starting parallel tasks");
    
            ArrayList<Future<Integer>> tasks = new ArrayList<Future<Integer>>();
            for (int i = 0; i < totalTasks; i++) {
                tasks.add(newFixedThreadPool.submit(new Process(i)));
            }
    
            long ts = System.currentTimeMillis();
            System.out.println("Step 2 - Wait for processes to finish...");
    
            boolean tasksCompleted;
            do {
                tasksCompleted = true;
    
                for (Future<Integer> task : tasks) {
                    if (!task.isDone()) {
                        tasksCompleted = false;
                        Thread.sleep(10);
                        break;
                    }
                }
    
            } while (!tasksCompleted);
    
            System.out.println(String.format("Step 2 - End in '%.3f' seconds", (System.currentTimeMillis() - ts) / 1000.0));
    
            System.out.println("Step 3 - All processes finished to run, let's collect results...");
    
            Integer sum = 0;
    
            for (Future<Integer> task : tasks) {
                sum += task.get();
            }
    
            System.out.println(String.format("Total final sum is: %d", sum));
        }
    }
    

    这里是 Process 类:

    package com.example;
    
    import java.util.concurrent.Callable;
    
    public class Process implements Callable<Integer>
    {
        private Integer value;
    
        public Process(Integer value)
        {
            this.value = value;
        }
    
        public Integer call() throws Exception
        {
            Long sleepTime = (long)(Math.random() * 2000);
    
            System.out.println(String.format("Starting process with value %d, sleep time %d", this.value, sleepTime));
    
            Thread.sleep(sleepTime);
    
            System.out.println(String.format("Stopping process with value %d", this.value));
    
            return value * 2;
        }
    }
    

    希望这会有所帮助。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-11-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多