【发布时间】:2021-12-30 10:56:28
【问题描述】:
我在一个 webflux 项目中遇到了一个场景,这让我很困惑,因为我还是 project-reactor 的新手。
代码如下:
class InputArg {
private List<String> list;
private ConcurrentMap<String, Integer> map;
// getter setters...
}
final int CONCURRENCY_LEVEL = 4;
0. var someInputArgwithNullMap = ... // init one InputArg object with null map
1. Mono.just(someInputArgwithNullMap)
2. .flatMapMany((InputArg inputArg) -> {
3. // inputArg.map == null
4. inputArg.setMap(new ConcurrentHashMap<>());
5
6. return Flux.fromIterable(inputArg.getList()).
7. .flatMap(str -> {
8. // some async call to external service with this str argument
9. Mono<String> message = ...
10. int randomInteger = ... // code to get random
11. return message.map(msg -> {
12. inputArg.getMap().putIfAbsent(msg, randomInteger);
13.
14. return msg
15. });
16. }, CONCURRENCY_LEVEL);
17. })
18. .map(...) //some operation that modify inputArg again
19. .subscribe(...); //subscriber which needs updated someInputArgwithNullMap and also print
20. //every msg.
我的问题是:
由于第 7 行中的 Flux.flatMap 是异步处理,并且可能使用不同于父 Mono 流的多线程运行,因此我在第 4 行中使用 concurrentHashMap 并在第 12 行中使用原子操作来保证在 inputArg 中对这个共享映射进行同步操作。
但是,
-
在多线程上下文下,第 4 行和第 12 行之间的对象 inputArg 本身怎么样。所以第 4 行和第 12 行中的线程可能仍然不同,第 4 行中的更改是否对第 12 行可见,并且即使在 DEC 或 ARM 等弱排序机器中也不可能在任何 PC 模型下抛出 NPE 或其他竞争条件?
-
使 Mono.flatMap(第 2 行)下的更改对转换后的 Flux 的第 17 行地图可见的内部机制?
我已经通过互联网阅读了几篇文章,现在我知道 onNext 信号之间存在一些顺序保证,因此某些同步/易失性或内存屏障的设置可与新的 Java 9+ VarHandle api 与 onNext 调用兼容。但是如何从我的代码sn-p在一些运算符之间的这个先决条件中推断出确切的内存排序效果。还有一个约束:作为兼容的原因,我们不能将 InputArg 修改为不可变的。等待并感谢您的回答。
2021-11-20 更新:
因为第 9 行可能会影响多线程,所以进一步澄清这个异步调用:
一个。使用 webClient 对另一个响应式服务(如 WebFlux)的异步调用;或
b.使用 Spring Data Redis 对外部缓存服务的异步调用。
2021 年 11 月 21 日更新:
经过一番研究,对于问题2,我自己的理解:
a) 根据 WebFlux doc 的声明(请参阅索引页上 Pivotal, Inc 的版权https://docs.spring.io/spring-framework/docs/current/reference/html/index.html)https://docs.spring.io/spring-framework/docs/current/reference/html/web-reactive.html:
“在运行时,会形成一个反应式管道,其中数据在不同的阶段按顺序处理。这样做的一个关键好处是,它使应用程序不必保护可变状态,因为该管道中的应用程序代码永远不会同时调用。”
似乎这句话也暗示了父/子流之间也存在发生前的关系(好像在不同的阶段)。
b.我检查了地图运算符的内部源代码,这种方法(参见 Apache 2 许可证https://www.apache.org/licenses/LICENSE-2.0)
public Object scanUnsafe(Attr key)
返回 Attr.RunStyle.SYNC 以便此运算符不会发生线程更改,这也是大多数其他运算符的默认值。这样地图操作员就会在第 2 行看到该 flatMapMany 所做的任何更改。
对于我的问题 1:
一个。第 9 行的调用可能仅在它具有一个 publishOn 运算符时才更改线程(subscribeOn 将通过其调度使用一个线程制作整个父子流),并且在我检查了 publishOn 的内部代码并且它在一个易失性值上有一个写入/获取关系其获取父 onNext 信号与其下游发送信号(可能会更改线程)之间的字段,因此第 4 行的更改对第 12 行的 getMap() 方法可见,因此我的代码是线程安全的。
b.来自 WebFlux 官方文档的声明以及反应流规范的声明(参见其 MIT No Attribution 版权所有https://github.com/aws/mit-0)https://github.com/reactive-streams/reactive-streams-jvm#1.3
“发送给订阅者的 onSubscribe、onNext、onError 和 onComplete 必须连续发送。”
似乎还可以保证父/子 onNext 信号之间的先发生关系。
c。无法重新订阅相同的订阅。
如果我的理解有错误,请纠正我,谢谢。
【问题讨论】:
标签: java multithreading reactive-programming spring-webflux project-reactor