【发布时间】:2019-05-28 12:15:48
【问题描述】:
我有一个基于 Kafka Streams 的简单 Spring 应用程序,它使用来自传入主题的消息,进行 map 转换并打印此消息。 KStream这样配置
@Bean
public KStream<?, ?> processingPipeline(StreamsBuilder builder, MyTransformer myTransformer,
PrintAction printAction, String topicName) {
KStream<String, JsonNode> source = builder.stream(topicName,
Consumed.with(Serdes.String(), new JsonSerde<>(JsonNode.class)));
// @formatter:off
source
.map(myTransformer)
.foreach(printAction);
// @formatter:on
return source;
}
在MyTransformer 内部,我正在调用此时可能会关闭的外部微服务。如果调用失败(通常抛出RuntimeException),我无法进行转换。
这里的问题是,如果在之前的处理过程中发生任何错误,有什么方法可以再次在 Streams 应用程序中重新处理消息?
根据我目前的研究,这里没有办法这样做,我唯一的可能是将消息推送到死信主题并尝试在将来再次失败时对其进行处理我再次将其推送到 DLT 并执行以这种方式重试。
【问题讨论】:
标签: java apache-kafka apache-kafka-streams spring-kafka