【问题标题】:How to Handle a Kafka Record with a Class-Level @KafkaListener with no @KafkaHandler for the Record Value如何处理具有类级别 @KafkaListener 的 Kafka 记录,而记录值没有 @KafkaHandler
【发布时间】:2021-01-28 18:30:06
【问题描述】:

通常,当我们定义类级别的@KafkaListener 和方法级别的@KafkaHandlers 时,我们可以定义一个默认的@KafkaHandler 来处理意外的负载。

https://docs.spring.io/spring-kafka/docs/current/reference/html/#class-level-kafkalistener

但是,如果我们没有默认方法怎么办?

【问题讨论】:

标签: apache-kafka spring-kafka


【解决方案1】:

在 2.6 及更高版本中,您可以配置 SeekToCurrentErrorHandler 以通过检查异常立即将此类消息发送到死信主题。

这是一个演示该技术的简单 Spring Boot 应用程序:

@SpringBootApplication
public class So59256214Application {

    public static void main(String[] args) {
        SpringApplication.run(So59256214Application.class, args);
    }

    @Bean
    public NewTopic topic1() {
        return TopicBuilder.name("so59256214").partitions(1).replicas(1).build();
    }

    @Bean
    public NewTopic topic2() {
        return TopicBuilder.name("so59256214.DLT").partitions(1).replicas(1).build();
    }

    @KafkaListener(id = "so59256214.DLT", topics = "so59256214.DLT")
    void listen(ConsumerRecord<?, ?> in) {
        System.out.println("dlt: " + in);
    }

    @Bean
    public ApplicationRunner runner(KafkaTemplate<String, Object> template) {
        return args -> {
            template.send("so59256214", 42);
            template.send("so59256214", 42.0);
            template.send("so59256214", "No handler for this");
        };
    }

    @Bean
    ErrorHandler eh(KafkaOperations<String, Object> template) {
        SeekToCurrentErrorHandler eh = new SeekToCurrentErrorHandler(new DeadLetterPublishingRecoverer(template));
        BackOff neverRetryOrBackOff = new FixedBackOff(0L, 0);
        BackOff normalBackOff = new FixedBackOff(2000L, 3);
        eh.setBackOffFunction((rec, ex) -> {
            if (ex.getMessage().contains("No method found for class")) {
                return neverRetryOrBackOff;
            }
            else {
                return normalBackOff;
            }
        });
        return eh;
    }

}

@Component
@KafkaListener(id = "so59256214", topics = "so59256214")
class Listener {

    @KafkaHandler
    void integerHandler(Integer in) {
        System.out.println("int: " + in);
    }

    @KafkaHandler
    void doubleHandler(Double in) {
        System.out.println("double: " + in);
    }

}
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer

结果:

int: 42
double: 42.0
dlt: ConsumerRecord(topic = so59256214.DLT, ...

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-02-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多