【问题标题】:How to implement delayed retry mechanism in a Kafka streaming (kstreams) application built in spring cloud stream?如何在spring cloud stream内置的Kafka流(kstreams)应用中实现延迟重试机制?
【发布时间】:2021-06-22 04:39:29
【问题描述】:
我有一个类似的场景。
有4个主题处理主题、异常主题、重试主题和拒绝主题。
我有一个使用 Kstream 的处理器的 spring 云流应用程序。该处理器从异常主题中读取消息,并基于每条消息中可用的标志
为 retry Topic 和 Reject topic 创建两个 kstream 分支。
现在需要做的是重试主题中存在的任何消息都必须等待特定的时间段,然后才能将其推回处理主题。
任何人都可以帮助我在春季云流中的kafka流应用程序中执行此操作的最佳设计或解决方案是什么。
是否可以使用 Flux 和 Mono 设计异步机制。任何资源或指导都会有很大帮助。谢谢
【问题讨论】:
标签:
java
apache-kafka
reactive-programming
apache-kafka-streams
spring-cloud-stream
【解决方案1】:
我们不能在 Spring Cloud Stream Kafka Streams 应用程序中混合反应类型。基本上,您只能将 Kafka Streams 类型(例如 KStream 或 KTable)作为输入/输出绑定。如果您想在将其发送到处理主题之前引入延迟,为什么不能使用ScheduleExecutorService 之类的东西,然后仅在初始延迟后调用服务?这是一个潜在解决方案的示例伪代码。
ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor();
...
...
.branch(
(k, v) -> {
if (reject flag found) {
return true;
}
}
(k, v) -> executorService.schedule(() -> true, initialDelay, TimeUnit.SECONDS));
在第一个分支中,我们检查记录是否绑定到拒绝状态,如果是,则立即将其发送到被拒绝的主题。否则,它针对重试主题,延迟一些值(我猜这是在您的应用程序外部配置的),然后返回布尔标志。延迟完成后,布尔值true将返回,到达分支的记录将被发送到重试主题。
请记住,我没有尝试过此代码,因此请注意第二个过滤器中的任何竞争条件,我们以异步方式将事情交给调度程序,但该线程上下文仍应保留当前的 @987654326 @ 一对。请用一些测试数据试试这个,看看它是否符合您的要求。