【问题标题】:Reduce returns unpredictable results for parallel stream减少并行流的返回不可预测的结果
【发布时间】:2015-09-28 19:26:48
【问题描述】:

我用 java 流 reduce 编写了以下代码示例:

Person reducedPerson = Person.getPersons().stream()
                .parallel()  //will return surprising result
                .reduce(new Person(), (intermediateResult, p2) -> {
                            intermediateResult.setAge(intermediateResult.getAge() + p2.getAge());
                            return intermediateResult;
                        },
                        (ir1, ir2) -> {
                            ir1.setAge(ir1.getAge() + ir2.getAge());
                            return ir1;
                        });
        System.out.println(reducedPerson);

型号:

public class Person {

    String name;

    Integer age;

    public Person() {
        age = 0;
        name = "default";
    }

    //...
    public Person(String name, Integer age) {
        this.name = name;
        this.age = age;
    }

    public static Collection<Person> getPersons() {
        List<Person> persons = new ArrayList<>();
        persons.add(new Person("Vasya", 12));
        persons.add(new Person("Petya", 32));
        persons.add(new Person("Serj", 10));
        persons.add(new Person("Onotole", 18));
        return persons;
    }
}

每个代码示例执行返回不同的结果:

示例: Person{name='default', age=256}

Person{name='default', age=248}

我已经在combiner 中定位了这个问题,因为在顺序流中代码可以正确执行。

请帮助纠正组合器。

附言

预期结果:名为“default”且年龄 72 岁的人(列表中所有 pepsons 的总和)

附言

与 reduce 结果相同的 Integer 代码可以正常工作:

Integer age = Person.getPersons().stream()
                .parallel()
                .reduce(0, (intermediateResult, p2) -> {
                    intermediateResult = intermediateResult + p2.getAge();
                    return intermediateResult;
                }, (ir1, ir2) -> {
                    System.out.println("combiner");
                    ir1 = ir1 + ir2;
                    return ir1;
                });
        System.out.println(age);

【问题讨论】:

  • 那是因为您的reduce 调用违反了合同的所有可能部分。请阅读文档。你不应该改变流中的对象。
  • 是的。它应该。你不合并,你改变其中一个对象然后返回它。
  • 不清楚为什么您使用reduce 的三参数版本而不是两个参数,因为累加器与您的输入类型相同。
  • @gstackoverflow 是的,这是一个独立的问题,但无论如何你都应该这样做。也就是说,为什么不写更合乎逻辑、更高效、更直接的 Person.getPersons().stream().parallel().mapToInt(Person::getAge).sum() 呢?为什么需要生成 Person 而不是总年龄,这样更合乎逻辑?
  • 您是否在同一个程序中多次调用它?还是您只是重新运行整个程序并获得不同的结果?

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


【解决方案1】:

要执行可变归约,请使用collect

reducedPerson = Person.getPersons().parallelStream()
        .collect(
                Person::new,
                (p, q) -> p.setAge(p.getAge() + q.getAge()),
                (p, q) -> p.setAge(p.getAge() + q.getAge())
        );

collect 专门设计用于安全地累积到可变容器中,即使是并行。

【讨论】:

    【解决方案2】:

    正如 Boris 所说,问题在于流中的突变。

    大多数流操作接受描述用户指定的参数 行为,例如传递给的 lambda 表达式 w -> w.getWeight() 上例中的 mapToInt。为了保持正确的行为,这些 行为参数:

    • 必须不干扰(它们不会修改流源);并在
    • 大多数情况必须是无状态的(它们的结果不应依赖于任何 在流管道执行期间可能会发生变化的状态)。

    https://docs.oracle.com/javase/8/docs/api/java/util/stream/Stream.html

    这里是使用reduce 的版本,更直接的版本是使用maptoint 和sum。

    class gstackoverflow{
      public static void main(String... args) {
        Person reducedPerson = Person.getPersons().stream()
            .parallel()  //will NOT return surprising result
            .reduce(new Person("default",0),
                (ir1, ir2) -> //no longer mutates
                    new Person(String.join(",", ir1.getName(), ir2.getName()), ir1.getAge() + ir2.getAge())
            );
        System.out.println(reducedPerson);
    
        //here is a clean(er) way to do it:
        int totalAge = Person.getPersons().stream()
            .parallel()  //will NOT return surprising result
            .mapToInt(Person::getAge)
            .sum();
        System.out.println(totalAge);
      }
    }
    
    class Person {//no longer mutable
    
      public String getName() {
        return name;
      }
    
      public Integer getAge() {
        return age;
      }
    
      final String name;
    
      final Integer age;
    
      //no args constructor removed
      public Person(String name, Integer age) {
        this.name = name;
        this.age = age;
      }
    
      public static Collection<Person> getPersons() {
        List<Person> persons = new ArrayList<>();
        persons.add(new Person("Vasya", 12));
        persons.add(new Person("Petya", 32));
        persons.add(new Person("Serj", 10));
        persons.add(new Person("Onotole", 18));
        return persons;
      }
      @Override
      public String toString() {
        final StringBuilder sb = new StringBuilder("Person{");
        sb.append("name='").append(name).append('\'');
        sb.append(", age=").append(age);
        sb.append('}');
        return sb.toString();
      }
    }
    

    【讨论】:

      猜你喜欢
      • 2011-11-30
      • 2015-08-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-01-07
      • 1970-01-01
      相关资源
      最近更新 更多