1

我有一个共同的休息控制器:

    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?

4

1 回答 1

3

该方法有多个重载timeout()- 您正在使用标准的在超时时引发异常的方法。

相反,只需使用重载的超时方法来提供一个空的默认发布者以回退到:

timeout(Duration.ofSeconds(2), Mono.empty())

(请注意,在一般情况下,您可以使用 显式捕获 TimeoutException并回退到空发布onErrorResume(TimeoutException.class, e -> Mono.empty())者,但这远不如在可能的情况下使用上述选项。)

于 2021-03-03T12:44:03.800 回答