Project Reactor响应式管道中可变对象的内存可见性问题咨询
响应式管道中可变对象的内存可见性问题解答
1. 对可见性问题的理解是否正确?
是的,你的理解完全正确。
Reactor管道虽然是顺序执行的,但每个操作阶段可能在不同线程上调度执行(比如你的示例里doSomething在remote-calls-1、onErrorResume在db-calls-2、postprocess在local-processing-3)。根据Java内存模型(JMM),如果没有同步机制(比如volatile修饰字段、同步块),线程间的内存操作没有happens-before关系保证:
- 线程A修改
myclass.errors的操作,可能仅停留在A的本地缓存,没有同步到主内存 - 后续线程B读取
myclass.errors时,可能从自己的本地缓存读取到旧值,看不到线程A的修改
你无法复现问题,是因为现代CPU的缓存一致性协议(比如MESI)、JVM的即时编译优化,或者测试场景下线程切换的时机刚好让修改被同步了,但这属于未定义行为——依赖这种偶然的一致性会埋下生产环境的隐患,比如高负载下缓存同步延迟、JIT优化重排导致的可见性问题,这类问题通常难以复现和排查。
2. 可复现的可见性问题示例
要复现这个问题,需要放大JVM的优化效果,同时让线程的缓存隔离更明显。下面是一个调整后的测试用例,通过循环重复执行、禁用某些JIT优化,来触发可见性问题:
import org.junit.jupiter.api.Test; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; // 手动实现getter/setter,避免Lombok可能带来的额外同步逻辑 class MyClass { private String errors; // 非volatile字段 private String someField; public String getErrors() { return errors; } public void setErrors(String errors) { this.errors = errors; } public String getSomeField() { return someField; } public void setSomeField(String someField) { this.someField = someField; } } public class VisibilityTest { private final Scheduler remoteScheduler = Schedulers.newParallel("remote", 1); private final Scheduler errorScheduler = Schedulers.newParallel("error", 1); private final Scheduler postScheduler = Schedulers.newParallel("post", 1); private String toErrors(Throwable t) { return t.getMessage(); } private <T> Mono<T> runOnScheduler(Mono<T> mono, Scheduler scheduler) { return mono.publishOn(scheduler); } private Mono<String> doSomething(MyClass myClass) { return runOnScheduler(Mono.just(myClass.getSomeField()) .doOnNext(__ -> System.out.println("doSomething on " + Thread.currentThread().getName())), remoteScheduler); } private Mono<String> doSomethingElse(String s, MyClass myClass) { return runOnScheduler(Mono.just(s) .doOnNext(__ -> { System.out.println("doSomethingElse on " + Thread.currentThread().getName()); throw new RuntimeException("forced error"); }), errorScheduler); } private Mono<MyClass> postProcess(String fallback, MyClass myClass) { return runOnScheduler(Mono.fromCallable(() -> { System.out.println("postProcess on " + Thread.currentThread().getName()); // 循环读取,放大缓存不一致的概率 for (int i = 0; i < 100000; i++) { if (myClass.getErrors() != null) { break; } } System.out.println("postProcess read errors: " + myClass.getErrors()); return myClass; }), postScheduler); } public Mono<MyClass> process(MyClass myClass) { return doSomething(myClass) .flatMap(result -> doSomethingElse(result, myClass)) .onErrorResume(e -> { System.out.println("onErrorResume on " + Thread.currentThread().getName()); myClass.setErrors(toErrors(e)); return Mono.just("fallback"); }) .flatMap(fallback -> postProcess(fallback, myClass)); } @Test void testVisibilityIssue() throws InterruptedException { // 重复执行多次,增加触发概率 for (int i = 0; i < 100; i++) { CountDownLatch latch = new CountDownLatch(1); MyClass myClass = new MyClass(); myClass.setSomeField("test"); process(myClass).subscribe( res -> { if (res.getErrors() == null) { System.err.println("TEST FAILED: postProcess didn't see errors! Iteration: " + i); } latch.countDown(); }, err -> latch.countDown() ); latch.await(5, TimeUnit.SECONDS); } } }
运行注意事项
要提高复现概率,需要给JVM添加以下启动参数,禁用一些会自动同步内存的优化:
-XX:-UseCompressedOops -XX:-TieredCompilation -Xint
预期现象
运行一段时间后,会出现TEST FAILED: postProcess didn't see errors!的输出——这说明postProcess所在的线程没有看到onErrorResume线程对errors字段的修改,也就是发生了内存可见性问题。
内容的提问来源于stack exchange,提问作者Evgeny
相关产品推荐
相关产品推荐

