【发布时间】: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