微信公众号搜"智元新知"关注
微信扫一扫可直接关注哦!

超时后停止从 KafkaReceiver 消费

如何解决超时后停止从 KafkaReceiver 消费

我有一个通用的休息控制器:

    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()) 并回退到空发布者,但这比尽可能使用上述选项要好得多。)

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


Selenium Web驱动程序和Java。元素在(x,y)点处不可单击。其他元素将获得点击?
Python-如何使用点“。” 访问字典成员?
Java 字符串是不可变的。到底是什么意思?
Java中的“ final”关键字如何工作?(我仍然可以修改对象。)
“loop:”在Java代码中。这是什么,为什么要编译?
java.lang.ClassNotFoundException:sun.jdbc.odbc.JdbcOdbcDriver发生异常。为什么?
这是用Java进行XML解析的最佳库。
Java的PriorityQueue的内置迭代器不会以任何特定顺序遍历数据结构。为什么?
如何在Java中聆听按键时移动图像。
Java“Program to an interface”。这是什么意思?