【发布时间】:2019-06-06 15:57:25
【问题描述】:
我对响应式编程和 Reactor 比较陌生。我有一种情况,我想在我的流中 bufferTimeout 值,同时 将其保持在我的控制之下(没有无限制的请求),因此我可以手动请求批量值。
以下示例说明了这一点:
BlockingQueue<Integer> queue = new LinkedBlockingQueue<>();
Flux<Object> flux = Flux.generate(sink -> {
try {
sink.next(queue.poll(10, TimeUnit.DAYS));
}
catch (InterruptedException e) {}
});
BaseSubscriber<List<Object>> subscriber = new BaseSubscriber<List<Object>>() {
protected void hookOnSubscribe(Subscription subscription) {
// Don't request unbounded
}
protected void hookOnNext(List<Object> value) {
System.out.println(value);
}
};
flux.subscribeOn(parallel())
.log()
.bufferTimeout(10, ofMillis(200))
.subscribe(subscriber);
subscriber.request(1);
// Offer a partial batch of values
queue.offer(1);
queue.offer(2);
queue.offer(3);
queue.offer(4);
queue.offer(5);
// Wait for timeout, expect [1, 2, 3, 4, 5] to be printed
Thread.sleep(500);
// Offer more values
queue.offer(6);
queue.offer(7);
queue.offer(8);
queue.offer(9);
queue.offer(10);
Thread.sleep(1000);
这是输出:
[DEBUG] (main) Using Console logging
[ INFO] (main) onSubscribe(FluxSubscribeOn.SubscribeOnSubscriber)
[ INFO] (main) request(10)
[ INFO] (parallel-1) onNext(1)
[ INFO] (parallel-1) onNext(2)
[ INFO] (parallel-1) onNext(3)
[ INFO] (parallel-1) onNext(4)
[ INFO] (parallel-1) onNext(5)
[1, 2, 3, 4, 5]
[ INFO] (parallel-1) onNext(6)
[ INFO] (parallel-1) onNext(7)
[ INFO] (parallel-1) onNext(8)
[ INFO] (parallel-1) onNext(9)
[ INFO] (parallel-1) onNext(10)
reactor.core.Exceptions$ErrorCallbackNotImplemented: reactor.core.Exceptions$OverflowException: Could not emit buffer due to lack of requests
我实际上是预料到的,因为我知道缓冲区订阅者将在上游请求 10 个值,这不知道超时并且无论如何都会产生所有这些值。 由于一旦超时完成,唯一的请求就完成了,稍后提供的值仍然会产生并溢出。
我想知道是否有可能防止在超时完成后产生剩余的值,或者在不失去控制的情况下缓冲它们。我试过了:
-
limitRate(1)在bufferTimeout之前,试图使缓冲区请求值“按需”。它确实请求了 10 次,因为缓冲区请求了 10 个值。 -
onBackpressureBuffer(10),因为如果我做对了,问题基本上是背压的定义。试图从超时请求中缓冲溢出的值,但这会请求无限的值,我想避免这种情况。
看起来我必须实现另一个bufferTimeout 实现,但有人告诉我编写发布者很难。我错过了什么吗?还是我做错了反应?
【问题讨论】:
-
我想知道是否是生产者的模拟(Flux.generate)在这里造成了混乱。你真正的制片人会是什么样子?您是真的从代码中生成数据,还是从数据库等外部来源获取数据?
-
来源是 Kafka,我使用reactor-kafka 连接到它。使用这个库,数据会定期从 Kafka 中获取。我可以通过设置较长的获取周期并手动发送一条消息来重现此问题,等待超时,然后发送另一条消息。
-
我举的例子不好吗?我应该把真实情况放在问题中吗?
-
这个例子还不错。只是不清楚通量的数据来自哪里。您是完全在本地生成元素,还是从其他地方获取它们?看看真实案例可能会有所帮助。
-
我已经在本地和远程进行了测试,包括 Kafka 模拟和真正的 Kafka。正如我可以通过示例展示的那样,即使没有 Kafka,我也可以始终如一地重现该问题。看来我必须改进处理流的方式。
标签: java reactive-programming project-reactor