【问题标题】:How can you check when a Java 8 Stream.forEach() finishes iterating?如何检查 Java 8 Stream.forEach() 何时完成迭代?
【发布时间】:2021-04-07 22:09:28
【问题描述】:

我想使用 Java 8 Streams 提供的并行性,但我还需要以特定顺序完成某些操作,否则一切都会中断。问题是使用流意味着代码现在是异步的 [注意:这是不正确的],我无法弄清楚如何仅在完成对整个集合的迭代时才发生某些事情。 [注意:这是自动发生的]

现在,我的代码如下:

    public void iterateOverMap(Map<String, String> m)
    {
        AtomicInteger count = new AtomicInteger(0);
        
        m.keySet().stream().forEach((k) -> {
                    Object o = m.get(k);
                    
                    // do stuff with o
                    
                    count.incrementAndGet();
                });
        
        // spin up a new thread to check if the Stream is done
        new Thread(() -> {
            for (;;)
            {
                if (count.intValue() >= map.size())
                    break;
            }
            afterFinishedIterating();
        }).start();
    }

我不喜欢为了跟踪这件事而不得不启动一个新线程或阻塞主线程的想法,但我想不出我还能怎么做。有谁知道更好的选择吗?

谢谢!

【问题讨论】:

  • 它不是异步的,它是并行的。同步和并行。它会在迭代完成后返回。另外,不要忙着等待。
  • @BoristheSpider 我只实现了计数器,因为在forEach 完成之前执行afterFinishedIterating() 的程序出现问题。另外,显然我不应该忙着等待,但我应该怎么做呢?
  • 没有什么是异步的;我不知道你想解决什么问题。
  • 另外,使用Map#entrySet() 而不是遍历键并重新获取每个条目。
  • @ABadHaiku 重读Map#entrySet() 的定义。您会看到每个Map.Entry&lt;&gt; 都包含键和值。

标签: java foreach parallel-processing java-stream thread-synchronization


【解决方案1】:

Stream 处理是同步的。

如果您想要一个如何跟踪Stream 进度的示例,您可以使用peek() 中间操作,但请记住,它最好用于调试目的

示例取自my other answer

Stream<MyData> myStream = readData();
final AtomicInteger loader = new AtomicInteger();
int fivePercent = elementsCount / 20;
MyResult result = myStream
    .map(row -> process(row))
    .peek(stat -> {
        if (loader.incrementAndGet() % fivePercent == 0) {
            System.out.println(loader.get() + " elements on " + elementsCount + " treated");
            System.out.println((5*(loader.get() / fivePercent)) + "%");
        }
    })
    .reduce(MyStat::aggregate);

【讨论】:

  • @BoristheSpider 我改了,希望现在措辞没问题
  • 我刚刚又查看了Collections.synchronizedList() 的 javadocs,您不知道吗,但是“当通过迭代器、拆分器或溪流。”事实证明,我毕竟不是疯了!相反,我只是使用了错误的列表......
  • @ABadHaiku 不用担心,编程让我们认为自己疯了,但我们中的大多数人完全理智:D 很高兴你能到达那里
  • @ABadHaiku 1) Don’t use LinkedList 除非您有信心成为百万分之一的真正需要它的人。 2) 不要将列表包装在synchronizedList 中,而是使用正确的终端操作(99.9% 的情况下:not forEach)。在这里,.filter(x -&gt; /* return whether error has been found */) .collect(Collectors.toList())。这提供了正确的结果,即使以正确的遇到顺序(与 forEach 添加到同步列表不同)没有同步开销。
  • @ABadHaiku 按照我在上一条评论中提供的链接并仔细阅读答案。即使在LinkedList 中添加和删除也只有在非常特殊的情况下才会更好(这与Java 的Collection API 与之交互的方式有关,链表的概念 很好)。具有多个链接操作的流管道仍然是 一个 循环。见Loop fusion of Stream in Java-8 (how it works internally)。和Java streams lazy vs fusion vs short-circuiting
【解决方案2】:

原来我遇到的问题是我正在使用的线程安全但不确定的迭代“Collections.synchronizedList(new LinkedList&lt;&gt;())”列表,而不是Stream 使用。 Stream 不像我想象的那样是异步的,所以答案只是“你不必这样做;它会为你做的”。

【讨论】:

    猜你喜欢
    • 2016-04-09
    • 2012-08-20
    • 2019-02-04
    • 2017-08-08
    • 1970-01-01
    • 2011-03-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多