Logo Questions Linux Laravel Mysql Ubuntu Git Menu

How to create unit test with kafka embedded in the spring cloud stream

Sorry for the question being too generic, but someone has some tutorial or guide on how to perform producer and consumer testing with kafka embedded. I've tried several, but there are several versions of dependencies and none actually works =/

I'm using spring cloud stream kafka.

like image 268
Tiago Costa Avatar asked Apr 10 '17 18:04

Tiago Costa

People also ask

How do you write a unit test for Kafka consumer?

Testing a Kafka Consumer Consuming data from Kafka consists of two main steps. Firstly, we have to subscribe to topics or assign topic partitions manually. Secondly, we poll batches of records using the poll method. The polling is usually done in an infinite loop.

How do you do a Kafka test?

To test Kafka APIs, you use the API Connection test step. To add it to a test case, you will need a ReadyAPI Test Pro license. If you do not have it, try a ReadyAPI trial.

How do you write test cases for Kafka listener?

If you want to write integration tests using EmbeddedKafka , then you can do something like this. Assume we have some KafkaListener , which accepts a RequestDto as a Payload . In your test class you should create a TestConfiguration in order to create producer beans and to autowire KafkaTemplate into your test.

How do you test a Kafka Consumer Spring boot?

First, we start by decorating our test class with two pretty standard Spring annotations: The @SpringBootTest annotation will ensure that our test bootstraps the Spring application context. We also use the @DirtiesContext annotation, which will make sure this context is cleaned and reset between different tests.

1 Answers

We generally recommend using the Test Binder in tests but if you want to use an embedded kafka server, it can be done...

Add this to your POM...


Test app...

public class So43330544Application {

    public static void main(String[] args) {
        SpringApplication.run(So43330544Application.class, args);

    public byte[] handle(byte[] in){
        return new String(in).toUpperCase().getBytes();




Test case...

public class So43330544ApplicationTests {

    public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1);

    private KafkaTemplate<byte[], byte[]> template;

    private KafkaProperties properties;

    public static void setup() {
        System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString());

    public void testSendReceive() {
        template.send("so0544in", "foo".getBytes());
        Map<String, Object> configs = properties.buildConsumerProperties();
        configs.put(ConsumerConfig.GROUP_ID_CONFIG, "test0544");
        configs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        ConsumerFactory<byte[], byte[]> cf = new DefaultKafkaConsumerFactory<>(configs);
        Consumer<byte[], byte[]> consumer = cf.createConsumer();
        ConsumerRecords<byte[], byte[]> records = consumer.poll(10_000);
        assertThat(new String(records.iterator().next().value())).isEqualTo("FOO");

like image 53
Gary Russell Avatar answered Jan 04 '23 11:01

Gary Russell