【问题标题】:Java parallel stream internalsJava 并行流内部结构
【发布时间】:2017-11-26 19:40:43
【问题描述】:

我注意到,取决于 doSth() 方法的实现(如果线程休眠了恒定或随机的时间),并行流的执行方式不同。

例子:

import java.util.Random;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.IntStream;

import static java.lang.System.out;

public class AtomicInt {

    public static void main(String[] args) throws ExecutionException, InterruptedException {
        out.println("Result: " + count());
    }

    public static int count() throws ExecutionException, InterruptedException {
        ForkJoinPool forkJoinPool = new ForkJoinPool(10);

        AtomicInteger counter = new AtomicInteger(0);

        forkJoinPool.submit(() -> IntStream
                .rangeClosed(1, 20)
                .parallel()
                .map(i -> doSth(counter))
                .forEach(i -> out.println(">>>forEach: " + Thread.currentThread().getName() + " value: " + i))
        ).get();

        return counter.get();
    }

    private static int doSth(AtomicInteger counter) {
        try {
            out.println(">>doSth1: " + Thread.currentThread().getName());
            Thread.sleep(100 + new Random().nextInt(1000));
//            Thread.sleep(1000);
        } catch (InterruptedException e) {
            throw new RuntimeException(e);
        }

        int counterValue = counter.incrementAndGet();
        out.println(">>doSth2: " + Thread.currentThread().getName() + " value: " + counterValue);

        return counterValue;
    }
}

每个数字都按顺序处理:

>>doSth1: ForkJoinPool-1-worker-9
>>doSth1: ForkJoinPool-1-worker-8
>>doSth1: ForkJoinPool-1-worker-2
>>doSth1: ForkJoinPool-1-worker-1
>>doSth1: ForkJoinPool-1-worker-6
>>doSth1: ForkJoinPool-1-worker-11
>>doSth1: ForkJoinPool-1-worker-15
>>doSth1: ForkJoinPool-1-worker-4
>>doSth1: ForkJoinPool-1-worker-13
>>doSth1: ForkJoinPool-1-worker-10
>>doSth2: ForkJoinPool-1-worker-8 value: 1
>>>forEach: ForkJoinPool-1-worker-8 value: 1
>>doSth1: ForkJoinPool-1-worker-8
>>doSth2: ForkJoinPool-1-worker-15 value: 2
>>>forEach: ForkJoinPool-1-worker-15 value: 2
>>doSth1: ForkJoinPool-1-worker-15
>>doSth2: ForkJoinPool-1-worker-11 value: 3
>>>forEach: ForkJoinPool-1-worker-11 value: 3
>>doSth1: ForkJoinPool-1-worker-11
>>doSth2: ForkJoinPool-1-worker-2 value: 4
>>>forEach: ForkJoinPool-1-worker-2 value: 4
>>doSth1: ForkJoinPool-1-worker-2
>>doSth2: ForkJoinPool-1-worker-9 value: 5
>>>forEach: ForkJoinPool-1-worker-9 value: 5
>>doSth1: ForkJoinPool-1-worker-9
>>doSth2: ForkJoinPool-1-worker-11 value: 6
>>>forEach: ForkJoinPool-1-worker-11 value: 6
>>doSth1: ForkJoinPool-1-worker-11
>>doSth2: ForkJoinPool-1-worker-1 value: 7
>>>forEach: ForkJoinPool-1-worker-1 value: 7
>>doSth1: ForkJoinPool-1-worker-1
>>doSth2: ForkJoinPool-1-worker-15 value: 8
>>>forEach: ForkJoinPool-1-worker-15 value: 8
>>doSth1: ForkJoinPool-1-worker-15
>>doSth2: ForkJoinPool-1-worker-8 value: 9
>>>forEach: ForkJoinPool-1-worker-8 value: 9
>>doSth1: ForkJoinPool-1-worker-8
>>doSth2: ForkJoinPool-1-worker-13 value: 10
>>>forEach: ForkJoinPool-1-worker-13 value: 10
>>doSth1: ForkJoinPool-1-worker-13
>>doSth2: ForkJoinPool-1-worker-9 value: 11
>>>forEach: ForkJoinPool-1-worker-9 value: 11
>>doSth2: ForkJoinPool-1-worker-15 value: 12
>>>forEach: ForkJoinPool-1-worker-15 value: 12
>>doSth2: ForkJoinPool-1-worker-10 value: 13
>>>forEach: ForkJoinPool-1-worker-10 value: 13
>>doSth2: ForkJoinPool-1-worker-4 value: 14
>>>forEach: ForkJoinPool-1-worker-4 value: 14
>>doSth2: ForkJoinPool-1-worker-6 value: 15
>>>forEach: ForkJoinPool-1-worker-6 value: 15
>>doSth2: ForkJoinPool-1-worker-11 value: 16
>>>forEach: ForkJoinPool-1-worker-11 value: 16
>>doSth2: ForkJoinPool-1-worker-2 value: 17
>>>forEach: ForkJoinPool-1-worker-2 value: 17
>>doSth2: ForkJoinPool-1-worker-13 value: 18
>>>forEach: ForkJoinPool-1-worker-13 value: 18
>>doSth2: ForkJoinPool-1-worker-1 value: 19
>>>forEach: ForkJoinPool-1-worker-1 value: 19
>>doSth2: ForkJoinPool-1-worker-8 value: 20
>>>forEach: ForkJoinPool-1-worker-8 value: 20
Result: 20

