【发布时间】:2020-12-16 04:52:45
【问题描述】:
我正在使用 Spring Cloud Stream Kafka Streams 编写一个 Java 应用程序。这是我正在使用的函数方法 sn-p:
@Bean
public Function<KStream<String, String>, KStream<String, String>> process() {
return input ->
input.transform(
() ->
new Transformer<String, String, KeyValue<String, String>>() {
ProcessorContext context;
@Override
public void init(ProcessorContext context) {
this.context = context;
}
@Override
public void close() {}
@Override
public KeyValue<String, String> transform(String key, String value) {
String result = fetch_data_from_database(key, value);
return new KeyValue<>(key, result);
}
});
fetch_data_from_database() 可以抛出异常。
如果 fetch_from_database() 出现异常,我如何停止处理入站 KStream(偏移量不应被提交)并使用相同的偏移量数据重试处理?
【问题讨论】:
-
stackoverflow.com/questions/51299528/… 这个问题也引起了同样的关注,但提供的答案并没有说明如何处理处理异常
标签: apache-kafka apache-kafka-streams spring-cloud-stream