【问题标题】:Nonblocking ReentrantLock with Reactor使用 Reactor 的非阻塞 ReentrantLock
【发布时间】:2018-10-25 22:09:59
【问题描述】:

我需要限制同时处理相同资源的客户端数量
所以我试图实现模拟到

lock.lock();
try {
     do work
} finally {
    lock.unlock();
}

但使用 Reactor 库以非阻塞方式。 我有类似的东西。

但我有一个问题:
有没有更好的方法来做到这一点
或者也许有人知道已实施的解决方案
或者也许这不是在响应式世界中应该这样做的,对于此类问题还有另一种方法?

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.EmitterProcessor;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;

import javax.annotation.Nullable;
import java.time.Duration;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;

public class NonblockingLock {
    private static final Logger LOG = LoggerFactory.getLogger(NonblockingLock.class);

    private String currentOwner;
    private final AtomicInteger lockCounter = new AtomicInteger();
    private final FluxSink<Boolean> notifierSink;
    private final Flux<Boolean> notifier;
    private final String resourceId;

    public NonblockingLock(String resourceId) {
        this.resourceId = resourceId;
        EmitterProcessor<Boolean> processor = EmitterProcessor.create(1, false);
        notifierSink = processor.sink(FluxSink.OverflowStrategy.LATEST);
        notifier = processor.startWith(true);
    }

    /**
     * Nonblocking version of
     * <pre><code>
     *     lock.lock();
     *     try {
     *         do work
     *     } finally {
     *         lock.unlock();
     *     }
     * </code></pre>
     * */
    public <T> Flux<T> processWithLock(String owner, @Nullable Duration tryLockTimeout, Flux<T> work) {
        Objects.requireNonNull(owner, "owner");
        return notifier.filter(it -> tryAcquire(owner))
                .next()
                .transform(locked -> tryLockTimeout == null ? locked : locked.timeout(tryLockTimeout))
                .doOnSubscribe(s -> LOG.debug("trying to obtain lock for resourceId: {}, by owner: {}", resourceId, owner))
                .doOnError(err -> LOG.error("can't obtain lock for resourceId: {}, by owner: {}, error: {}", resourceId, owner, err.getMessage()))
                .flatMapMany(it -> work)
                .doFinally(s -> {
                    if (tryRelease(owner)) {
                        LOG.debug("release lock resourceId: {}, owner: {}", resourceId, owner);
                        notifierSink.next(true);
                    }
                });
    }

    private boolean tryAcquire(String owner) {
        boolean acquired;
        synchronized (this) {
            if (currentOwner == null) {
                currentOwner = owner;
            }
            acquired = currentOwner.equals(owner);
            if (acquired) {
                lockCounter.incrementAndGet();
            }
        }
        return acquired;
    }

    private boolean tryRelease(String owner) {
        boolean released = false;
        synchronized (this) {
            if (currentOwner.equals(owner)) {
                int count = lockCounter.decrementAndGet();
                if (count == 0) {
                    currentOwner = null;
                    released = true;
                }
            }
        }
        return released;
    }
}

这就是我认为它应该如何工作的方式

@Test
public void processWithLock() throws Exception {
    NonblockingLock lock = new NonblockingLock("work");
    String client1 = "client1";
    String client2 = "client2";
    Flux<String> requests = getWork(client1, lock)
            //emulate async request for resource by another client
            .mergeWith(Mono.delay(Duration.ofMillis(300)).flatMapMany(it -> getWork(client2, lock)))
            //emulate async request for resource by the same client
            .mergeWith(Mono.delay(Duration.ofMillis(400)).flatMapMany(it -> getWork(client1, lock)));
    StepVerifier.create(requests)
            .expectSubscription()
            .expectNext(client1)
            .expectNext(client1)
            .expectNext(client1)
            .expectNext(client1)
            .expectNext(client1)
            .expectNext(client1)
            .expectNext(client2)
            .expectNext(client2)
            .expectNext(client2)
            .expectComplete()
            .verify(Duration.ofMillis(5000));
}
private static Flux<String> getWork(String client, NonblockingLock lock) {
    return lock.processWithLock(client, null,
            Flux.interval(Duration.ofMillis(300))
                    .take(3)
                    .map(i -> client)
                    .log(client)
    );
}

