【问题标题】:Memory Ordering Effects for Mono/Flux StreamMono/Flux 流的内存排序效果
【发布时间】: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 中对这个共享映射进行同步操作。

但是,

  1. 在多线程上下文下,第 4 行和第 12 行之间的对象 inputArg 本身怎么样。所以第 4 行和第 12 行中的线程可能仍然不同,第 4 行中的更改是否对第 12 行可见,并且即使在 DEC 或 ARM 等弱排序机器中也不可能在任何 PC 模型下抛出 NPE 或其他竞争条件?

  2. 使 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


    【解决方案1】:

    flatMap 将生成一个遵循 Reactive Streams 规范的 Flux,其中每个传出的 onNext 与另一个之间具有发生前的关系。这是通过排空并发 MPSC 队列在内部完成的。

    不过,这种保证会在每个内部 Mono 将其值发布到 flatMap 协调器时停止。

    例如,完全有可能 2 个内部 Monos 被协调器订阅并在两个不同的线程中并行执行。这些 Mono 使用的共享资源可能会处于争用状态。

    从传入的InputArg 产生内部Mono 的Function 也是串行执行的,但这仅代表组装阶段,而不是所述内部 Monos 的执行阶段。

    注意:我在这里使用Mono,因为这与您的示例相匹配。任何Publisher都是如此

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-06-07
      • 2020-01-28
      • 1970-01-01
      • 2021-09-30
      • 2020-09-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多