当我将 doSth() 方法更改为始终睡眠 1 秒而不是随机时间时,结果是按顺序计算的:

>>doSth1: ForkJoinPool-1-worker-6
>>doSth1: ForkJoinPool-1-worker-1
>>doSth1: ForkJoinPool-1-worker-10
>>doSth1: ForkJoinPool-1-worker-2
>>doSth1: ForkJoinPool-1-worker-13
>>doSth1: ForkJoinPool-1-worker-15
>>doSth1: ForkJoinPool-1-worker-8
>>doSth1: ForkJoinPool-1-worker-11
>>doSth1: ForkJoinPool-1-worker-4
>>doSth1: ForkJoinPool-1-worker-9
>>doSth2: ForkJoinPool-1-worker-1 value: 1
>>doSth2: ForkJoinPool-1-worker-10 value: 2
>>doSth2: ForkJoinPool-1-worker-6 value: 3
>>>forEach: ForkJoinPool-1-worker-6 value: 3
>>>forEach: ForkJoinPool-1-worker-10 value: 2
>>>forEach: ForkJoinPool-1-worker-1 value: 1
>>doSth1: ForkJoinPool-1-worker-10
>>doSth1: ForkJoinPool-1-worker-6
>>doSth1: ForkJoinPool-1-worker-1
>>doSth2: ForkJoinPool-1-worker-15 value: 4
>>doSth2: ForkJoinPool-1-worker-9 value: 10
>>doSth2: ForkJoinPool-1-worker-8 value: 7
>>>forEach: ForkJoinPool-1-worker-8 value: 7
>>doSth1: ForkJoinPool-1-worker-8
>>doSth2: ForkJoinPool-1-worker-4 value: 9
>>>forEach: ForkJoinPool-1-worker-4 value: 9
>>doSth2: ForkJoinPool-1-worker-11 value: 8
>>doSth2: ForkJoinPool-1-worker-13 value: 6
>>doSth2: ForkJoinPool-1-worker-2 value: 5
>>>forEach: ForkJoinPool-1-worker-13 value: 6
>>>forEach: ForkJoinPool-1-worker-11 value: 8
>>doSth1: ForkJoinPool-1-worker-4
>>>forEach: ForkJoinPool-1-worker-9 value: 10
>>>forEach: ForkJoinPool-1-worker-15 value: 4
>>doSth1: ForkJoinPool-1-worker-9
>>doSth1: ForkJoinPool-1-worker-11
>>doSth1: ForkJoinPool-1-worker-13
>>>forEach: ForkJoinPool-1-worker-2 value: 5
>>doSth1: ForkJoinPool-1-worker-15
>>doSth1: ForkJoinPool-1-worker-2
>>doSth2: ForkJoinPool-1-worker-10 value: 12
>>>forEach: ForkJoinPool-1-worker-10 value: 12
>>doSth2: ForkJoinPool-1-worker-6 value: 11
>>doSth2: ForkJoinPool-1-worker-1 value: 13
>>>forEach: ForkJoinPool-1-worker-6 value: 11
>>>forEach: ForkJoinPool-1-worker-1 value: 13
>>doSth2: ForkJoinPool-1-worker-9 value: 15
>>doSth2: ForkJoinPool-1-worker-2 value: 20
>>>forEach: ForkJoinPool-1-worker-2 value: 20
>>doSth2: ForkJoinPool-1-worker-15 value: 19
>>>forEach: ForkJoinPool-1-worker-15 value: 19
>>doSth2: ForkJoinPool-1-worker-8 value: 14
>>doSth2: ForkJoinPool-1-worker-11 value: 17
>>doSth2: ForkJoinPool-1-worker-4 value: 16
>>doSth2: ForkJoinPool-1-worker-13 value: 18
>>>forEach: ForkJoinPool-1-worker-4 value: 16
>>>forEach: ForkJoinPool-1-worker-11 value: 17
>>>forEach: ForkJoinPool-1-worker-8 value: 14
>>>forEach: ForkJoinPool-1-worker-9 value: 15
>>>forEach: ForkJoinPool-1-worker-13 value: 18
Result: 20