【问题讨论】:

  • 您能否描述一个您打算使用这种锁的真实场景?我的意思是,您尝试实现的目标或多或少很清楚,但为什么呢?
  • 我有一个带有内存存储的 Web 应用程序,我需要提供其数据的一致性。因此,只有一个客户端可以对“事务”中的数据应用更改是必要的。另一个用例 - 是制作资源池。因此,如果目前没有可用资源 - 只需等到有可用资源即可
  • 也可以用于非阻塞缓存。由于 Mono.cache() 具有在没有值信号的情况下保留 Error 或 Complete 的特殊性,因此如果我只想用数据缓存成功的结果,这是不可取的行为。而且 Mono.cache() 不如阻塞缓存(如番石榴缓存)灵活。因此,有了这样的锁,我可以使用阻塞缓存进行数据存储,并在成功进行昂贵操作的非阻塞重新计算后填充它。我认为还有一些用例,但令我惊讶的是它还没有实现。所以我已经填写了我做错了什么。
  • 我从@brian-clozel 和 alexander-pankin 看到了来自Cache the result of a Mono from a WebClient call... 的答案,但是如果有 10 个同时请求,他们的解决方案将进行 10 次重新计算(调用远程服务),这是对服务器的浪费和客户资源,如果他们最终得到相同的结果。但是使用这样的 Lock 可以只进行 1 次昂贵的调用,而其他订阅者将等待结果
  • 我确实在我的一个项目中使用 CacheMono 锁定了具有相同参数的远程服务的独占调用。不要认为这会是您更慷慨的问题的好答案,但我可以在几天内分享。

标签: java project-reactor


【解决方案1】:

我有一个解决方案,可以独占调用具有相同参数的远程服务。也许它对你的情况有帮助。

它基于即时tryLock,如果资源忙则错误,Mono.retryWhen“等待”释放。

所以我有 LockData 类用于锁定的元数据

public final class LockData {
    // Lock key to identify same operation (same cache key, for example).
    private final String key;
    // Unique identifier for equals and hashCode.
    private final String uuid;
    // Date and time of the acquiring for lock duration limiting.
    private final OffsetDateTime acquiredDateTime;
    ...
}

LockCommand 接口是对 LockData 的阻塞操作的抽象

public interface LockCommand {

    Tuple2<Boolean, LockData> tryLock(LockData lockData);

    void unlock(LockData lockData);
    ...
}

UnlockEventsRegistry 接口是解锁事件监听器收集器的抽象。

public interface UnlockEventsRegistry {
    // initialize event listeners collection when acquire lock
    Mono<Void> add(LockData lockData);

    // notify event listeners and remove collection when release lock
    Mono<Void> remove(LockData lockData);

    // register event listener for given lockData
    Mono<Boolean> register(LockData lockData, Consumer<Integer> unlockEventListener);
}

Lock 类可以用锁包装源 Mono,解锁和用解锁包装 CacheMono 编写器。

public final class Lock {
    private final LockCommand lockCommand;
    private final LockData lockData;
    private final UnlockEventsRegistry unlockEventsRegistry;
    private final EmitterProcessor<Integer> unlockEvents;
    private final FluxSink<Integer> unlockEventSink;

    public Lock(LockCommand lockCommand, String key, UnlockEventsRegistry unlockEventsRegistry) {
        this.lockCommand = lockCommand;
        this.lockData = LockData.builder()
                .key(key)
                .uuid(UUID.randomUUID().toString())
                .build();
        this.unlockEventsRegistry = unlockEventsRegistry;
        this.unlockEvents = EmitterProcessor.create(false);
        this.unlockEventSink = unlockEvents.sink();
    }

    ...

    public final <T> Mono<T> tryLock(Mono<T> source, Scheduler scheduler) {
        return Mono.fromCallable(() -> lockCommand.tryLock(lockData))
                .subscribeOn(scheduler)
                .flatMap(isLocked -> {
                    if (isLocked.getT1()) {
                        return unlockEventsRegistry.add(lockData)
                                .then(source
                                        .switchIfEmpty(unlock().then(Mono.empty()))
                                        .onErrorResume(throwable -> unlock().then(Mono.error(throwable))));
                    } else {
                        return Mono.error(new LockIsNotAvailableException(isLocked.getT2()));
                    }
                });
    }

