Quarkus中无需@Incoming注解消费Kafka消息的测试方案咨询
在Quarkus测试中不使用@Incoming消费Kafka消息的几种方案
如果你不想依赖@Incoming注解,完全可以通过以下几种方式手动消费Kafka消息,同时复用application.yaml中配置的SmallRye Kafka连接器参数:
1. 直接使用原生Kafka消费者(结合SmallRye配置加载)
利用SmallRye的KafkaConsumerConfigurer自动读取配置文件中的Kafka参数,再创建原生KafkaConsumer实例手动消费:
import io.smallrye.reactive.messaging.kafka.KafkaConsumerConfigurer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import jakarta.inject.Inject; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class KafkaPublishTest { @Inject KafkaConsumerConfigurer configurer; public void verifyPublishedMessage() { Properties consumerProps = new Properties(); // 加载application.yaml中指定连接器的配置 configurer.configure("your-kafka-connector-name", consumerProps); // 覆盖测试所需的特定配置 consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "test-publish-verifier-group"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 自动关闭消费者 try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps)) { consumer.subscribe(Collections.singletonList("your-target-topic")); // 轮询获取消息,设置合理超时时间适配测试场景 var records = consumer.poll(Duration.ofSeconds(3)); records.forEach(record -> { // 在这里断言消息内容是否符合预期 System.out.println("Received test message: " + record.value()); }); // 手动提交偏移量(按需选择) consumer.commitSync(); } } }
2. 通过KafkaConsumerRegistry获取已配置的消费者
如果你的应用已经在配置文件中定义了消费者连接器,可以直接通过KafkaConsumerRegistry获取预配置的消费者实例:
import io.smallrye.reactive.messaging.kafka.KafkaConsumerRegistry; import org.apache.kafka.clients.consumer.KafkaConsumer; import jakarta.inject.Inject; import java.time.Duration; import java.util.Collections; public class KafkaPublishTest { @Inject KafkaConsumerRegistry consumerRegistry; public void checkPublishedMessages() { // 通过连接器名称获取已配置的消费者 KafkaConsumer<String, String> consumer = consumerRegistry.getConsumer("your-kafka-connector-name"); consumer.subscribe(Collections.singletonList("your-target-topic")); var records = consumer.poll(Duration.ofSeconds(3)); // 处理并验证消息 records.forEach(record -> System.out.println("Test message received: " + record.value())); consumer.commitSync(); } }
3. 使用MicroProfile Reactive Messaging的MessagingProvider构建消费者通道
通过MessagingProvider手动创建消费者订阅,以响应式方式处理消息:
import org.eclipse.microprofile.reactive.messaging.Message; import org.eclipse.microprofile.reactive.messaging.spi.MessagingProvider; import jakarta.inject.Inject; import java.util.concurrent.TimeUnit; public class KafkaPublishTest { @Inject MessagingProvider messagingProvider; public void validatePublishedMessage() throws InterruptedException { // 构建消费者订阅器 var publisher = messagingProvider.createSubscriberBuilder() .fromChannel("your-target-topic") .buildRs(); // 订阅消息并处理 publisher.subscribe(new java.util.concurrent.Subscriber<>() { @Override public void onSubscribe(java.util.concurrent.Subscription subscription) { subscription.request(5); // 请求指定数量的消息 } @Override public void onNext(Message<String> message) { // 断言消息内容 System.out.println("Received message payload: " + message.getPayload()); message.ack(); // 确认消息处理完成 } @Override public void onError(Throwable throwable) { throwable.printStackTrace(); } @Override public void onComplete() { System.out.println("Test consumption completed"); } }); // 等待消息处理完成 TimeUnit.SECONDS.sleep(3); publisher.close(); } }
关键注意事项
- 测试类需要添加
@QuarkusTest注解,确保Quarkus加载application.yaml中的Kafka配置 - 测试用的消费者组ID建议单独设置,避免与生产/其他测试的消费者组冲突
- 根据实际消息的序列化类型,调整
KafkaConsumer的泛型参数和反序列化器配置(比如用JSON反序列化器处理结构化消息) - 轮询超时时间需要根据测试环境的Kafka响应速度合理设置,避免测试超时或漏收消息
内容的提问来源于stack exchange,提问作者Bull
相关产品推荐
相关产品推荐

