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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 06:42:42