【发布时间】: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