You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何测试带自定义元素的扩展KafkaListener并实现死信消息发送?

解决方案

一、测试自定义@myListener注解的Kafka消费者

基于spring-kafka-test提供的EmbeddedKafka模拟Kafka环境,核心是确保自定义元注解@myListener被Spring正确识别生效,具体步骤如下:

  1. 编写基础测试用例
    在测试类上添加@SpringBootTest和@EmbeddedKafka注解,指定测试Topic,通过KafkaTemplate发送消息,用CountDownLatch等待消费完成后验证结果:

    @SpringBootTest
    @EmbeddedKafka(partitions = 1, topics = {"test-topic", "dlq-topic"})
    public class MyListenerTest {
    
        @Autowired
        private KafkaTemplate<String, String> kafkaTemplate;
    
        private CountDownLatch latch = new CountDownLatch(1);
        private String receivedMsg;
    
        // 模拟业务消费者的自定义监听方法
        @myListener(topics = "test-topic", myattr = "dlq-topic")
        public void consume(String message) {
            receivedMsg = message;
            latch.countDown();
        }
    
        @Test
        public void testNormalConsume() throws InterruptedException {
            // 发送测试消息
            kafkaTemplate.send("test-topic", "hello-test");
            // 等待消费完成
            latch.await(5, TimeUnit.SECONDS);
            // 断言消息被正确接收
            assertEquals("hello-test", receivedMsg);
        }
    }
    
  2. 验证自定义注解属性解析
    如果需要确认myattr属性是否被正确读取,可以通过KafkaListenerEndpointRegistry获取注册的监听端点,提取注解信息:

    @Autowired
    private KafkaListenerEndpointRegistry registry;
    
    @Test
    public void testMyListenerAttr() {
        // 对应@myListener的id属性值
        MessageListenerContainer container = registry.getListenerContainer("my-test-listener");
        MethodKafkaListenerEndpoint<?, ?> endpoint = (MethodKafkaListenerEndpoint<?, ?>) container.getEndpoint();
        MyListener myListener = endpoint.getMethod().getAnnotation(MyListener.class);
        assertEquals("dlq-topic", myListener.myattr());
    }
    

二、消费异常时从@myListener的myattr获取死信主题并发送

通过自定义Kafka错误处理器,在消费异常时提取@myListener的myattr属性,将异常消息转发到指定死信主题,具体实现:

  1. 自定义错误处理器
    继承SeekToCurrentErrorHandler,重写handle方法,从监听方法中获取@myListener注解的死信主题,完成消息转发:

    @Component
    public class MyDlqErrorHandler extends SeekToCurrentErrorHandler {
    
        private final KafkaTemplate<String, Object> kafkaTemplate;
    
        public MyDlqErrorHandler(KafkaTemplate<String, Object> kafkaTemplate) {
            // 保留父类的重试逻辑
            super(new DeadLetterPublishingRecoverer(kafkaTemplate), new FixedBackOff(1000L, 2L));
            this.kafkaTemplate = kafkaTemplate;
        }
    
        @Override
        public void handle(Exception thrownException, List<ConsumerRecord<?, ?>> records, Consumer<?, ?> consumer, MessageListenerContainer container) {
            // 获取当前监听方法的@myListener注解
            MethodKafkaListenerEndpoint<?, ?> endpoint = (MethodKafkaListenerEndpoint<?, ?>) container.getEndpoint();
            MyListener myListener = endpoint.getMethod().getAnnotation(MyListener.class);
            String dlqTopic = myListener.myattr();
    
            // 将异常消息发送到死信主题
            records.forEach(record -> {
                kafkaTemplate.send(dlqTopic, record.key(), record.value());
            });
    
            // 调用父类逻辑完成seek,避免重复消费
            super.handle(thrownException, records, consumer, container);
        }
    }
    
  2. 关联错误处理器到消费者容器
    在配置类中将自定义错误处理器绑定到Kafka监听容器工厂:

    @Configuration
    public class KafkaConfig {
    
        @Bean
        public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
                ConsumerFactory<Object, Object> consumerFactory,
                MyDlqErrorHandler myDlqErrorHandler) {
            ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
            factory.setConsumerFactory(consumerFactory);
            // 设置自定义错误处理器
            factory.setErrorHandler(myDlqErrorHandler);
            return factory;
        }
    }
    
  3. 确保自定义元注解正确继承@KafkaListener
    你的@myListener需要正确标注元注解继承关系:

    @Target({ElementType.METHOD, ElementType.ANNOTATION_TYPE})
    @Retention(RetentionPolicy.RUNTIME)
    @KafkaListener
    public @interface myListener {
        String[] topics() default {};
        String myattr() default "";
        // 可按需继承@KafkaListener的其他属性,如groupId、id等
        String groupId() default "";
        String id() default "";
    }
    

内容的提问来源于stack exchange,提问作者Learner

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 20:05:32