【问题标题】:Multiple Threads working within single Stream多个线程在单个 Stream 中工作
【发布时间】:2018-06-18 10:25:26
【问题描述】:

我想知道以下是否有效:

class SomeCalc {

    AtomicLong   someResultA;
    AtomicDouble someResultB;

    public SomeCalc(List<Something> someList) {

        final Stream<Something> someStream = someList.stream();

        new Thread(() -> {
            long result = someStream.mapToLong(something -> something.getSomeLong()).sum();
            someResultA.set(result);
        }).start();

        new Thread(() -> {
            double result = someStream.mapToDouble(something -> (double) something.getSomeLong() / something.getAnotherLong()).sum();
            someResultB.set(result);
        }).start();
    }
}

只要结果(如本例中所示)是原子的,那么这会很好吗?或者在此过程中会出现任何ConcurrentAccessModification 异常,否则会出现其他问题?

【问题讨论】:

    标签: java multithreading concurrency java-8 java-stream


    【解决方案1】:

    它不会工作。你会得到一个:

    IllegalStateException: Stream has already been operated upon or closed
    

    因为每个流只能遍历一次。这也可以在java.util.stream.Stream&lt;T&gt;javadoc中读到:

    一个流应该操作 [...] 只操作一次。例如,这排除了“分叉”流,其中相同的源提供两个或多个管道,或者同一流的多次遍历。如果流实现检测到流正在被重用,它可能会抛出 IllegalStateException。但是,由于某些流操作可能会返回其接收者而不是新的流对象,因此可能无法在所有情况下都检测到重用。

    我猜你的例子是简化的。因此,您正在准备流,然后让 2 个线程执行其他操作。

    如果你还想保留它。您可以创建该流的供应商,然后让线程获得 2 个不同的实例:

    final Supplier<Stream<Something>> provider = () -> someList.stream(); // and maybe more operations
    
    new Thread(() -> {
        long result = provider.get()  // get an instance
           .mapToLong(something -> something.getSomeLong())
           .sum();
        someResultA.set(result);
    }).start();
    
    new Thread(() -> {
        double result = provider.get() // get another instance
             .mapToDouble(something -> (double) something.getSomeLong() / something.getAnotherLong())
             .sum();
        someResultB.set(result);
    }).start();
    

    【讨论】:

    • 解释了问题,但没有提供解决方案或替代方案。你可以改进一下。尽管如此,很好的答案。
    【解决方案2】:

    并发修改的可能性更不用说,这段代码不会工作,因为它在同一个流上调用多个终端操作。

    两个线程最终会在同一个流上调用sum 两次,这是不允许的。但是不必担心并发修改原子长对象,因为它旨在处理来自多个线程的修改调用。

    由于您已经在使用原子数据结构,您可以将其更改为类似这样,以使其在一次遍历中工作:

    someList.stream()
    .mapToLong(something -> something.getSomeLong())
    .forEach(entry -> {
        someResultA.incrementAndGet(something.getSomeLong(entry));
    
        //I assume the below call exists...
        someResultB.incrementAndGet((double) something.getSomeLong() / something.getAnotherLong());
    });
    

    如果使用不同线程的原因是并行性,那么您可以将其设为并行流:

    someList.stream().parallel()
     .map...
    

    由于同步,维护原子值最终可能会更慢(与求和和设置相比),但它为您提供了并行流的单次遍历。

    【讨论】:

    • 好答案 :),您可能想写 .parallelStream() 而不是 .stream().parallel()
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-07
    • 1970-01-01
    相关资源
    最近更新 更多