这是巧合还是对这种行为有解释?

【问题讨论】:

  • 这确实很有趣
  • 这不是必须使用Random 的方式。为流的每个元素生成一个新的随机序列以仅从中选择一个值是一种非常糟糕的模式。
  • 同意,我想你不会看到 ThreadLocalRandom 的这种行为。
  • @dimo414 所以我也想...但是你会得到与ThreadLocalRandom 相同的输出
  • @Eugene:确实,随机的类型完全无关紧要。您甚至可以改用Thread.currentThread().hashCode()%100。这只是分散了实际发生的事情的注意力。

标签: java multithreading parallel-processing java-8 java-stream


【解决方案1】:

当您执行sleep 语句时,没有定义顺序。虽然通过IntStream.range 创建的流具有已定义的遇到顺序,但您正在通过忽略实际的int 值将操作变成无序操作。

产生可感知订单的第一个操作是counter.incrementAndGet()。在此之前,哪个线程到达该点以及它与哪个流元素相关联都无关紧要。此时它从AtomicInteger 获得它的号码。之后,仅使用该号码执行两个附加操作,使用该号码打印两条消息。对于这两个不同的结果,重要的是这三个动作,counter.incrementAndGet() 和打印这两个消息,是否被另一个线程拦截。

我们可以轻松地将这个场景分解为

AtomicInteger counter = new AtomicInteger();
ExecutorService es = Executors.newFixedThreadPool(20);
es.invokeAll(Collections.nCopies(20, () -> {
    out.println("1st: " + Thread.currentThread().getName());
    Thread.sleep(100 + new Random().nextInt(1000));
//    Thread.sleep(1000);
    int counterValue = counter.incrementAndGet();
    out.println("2nd: " + Thread.currentThread().getName() + " value: " + counterValue);
    out.println("3rd: " + Thread.currentThread().getName() + " value: " + counterValue);
    return null;
}));
es.shutdown();

请注意,对于invokeAll,根本没有定义的顺序,但正如前面所说,这并不重要。直到调用incrementAndGet() 时,任务才获得分配的序列号。行为与流示例相同。


虽然我一直强调并发并不意味着并行,但由于未指定的执行时间和线程调度行为,仍然很有可能启动简短的相同代码同时在没有后台活动给 CPU 内核带来不可预测的工作负载时真正并行运行。

当所有线程并行运行时,它们同时达到out.println的内部同步,只有一个线程可以继续,其他线程进入队列。然后,synchronized 的不公平性质开始发挥作用。任意线程将获胜,之后任意线程将被放回调度。这会导致数字以随机顺序打印。

当您让线程随机休眠一段时间时,它们不再完全并行运行,从而提高了在不同时间到达 print 语句的机会,从而能够无竞争地执行它们。哪个线程将首先到达这一点,这是随机的,但由于它们在睡眠后分配了编号,因此到达该点的第一个线程将获得第一号,依此类推。

【讨论】:

  • 这是否意味着 OP(例如我)非常幸运地看到了有序的输出?我真的试着理解你的这个答案:|
  • @Eugene:不,代码基本上是说“给我顺序的下一个数字”,然后立即打印该数字。看到它们乱序实际上是幸运的问题,因为只有当线程对齐、同时做同样的事情、在 print 语句中发生争用时才会发生这种情况。取消对齐它们,例如通过随机等待,比让它们对齐更容易。
  • 在路上或其他 - 据我所知,这是不确定的
【解决方案2】:

我怀疑这种行为的原因是在 java.util.Random 类 (Java 8) 中对 nextInt 的调用。

>     protected int next(int bits) {
>         long oldseed, nextseed;
>         AtomicLong seed = this.seed;
>         do {
>             oldseed = seed.get();
>             nextseed = (oldseed * multiplier + addend) & mask;
>         } while (!seed.compareAndSet(oldseed, nextseed));
>         return (int)(nextseed >>> (48 - bits));
>     }

