【问题标题】:How to make Spring cloud stream Kafka streams binder retry processing a message if a failure occurs during the processing step?如果在处理步骤中发生故障,如何使 Spring Cloud Stream Kafka Stream binder 重试处理消息?
【发布时间】:2020-06-03 16:21:09
【问题描述】:

我正在使用 Spring Cloud Stream 开发 Kafka Streams。在消息处理应用程序中,可能会产生错误。所以消息不应该再次提交和重试。

我的申请方法-

@Bean
public Function<KStream<Object, String>, KStream<String, Long>> process() {
return (input) -> {
KStream<Object, String> kt = input.flatMapValues(v -> Arrays.asList(v.toUpperCase().split("\\W+")));
KGroupedStream<String, String> kgt =kt.map((k, v) -> new KeyValue<>(v, v)).groupByKey(Grouped.with(Serdes.String(), Serdes.String()));
KTable<Windowed<String>, Long> ktable = kgt.windowedBy(TimeWindows.of(500)).count();
KStream<String, WordCount> kst =ktable.toStream().map((k,v) -> {
WordCount wc = new WordCount();
wc.setWord(k.key());
wc.setCount(v);
wc.setStart(new Date(k.window().start()));
wc.setEnd(new Date(k.window().end()));

dao.insert(wc);

return new KeyValue<>(k.key(),wc);
});
return kst.map((k,v) -> new KeyValue<>(k, v.getCount()));
};
}

这里如果DAO插入方法失败,消息不应该发布到输出主题,并且应该重试相同消息的处理。

我们如何配置 kafka 流绑定器来做到这一点?非常感谢您对此提供任何帮助。

【问题讨论】:

    标签: java apache-kafka-streams spring-cloud-stream event-driven-design spring-cloud-stream-binder-kafka


    【解决方案1】:

    Spring Cloud Stream Kafka Streams binder 本身在执行业务逻辑时不提供这种重试机制。但是,解决此用例的一种方法可能是将您的关键调用(在这种情况下为dao.insert())包装在您在本地定义的RetryTemplate 中。这是一个可能的实现,它使用 1 秒的退避策略重试 10 次。如果您正在尝试此解决方案,请确保从您的主要业务逻辑中提取与 RetryTemplate 相关的公共代码。我还没有尝试过,但它应该可以工作。

    KStream<String, WordCount> kst =ktable.toStream().map((k,v) -> {
      WordCount wc = new WordCount();
      ...
    
      org.springframework.retry.support.RetryTemplate retryTemplate = new 
       RetryTemplate();
    
      RetryPolicy retryPolicy = new SimpleRetryPolicy(10);
      FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
      backOffPolicy.setBackOffPeriod(1000);
    
      retryTemplate.setBackOffPolicy(backOffPolicy);
      retryTemplate.setRetryPolicy(retryPolicy);
    
      retryTemplate.execute(context -> {
        try {
          dao.insert(wc);
        }
        catch (Exception e) {
          throw new IllegalStateException(..);
       }
      });
    
      return new KeyValue<>(k.key(),wc);
    });
    
    

    dao insert 操作重试10次后的事件,如果仍然失败,则抛出异常终止应用程序,此时不会提交偏移量。在重新启动时,在修复了基础问题后,您的应用程序仍应从此偏移量继续。

    【讨论】:

    • 嗨,RetryTemplate 的执行方法中的“上下文”对象是什么?如何创建它?
    • 即由框架创建并传入的RetryContext。我们只是在应用程序中使用它作为 lambda 参数。您不需要显式创建它,框架会处理它。
    猜你喜欢
    • 1970-01-01
    • 2019-10-03
    • 1970-01-01
    • 2021-04-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-12-19
    • 2018-03-09
    相关资源
    最近更新 更多