    public Mono<Void> unlock(Scheduler scheduler) {
        return Mono.<Void>fromRunnable(() -> lockCommand.unlock(lockData))
                .then(unlockEventsRegistry.remove(lockData))
                .subscribeOn(scheduler);
    }

    public <KEY, VALUE> BiFunction<KEY, Signal<? extends VALUE>, Mono<Void>> unlockAfterCacheWriter(
            BiFunction<KEY, Signal<? extends VALUE>, Mono<Void>> cacheWriter) {
        Objects.requireNonNull(cacheWriter);
        return cacheWriter.andThen(voidMono -> voidMono.then(unlock())
                .onErrorResume(throwable -> unlock()));
    }

    public final <T> UnaryOperator<Mono<T>> retryTransformer() {
        return mono -> mono
                .doOnError(LockIsNotAvailableException.class,
                        error -> unlockEventsRegistry.register(error.getLockData(), unlockEventSink::next)
                                .doOnNext(registered -> {
                                    if (!registered) unlockEventSink.next(0);
                                })
                                .then(Mono.just(2).map(unlockEventSink::next)
                                        .delaySubscription(lockCommand.getMaxLockDuration()))
                                .subscribe())
                .doOnError(throwable -> !(throwable instanceof LockIsNotAvailableException),
                        ignored -> unlockEventSink.next(0))
                .retryWhen(errorFlux -> errorFlux.zipWith(unlockEvents, (error, integer) -> {
                    if (error instanceof LockIsNotAvailableException) return integer;
                    else throw Exceptions.propagate(error);
                }));
    }
}

现在如果我必须用 CacheMono 包裹我的 Mono 并锁定,我可以这样做:

private Mono<String> getCachedLockedMono(String cacheKey, Mono<String> source, LockCommand lockCommand, UnlockEventsRegistry unlockEventsRegistry) {
    Lock lock = new Lock(lockCommand, cacheKey, unlockEventsRegistry);

    return CacheMono.lookup(CACHE_READER, cacheKey)
            // Lock and double check
            .onCacheMissResume(() -> lock.tryLock(Mono.fromCallable(CACHE::get).switchIfEmpty(source)))
            .andWriteWith(lock.unlockAfterCacheWriter(CACHE_WRITER))
            // Retry if lock is not available
            .transform(lock.retryTransformer());
}

您可以在 GitHub 上找到带有示例的代码和测试

【讨论】:

  • 感谢@Alexander,这种方法可以解决我的问题,它与我第一次尝试创建非阻塞锁非常相似,但不是 Mono.error() - retry() on failed tryLock() 我只需返回 Mono.empty(),然后执行 Mono.repeatWhenEmpty(..)。所以它不会产生不必要的异常,这只是重试的信号。这种方法的第二个问题是它只是一个手工制作的非阻塞循环,可能不如事件驱动方法那么有效,而且它增加了获取结果的延迟,在最坏的情况下,它与实际计算的延迟大约为 100 ms结果。
  • 如果我在锁忙时返回空单声道,如果源单声道也是空的,我无法处理。可以增强重试功能以对锁释放做出反应,这是一个很好的评论。对于我的情况,非阻塞循环是可以的,因为并发执行很少,延迟并不重要。我将使用我的重试功能来对锁释放做出反应,并可能改变我对循环的想法。感谢您的反馈。
  • 更改了此解决方案以对解锁事件做出反应。
  • @AlexanderPankin 感谢分享这个,现在 wince EmitProcessor 已被弃用,你有没有 EmitProcessor 的等效示例(使用 Sinks.Many)?
  • @hmble,你好。我这周只有智能手机,所以我稍后会分享新的解决方案。现在您可以在 GitHub 上找到一些示例(链接在我的答案末尾)github.com/alex-pumpkin/reactor-lock
猜你喜欢
  • 2020-10-28
  • 1970-01-01
  • 1970-01-01
  • 2020-11-07
  • 1970-01-01
  • 2023-01-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多