如果你不执行java.util.Random.nextInt(int),而是预先计算随机值,结果仍然是乱序的——或者如果你创建自己的从java.util.Random派生的Random类并重写方法protected int next(int bits)返回任何常数整数。我试图预先计算随机值,如下所示:

private static final int[] randomIntervals = IntStream.range(0, 20)
            .map(i -> new Random().nextInt(1000) + 100)
            .toArray();

然后在doSth方法中使用:

private static int doSth(AtomicInteger counter) {
    try {
        Thread.sleep(randomIntervals[counter.intValue()]);
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }

    int counterValue = counter.incrementAndGet();
    out.println(">>doSth2: " + Thread.currentThread().getName() + " value: " + counterValue);

    return counterValue;
}

结果看起来类似于只使用Thread.sleep(1000) 的代码版本,例如:

>>doSth2: ForkJoinPool-1-worker-26 value: 8
>>doSth2: ForkJoinPool-1-worker-1 value: 4
>>>forEach: ForkJoinPool-1-worker-1 value: 4
>>doSth2: ForkJoinPool-1-worker-8 value: 7
>>>forEach: ForkJoinPool-1-worker-8 value: 7
>>doSth2: ForkJoinPool-1-worker-25 value: 1
>>>forEach: ForkJoinPool-1-worker-25 value: 1
>>doSth2: ForkJoinPool-1-worker-15 value: 6
>>>forEach: ForkJoinPool-1-worker-15 value: 6
>>doSth2: ForkJoinPool-1-worker-29 value: 5
>>>forEach: ForkJoinPool-1-worker-29 value: 5
>>doSth2: ForkJoinPool-1-worker-4 value: 2
>>>forEach: ForkJoinPool-1-worker-4 value: 2
>>doSth2: ForkJoinPool-1-worker-18 value: 3
>>>forEach: ForkJoinPool-1-worker-18 value: 3
>>>forEach: ForkJoinPool-1-worker-26 value: 8
>>doSth2: ForkJoinPool-1-worker-11 value: 9
>>>forEach: ForkJoinPool-1-worker-11 value: 9
>>doSth2: ForkJoinPool-1-worker-22 value: 10
>>>forEach: ForkJoinPool-1-worker-22 value: 10
>>doSth2: ForkJoinPool-1-worker-1 value: 11
>>doSth2: ForkJoinPool-1-worker-25 value: 12
>>>forEach: ForkJoinPool-1-worker-25 value: 12
>>doSth2: ForkJoinPool-1-worker-8 value: 13
>>>forEach: ForkJoinPool-1-worker-8 value: 13
>>>forEach: ForkJoinPool-1-worker-1 value: 11
>>doSth2: ForkJoinPool-1-worker-29 value: 14
>>>forEach: ForkJoinPool-1-worker-29 value: 14
>>doSth2: ForkJoinPool-1-worker-4 value: 15
>>>forEach: ForkJoinPool-1-worker-4 value: 15
>>doSth2: ForkJoinPool-1-worker-15 value: 16
>>>forEach: ForkJoinPool-1-worker-15 value: 16
>>doSth2: ForkJoinPool-1-worker-18 value: 17
>>>forEach: ForkJoinPool-1-worker-18 value: 17
>>doSth2: ForkJoinPool-1-worker-26 value: 18
>>>forEach: ForkJoinPool-1-worker-26 value: 18
>>doSth2: ForkJoinPool-1-worker-22 value: 19
>>>forEach: ForkJoinPool-1-worker-22 value: 19
>>doSth2: ForkJoinPool-1-worker-11 value: 20
>>>forEach: ForkJoinPool-1-worker-11 value: 20

另一个实验是只调用java.util.Random的构造函数,而不是nextInt。结果又是乱序的:

private static int doSth(AtomicInteger counter) {
    try {
        Random r = new Random();
        Thread.sleep(1000);
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }

    int counterValue = counter.incrementAndGet();
    out.println(">>doSth2: " + Thread.currentThread().getName() + " value: " + counterValue);

    return counterValue;
}

【讨论】:

  • 但是即使它被同步了,不同的线程也应该在不同的时间到达那个同步点
  • 其实自从 OP 一直在创建一个新的Random istance,这就更没有意义了。
  • 你是对的:它不是构造函数。这是对正在同步线程的 nextInt(1000) 的调用。我的假设不正确。
  • 更新回答指出问题不在java.util.Random的构造函数,而是在这个类的方法nextInt(int)。
猜你喜欢
  • 2013-01-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-02-03
  • 2012-12-12
  • 2011-10-03
  • 1970-01-01
相关资源
最近更新 更多