【问题标题】:using java streams in parallel with collect(supplier, accumulator, combiner) not giving expected results与 collect(supplier, accumulator, combiner) 并行使用 java 流没有给出预期的结果
【发布时间】:2017-05-29 01:53:58
【问题描述】:

我正在尝试查找给定字符串中的单词数。下面是它的顺序算法,效果很好。

public int getWordcount() {

        boolean lastSpace = true;
        int result = 0;

        for(char c : str.toCharArray()){
            if(Character.isWhitespace(c)){
                lastSpace = true;
            }else{
                if(lastSpace){
                    lastSpace = false;
                    ++result;
                }
            }
        }

        return result;

    }

但是,当我尝试使用 Stream.collect(supplier, accumulator, combiner) 方法“并行化”它时,我得到 wordCount = 0。我使用不可变类 (WordCountState) 只是为了保持字数的状态.

代码:

public class WordCounter {
    private final String str = "Java8 parallelism  helps    if you know how to use it properly.";

public int getWordCountInParallel() {
        Stream<Character> charStream = IntStream.range(0, str.length())
                                                .mapToObj(i -> str.charAt(i));

        WordCountState finalState = charStream.parallel()                                             
                                              .collect(WordCountState::new,
                                                        WordCountState::accumulate,
                                                        WordCountState::combine);

        return finalState.getCounter();
    }
}

public class WordCountState {
    private final boolean lastSpace;
    private final int counter;
    private static int numberOfInstances = 0;

public WordCountState(){
        this.lastSpace = true;
        this.counter = 0;
        //numberOfInstances++;
    }

    public WordCountState(boolean lastSpace, int counter){
        this.lastSpace = lastSpace;
        this.counter = counter;
        //numberOfInstances++;
    }

//accumulator
    public WordCountState accumulate(Character c) {


        if(Character.isWhitespace(c)){
            return lastSpace ? this : new WordCountState(true, counter);
        }else{
            return lastSpace ? new WordCountState(false, counter + 1) : this;
        }   
    }

    //combiner
    public WordCountState combine(WordCountState wordCountState) {  
        //System.out.println("Returning new obj with count : " + (counter + wordCountState.getCounter()));
        return new WordCountState(this.isLastSpace(), 
                                    (counter + wordCountState.getCounter()));
    }

我发现上述代码存在两个问题: 1. 创建的对象数(WordCountState)大于字符串中的字符数。 2. 结果始终为 0。 3.根据累加器/消费者文档,累加器不应该返回无效吗?即使我的累加器方法返回一个对象,编译器也不会抱怨。

任何线索我可能偏离了轨道?

更新: 使用的解决方案如下 -

public int getWordCountInParallel() {
        Stream<Character> charStream = IntStream.range(0, str.length())
                                                .mapToObj(i -> str.charAt(i));


        WordCountState finalState = charStream.parallel()
                                              .reduce(new WordCountState(),
                                                        WordCountState::accumulate,
                                                        WordCountState::combine);

        return finalState.getCounter();
    }

【问题讨论】:

  • @Holger/Malte/Eugene :感谢您的详细意见。我应该早点说清楚......这个程序的重点是强调并行性,如果不使用上下文实现(即在正确的位置而不是在单词之间分割字符串),会违背目的并给出不正确的结果。另一种解决方案是使用 Spliterator 来确保拆分不会发生在单词的中间。我发现 collect 中的方法引用不期望任何返回值。因此,累积返回的值落在黑洞中。将 collect() 更改为 reduce() 解决了问题
  • 接受 Holger 的解决方案,因为它与我试图实现的目标非常相似。
  • 这是一个有趣的方面。如果您只是计算单词,识别单词成为主要任务,实现者将尝试使用并行流(例如通过归约)来解决。但是如果你使用单词作为起点,即流的元素,在后续的中间流步骤中会经历大量的处理,自然会尝试在较低级别的单词边界处进行拆分,以创建高效的并行词流.与Pattern.compile("\\s+").splitAsStream(str) 类似,但具有更好的并行性能……

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


【解决方案1】:

您始终可以调用方法并忽略其返回值,因此在使用方法引用时允许这样做是合乎逻辑的。因此,在需要消费者时,只要参数匹配,创建对非void方法的方法引用是没有问题的。

您使用不可变的WordCountState 类创建的是一个reduction 操作,即它将支持像

这样的用例
Stream<Character> charStream = IntStream.range(0, str.length())
                                        .mapToObj(i -> str.charAt(i));

WordCountState finalState = charStream.parallel()
        .map(ch -> new WordCountState().accumulate(ch))
        .reduce(new WordCountState(), WordCountState::combine);

而collect 方法支持可变归约,其中容器实例(可能与结果相同)被修改。

您的解决方案中仍然存在逻辑错误,因为每个 WordCountState 实例都假设前面有一个空格字符,但不知道实际情况,也没有尝试在组合器中解决此问题。

解决和简化这个问题的方法是:

public int getWordCountInParallel() {
    return str.codePoints().parallel()
        .mapToObj(WordCountState::new)
        .reduce(WordCountState::new)
        .map(WordCountState::getResult).orElse(0);
}


public class WordCountState {
    private final boolean firstSpace, lastSpace;
    private final int counter;

