【问题标题】:Reprocessing a message when an error occurs while processing it in the Kafka Stream在 Kafka Stream 中处理消息时发生错误时重新处理消息
【发布时间】: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


    【解决方案1】:

    如果在 Kafka Streams 处理期间发生任何未捕获的异常,您的流将状态更改为 ERROR 并停止消费发生错误的分区的传入消息。 您需要自己捕获异常。重试可以通过以下方式实现:1)使用 Spring RetryTemplate 调用外部微服务(但请记住,您将延迟从特定分区消费消息),或 2)将失败的消息推送到另一个主题以供以后重新处理(如你建议)


    更新自kafka-streams2.8.0

    由于kafka-streams2.8.0,您可以自动替换失败的流线程(由未捕获的异常引起) 使用KafkaStreams 方法void setUncaughtExceptionHandler(StreamsUncaughtExceptionHandler eh);StreamThreadExceptionResponse.REPLACE_THREAD。更多详情请关注Kafka Streams Specific Uncaught Exception Handler

    kafkaStreams.setUncaughtExceptionHandler(ex -> {
        log.error("Kafka-Streams uncaught exception occurred. Stream will be replaced with new thread", ex);
        return StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse.REPLACE_THREAD;
    });
    

    【讨论】:

      猜你喜欢
      • 2018-02-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-09-22
      • 2020-02-15
      • 1970-01-01
      • 2019-06-21
      • 1970-01-01
      相关资源
      最近更新 更多