使用Spring Embedded Kafka测试@KafkaListener

ric*_*din 12 java apache-kafka spring-boot spring-kafka spring-boot-test

我正在尝试为使用Spring Boot 2.x开发的Kafka监听器编写单元测试.作为一个单元测试,我不想启动一个完整的Kafka服务器作为Zookeeper的一个实例.所以,我决定使用Spring Embedded Kafka.

我的听众的定义非常基本.

@Component
public class Listener {
    private final CountDownLatch latch;

    @Autowired
    public Listener(CountDownLatch latch) {
        this.latch = latch;
    }

    @KafkaListener(topics = "sample-topic")
    public void listen(String message) {
        latch.countDown();
    }
}
Run Code Online (Sandbox Code Playgroud)

此外,latch在接收消息后验证计数器等于零的测试非常容易.

@RunWith(SpringRunner.class)
@SpringBootTest
@DirtiesContext
@EmbeddedKafka(topics = { "sample-topic" })
@TestPropertySource(properties = { "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}" })
public class ListenerTest {

    @Autowired
    private KafkaEmbedded embeddedKafka;

    @Autowired
    private CountDownLatch latch;

    private KafkaTemplate<Integer, String> producer;

    @Before
    public void setUp() {
        this.producer = buildKafkaTemplate();
        this.producer.setDefaultTopic("sample-topic");
    }

    private KafkaTemplate<Integer, String> buildKafkaTemplate() {
        Map<String, Object> senderProps = KafkaTestUtils.producerProps(embeddedKafka);
        ProducerFactory<Integer, String> pf = new DefaultKafkaProducerFactory<>(senderProps);
        return new KafkaTemplate<>(pf);
    }

    @Test
    public void listenerShouldConsumeMessages() throws InterruptedException {
        // Given
        producer.sendDefault(1, "Hello world");
        // Then
        assertThat(latch.await(10L, TimeUnit.SECONDS)).isTrue();
    }
}
Run Code Online (Sandbox Code Playgroud)

不幸的是,测试失败了,我无法理解为什么.是否可以使用实例KafkaEmbedded来测试标注注释的方法@KafkaListener?

所有代码都在我的GitHub存储库kafka-listener中共享.

谢谢大家.

Gar*_*ell 8

您可能在为消费者分配了主题/分区之前发送消息.设置属性......

spring:
  kafka:
    consumer:
      auto-offset-reset: earliest
Run Code Online (Sandbox Code Playgroud)

...它默认为latest.

这就像使用--from-beginning控制台消费者一样.

编辑

哦; 你没有使用boot的属性.

加

props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
Run Code Online (Sandbox Code Playgroud)

EDIT2

顺便说一句,你应该get(10L, TimeUnit.SECONDS)对template.send()(a Future<>)断言发送成功的结果做一个.

EDIT3

要仅为测试覆盖偏移重置,您可以执行与代理地址相同的操作:

@Value("${spring.kafka.consumer.auto-offset-reset:latest}")
private String reset;

...

    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, this.reset);
Run Code Online (Sandbox Code Playgroud)

和

@TestPropertySource(properties = { "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}",
        "spring.kafka.consumer.auto-offset-reset=earliest"})
Run Code Online (Sandbox Code Playgroud)

但是,请记住,此属性仅适用于组首次消耗.要在每次应用程序启动时始终从最后开始,您必须在启动期间寻求结束.

此外,我建议设置enable.auto.commit为false使容器负责提交补偿,而不是仅仅依靠消费者客户端按时间表执行.