    public WordCountState(int character){
        firstSpace = lastSpace = Character.isWhitespace(character);
        this.counter = 0;
    }

    public WordCountState(WordCountState a, WordCountState b) {
        this.firstSpace = a.firstSpace;
        this.lastSpace = b.lastSpace;
        this.counter = a.counter + b.counter + (a.lastSpace && !b.firstSpace? 1: 0);
    }
    public int getResult() {
        return counter+(firstSpace? 0: 1);
    }
}

如果您担心WordCountState 实例的数量,请注意与您最初的方法相比,此解决方案没有创建多少Character 实例。

但是,如果你将WordCountState 重写为可变结果容器,这个任务确实适合可变归约:

public int getWordCountInParallel() {
    return str.codePoints().parallel()
        .collect(WordCountState::new, WordCountState::accumulate, WordCountState::combine)
        .getResult();
}


public class WordCountState {
    private boolean firstSpace, lastSpace=true, initial=true;
    private int counter;

    public void accumulate(int character) {
        boolean white=Character.isWhitespace(character);
        if(lastSpace && !white) counter++;
        lastSpace=white;
        if(initial) {
            firstSpace=white;
            initial=false;
        }
    }
    public void combine(WordCountState b) {
        if(initial) {
            this.initial=b.initial;
            this.counter=b.counter;
            this.firstSpace=b.firstSpace;
            this.lastSpace=b.lastSpace;
        }
        else if(!b.initial) {
            this.counter += b.counter;
            if(!lastSpace && !b.firstSpace) counter--;
            this.lastSpace = b.lastSpace;
        }
    }
    public int getResult() {
        return counter;
    }
}

注意如何使用 int 一致地表示 unicode 字符,允许使用 CharSequence 的 codePoint() 流,这不仅更简单,而且可以处理基本多语言平面之外的字符,并且可能更有效,因为它不需要装箱到Character 实例。

【讨论】:

  • 太棒了!为什么在世界上我也不认为建议进行可变减少?有你回来感觉很好
【解决方案2】:

当您实现stream().collect(supplier, accumulator, combiner) 时,它们确实返回void(组合器和累加器)。问题在于:

  collect(WordCountState::new,
          WordCountState::accumulate,
          WordCountState::combine)

在您的情况下实际上意味着(只是累加器,但组合器也是如此):

     (wordCounter, character) -> {
              WordCountState state = wc.accumulate(c);
              return;
     }

这确实不是一件容易的事。假设我们有两种方法:

public void accumulate(Character c) {
    if (!Character.isWhitespace(c)) {
        counter++;
    }
}

public WordCountState accumulate2(Character c) {
    if (Character.isWhitespace(c)) {
        return lastSpace ? this : new WordCountState(true, counter);
    } else {
        return lastSpace ? new WordCountState(false, counter + 1) : this;
    }
}

对于他们来说,下面的代码可以正常工作,但仅适用于方法引用,不适用于 lambda 表达式。

BiConsumer<WordCountState, Character> cons = WordCountState::accumulate;

BiConsumer<WordCountState, Character> cons2 = WordCountState::accumulate2;

您可以通过implementes BiConsumer 的类来想象它略有不同,例如:

 BiConsumer<WordCountState, Character> clazz = new BiConsumer<WordCountState, Character>() {
        @Override
        public void accept(WordCountState state, Character character) {
            WordCountState newState = state.accumulate2(character);
            return;
        }
    };

因此,您的 combine 和 accumulate 方法需要更改为:

public void combine(WordCountState wordCountState) {
    counter = counter + wordCountState.getCounter();
}


public void accumulate(Character c) {
    if (!Character.isWhitespace(c)) {
        counter++;
    }
}

【讨论】:

