超时后停止从 KafkaReceiver 消费
Stop consuming from KafkaReceiver after a timeout
我有一个通用的休息控制器:
private final KafkaReceiver<String, Domain> receiver;
@GetMapping(produces = MediaType.APPLICATION_STREAM_JSON_VALUE)
public Flux<Domain> produceFluxMessages() {
return receiver.receive().map(ConsumerRecord::value)
.timeout(Duration.ofSeconds(2));
}
我想要实现的是在一段时间内从 Kafka 主题收集消息,然后停止消费并认为此流量已完成。如果我删除超时并在浏览器中打开它,我将永远收到消息,下载永远不会停止。并且在 2 秒后超时消耗停止,但我得到一个例外:
java.util.concurrent.TimeoutException: Did not observe any item or terminal signal within 2000ms in 'map' (and no fallback has been configured)
有没有办法在超时后成功完成Flux?
timeout()
方法有多个重载 - 您使用的是在超时时抛出异常的标准方法。
相反,只需使用 overloaded timeout method 提供一个空的默认发布者以回退到:
timeout(Duration.ofSeconds(2), Mono.empty())
(请注意,在一般情况下,您可以显式捕获 TimeoutException
并使用 onErrorResume(TimeoutException.class, e -> Mono.empty())
回退到空发布者,但如果可能的话,这远不如使用上述选项更可取。)
我有一个通用的休息控制器:
private final KafkaReceiver<String, Domain> receiver;
@GetMapping(produces = MediaType.APPLICATION_STREAM_JSON_VALUE)
public Flux<Domain> produceFluxMessages() {
return receiver.receive().map(ConsumerRecord::value)
.timeout(Duration.ofSeconds(2));
}
我想要实现的是在一段时间内从 Kafka 主题收集消息,然后停止消费并认为此流量已完成。如果我删除超时并在浏览器中打开它,我将永远收到消息,下载永远不会停止。并且在 2 秒后超时消耗停止,但我得到一个例外:
java.util.concurrent.TimeoutException: Did not observe any item or terminal signal within 2000ms in 'map' (and no fallback has been configured)
有没有办法在超时后成功完成Flux?
timeout()
方法有多个重载 - 您使用的是在超时时抛出异常的标准方法。
相反,只需使用 overloaded timeout method 提供一个空的默认发布者以回退到:
timeout(Duration.ofSeconds(2), Mono.empty())
(请注意,在一般情况下,您可以显式捕获 TimeoutException
并使用 onErrorResume(TimeoutException.class, e -> Mono.empty())
回退到空发布者,但如果可能的话,这远不如使用上述选项更可取。)