【问题标题】:Project Reactor and the Java memory modelProject Reactor 和 Java 内存模型
【发布时间】:2019-04-06 00:07:46
【问题描述】:

我正在尝试了解 Project reactor 为应用程序代码提供的数据可见性保证。例如我希望下面的代码会失败,但在一百万次迭代后它不会。我正在更改线程 A 上典型 POJO 的状态并从线程 B 读取它。Reactor 是否保证 POJO 更改在线程中可见?

public class Main {
    public static void main(String[] args) {
        Integer result = Flux.range(1, 1_000_000)
                .map(i -> {
                    Data data = new Data();
                    data.setValue(i);
                    data.setValueThreeTimes(i);
                    data.setValueObj(i + i);
                    return data;
                })
                .parallel(250)
                .runOn(Schedulers.newParallel("par", 500))
                .map(d -> {
                    d.setValueThreeTimes(d.getValueThreeTimes() + d.getValue());
                    return d;
                })
                .sequential()
                .parallel(250)
                .runOn(Schedulers.newParallel("par", 500))
                .map(d -> {
                    d.setValueThreeTimes(d.getValueThreeTimes() + d.getValue());
                    return d;
                })
                //                .sequential()
                .map(d -> {
                    if (d.getValue() * 3 != d.getValueThreeTimes()) throw new RuntimeException("data corrupt error");
                    return d;
                })
                .reduce(() -> 0, (Integer sum, Data d) -> sum + d.getValueObj() + d.getValue())
                .sequential()
                .blockLast();
    }

    static class Data {
        private int value;
        private int valueThreeTimes;
        private Integer valueObj;

        public int getValueThreeTimes() {
            return valueThreeTimes;
        }

        public void setValueThreeTimes(int valueThreeTimes) {
            this.valueThreeTimes = valueThreeTimes;
        }

        public int getValue() {
            return value;
        }

        @Override
        public String toString() {
            return "Data{" +
                    "value=" + value +
                    ", valueObj=" + valueObj +
                    '}';
        }

        public void setValue(int value) {
            this.value = value;
        }

        public Integer getValueObj() {
            return valueObj;
        }

        public void setValueObj(Integer valueObj) {
            this.valueObj = valueObj;
        }
    }

    private static <T> T identityWithThreadLogging(T el, String operation) {
        System.out.println(operation + " -- " + el + " -- " +
                Thread.currentThread().getName());
        return el;
    }
}

【问题讨论】:

  • 您是否尝试过从链中删除sequential?
  • 刚试过没有sequential。返回正确的结果。当元素从一个运算符移动到下一个运算符时,似乎会应用内存围栏。
  • Idk 实现,但可能步骤是顺序的,并且步骤中的特定操作是并行的 - 因此你的“障碍”

标签: java reactive-programming project-reactor


【解决方案1】:

Reactive Streams 规范强制要求对于 Flux 或 Mono(Publisher),onNext 事件必须是连续的。

parallel() 是一个ParallelFlux,它通过分而治之的方式稍微放松了一点:你会得到多个“轨道”,每个轨道都单独遵守规范,但总体上不遵守规范(轨道之间的并行化)。

反过来,sequential() 回到 Flux 世界并引入内存屏障以确保生成的序列符合 RS 规范。

【讨论】:

  • 如果我要创建一个简单的 POJO 并在 map 阶段调用一些 setter,然后开始在 parallel rails 上进行处理,那么 setter 调用中的值是否会在并行 rails 中可见?在使用反应器时,我们是否必须使用 volatile 或 final 字段来获得一致的数据视图,因为它在流中发生变化?
  • parallel() 之前的阶段要么在同一个线程上执行,要么使用 volatile 触发内存屏障
猜你喜欢
  • 2018-11-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-01-07
  • 2013-12-31
  • 2018-09-26
  • 2017-05-18
相关资源
最近更新 更多