如何测试带自定义元素的扩展KafkaListener并实现死信消息发送?
解决方案
一、测试自定义@myListener注解的Kafka消费者
基于spring-kafka-test提供的EmbeddedKafka模拟Kafka环境,核心是确保自定义元注解@myListener被Spring正确识别生效,具体步骤如下:
编写基础测试用例
在测试类上添加@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); } }验证自定义注解属性解析
如果需要确认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属性,将异常消息转发到指定死信主题,具体实现:
自定义错误处理器
继承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); } }关联错误处理器到消费者容器
在配置类中将自定义错误处理器绑定到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; } }确保自定义元注解正确继承@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
相关产品推荐
相关产品推荐