    【解决方案3】:

    首先,使用input.split("\\s+").length 之类的东西来获取字数不是更容易吗?

    如果这是关于流和收集器的练习,让我们讨论一下您的实现。您已经指出了最大的错误:您的累加器和组合器不应返回新实例。 collect 的签名告诉您它期望 BiConsumer,它不返回任何内容。因为您在累加器中创建了新对象,所以您永远不会增加您的收集器实际使用的 WordCountState 对象的数量。通过在组合器中创建一个新对象,您将放弃您可能取得的任何进展。这也是您在输入中创建的对象多于字符的原因:每个字符一个,然后一些用于返回值。

    查看这个改编的实现:

    public static class WordCountState
    {
        private boolean lastSpace = true;
        private int     counter   = 0;
    
        public void accumulate(Character character)
        {
            if (!Character.isWhitespace(character))
            {
                if (lastSpace)
                {
                    counter++;
                }
                lastSpace = false;
            }
            else
            {
                lastSpace = true;
            }
        }
    
        public void combine(WordCountState wordCountState)
        {
            counter += wordCountState.counter;
        }
    }
    

    在这里,我们不是在每一步都创建新对象,而是改变我们拥有的对象的状态。我认为您尝试创建新对象是因为您的 Elvis 操作员强迫您返回某些内容和/或您无法更改实例字段,因为它们是最终的。不过,它们不需要是最终的,您可以轻松更改它们。

    现在按顺序运行这个经过调整的实现可以正常工作,因为我们很好地逐个查看字符并最终得到 11 个单词。

    但同时,它失败了。似乎它为每个字符创建了一个新的WordCountState,但并没有计算所有字符,最终为 29(至少对我而言)。这表明您的算法存在一个基本缺陷:拆分每个字符不能并行工作。想象一下输入abc abc,它应该得到2。如果你并行执行并且没有指定如何分割输入,你最终可能会得到这些块:ab, c a, bc,加起来会是4。

    问题在于,通过字符之间的并行化(即在单词中间),您可以使单独的 WordCountStates 相互依赖(因为他们需要知道哪个在他们之前以及它是否以空白字符)。这会破坏并行性并导致错误。

    除此之外,实现Collector 接口可能比提供三个方法更容易:

    public static class WordCountCollector
        implements Collector<Character, SimpleEntry<AtomicInteger, Boolean>, Integer>
    {
        @Override
        public Supplier<SimpleEntry<AtomicInteger, Boolean>> supplier()
        {
            return () -> new SimpleEntry<>(new AtomicInteger(0), true);
        }
    
        @Override
        public BiConsumer<SimpleEntry<AtomicInteger, Boolean>, Character> accumulator()
        {
            return (count, character) -> {
                if (!Character.isWhitespace(character))
                {
                    if (count.getValue())
                    {
                        String before = count.getKey().get() + " -> ";
                        count.getKey().incrementAndGet();
                        System.out.println(before + count.getKey().get());
                    }
                    count.setValue(false);
                }
                else
                {
                    count.setValue(true);
                }
            };
        }
    
        @Override
        public BinaryOperator<SimpleEntry<AtomicInteger, Boolean>> combiner()
        {
            return (c1, c2) -> new SimpleEntry<>(new AtomicInteger(c1.getKey().get() + c2.getKey().get()), false);
        }
    
        @Override
        public Function<SimpleEntry<AtomicInteger, Boolean>, Integer> finisher()
        {
            return count -> count.getKey().get();
        }
    
        @Override
        public Set<java.util.stream.Collector.Characteristics> characteristics()
        {
            return new HashSet<>(Arrays.asList(Characteristics.CONCURRENT, Characteristics.UNORDERED));
        }
    }
    

    我们使用一对 (SimpleEntry) 来保持计数和关于最后一个空格的知识。这样,我们不需要在收集器本身中实现状态或为其编写 param 对象。你可以像这样使用这个收集器:

    return charStream.parallel().collect(new WordCountCollector());
    

    此收集器的并行化比初始实现更好,但由于您的方法中提到的弱点,结果仍然不同(主要在 14 到 16 之间)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-05-14
      • 2020-05-29
      • 1970-01-01
      • 2019-06-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多