【问题标题】:Non-blocking lock without using EmitProcessor不使用 EmitProcessor 的非阻塞锁
【发布时间】:2021-07-25 12:45:10
【问题描述】:

我有一个资源 (r) 有两个非阻塞操作(这是来自我无法更改的外部库):

Mono<String> op1(String someInput);
Mono<String> op2(String otherInput);

op1() 和 op2() 可以在不同的线程上调用但不能同时执行。何时调用 op1 取决于外部因素,关于 op2 也是如此。两个操作无关(op1 可能被调用 100 次,op2 可能被调用 1 次)。如何保证op1和op2互斥处理。如果 op1 /op2 被阻塞,java 'synchronized' 会解决它。如果 op1 和 op2 处理同时发生,则外部库方法会失败。

如何在不使用 EmitProcessor(已弃用)的情况下强制执行此类同步,以便可以从不同的调度程序线程调用 op1 和 op2?或者 WebFlux api 中是否有内置的标准解决方案来解决这种情况?

(有一个使用 EventProcessor 的解决方案,但希望避免它,因为 EventProcessor 已被弃用 Nonblocking ReentrantLock with Reactor

【问题讨论】:

  • 您如何决定“反之亦然”部分?你不能决定总是先订阅op1然后订阅op2吗?在这种情况下,op1.then(op2) 可以解决问题。
  • @SimonBaslé op1 和 op2 是回调,它们的执行取决于消费者,即消费者1 控制何时调用 op1,消费者2 控制 op2。所以它们没有被并排调用,因此不能使用 then、map、delayUntil 等。我正在寻找stackoverflow.com/questions/52998809/… 的解决方案,但不使用 EventProcessor,不确定是否已经存在?谢谢(添加了更多相关信息以清除此问题)
  • 您是在谈论 mutual exclusion (mutex) 而非 Mono 订阅吗?我不确定反应堆有什么用的。但是,您可以尝试强制使用单线程调度程序以避免并发执行。它的好处是简化了锁定逻辑。请注意,我既不确定性能影响,也不确定下游的传播行为。您应该查看subscribeOnpublishOn 方法。注意:这正是 UI 技术人员所发生的事情。他们通常使用必须在其上进行所有显示更新的“UI 线程”。
  • 您的解释自相矛盾。首先,您说 op1 正在控制何时调用 op2。那么你说op1和op2是由消费者控制的。那么它会是哪一个呢? op2 是由 op1 还是由 consumer2 控制的?
  • op1 和 op2 调用的频率取决于 api 调用者,op1 和 op2 不相互依赖,但它们不能同时处理。很抱歉混淆了,但顺序调用并不是指先调用 op1 再调用 op2(顺序处理是指非并发处理)。我还将更新问题以澄清这一点。

标签: spring-webflux project-reactor


【解决方案1】:

目前 Reactor 中没有为此内置任何内容。

其他答案中的解决方案可能会更新为Sinks.Many新API,看来相关项目确实已更新:https://github.com/alex-pumpkin/reactor-lock

也可以使用https://github.com/reactor/reactor-pool,但这有点过头了。

【讨论】:

    【解决方案2】:

    以下解决方案确保 op1 和 op2 的互斥处理:

    public class Locker {
        private final AtomicBoolean locked = new AtomicBoolean(true);
        private final Flux<Boolean> notifier;
        private final Sinks.Many<Boolean> notifierSink;
    
        public Locker() {
            this.notifierSink = Sinks.many().multicast().onBackpressureBuffer(1, false);
            this.notifier = notifierSink.asFlux();
            this.notifierSink.emitNext(true, Sinks.EmitFailureHandler.FAIL_FAST);
        }
    
        public <T> Flux<T> lockThenProcess(Duration lockTimeout, Flux<T> job) {
            return notifier.filter(v -> obtainLock())
                    .next()
                    .transform(locked -> lockTimeout == null ? locked : locked.timeout(lockTimeout))
                    .doOnSubscribe(s -> log.debug("obtaining lock"))
                    .doOnError(th -> log.error("can't obtain lock: " + th.getMessage(), th))
                    .flatMapMany(v -> job)
                    .doFinally(s -> {
                        if (releaseLock()) {
                            log.debug("released lock");
                            notifierSink.emitNext(true, Sinks.EmitFailureHandler.FAIL_FAST);
                        }
                    });
        }
    
        private synchronized boolean obtainLock() {
            return locked.getAndSet(false);
        }
    
        private synchronized boolean releaseLock() {
            locked.set(true);
            return locked.get();
        }
    }
    

    然后,调用 op1 和 op2(在任何线程上)如下:

    op1Trigger.concatMap(v -> locker.lockThenProcess(Duration.ofMinutes(1), r.op1(input).flux()))
    

    op2Trigger.concatMap(v -> locker.lockThenProcess(Duration.ofMinutes(1), r.op2(input).flux()))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2013-03-13
      • 2014-04-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-19
      • 2016-07-06
      相关资源
      最近更新 更多