Spring for Apache Kafka

This tutorial demonstrates how to use Spring for Apache Kafka to publish and listen for messages in a Spring Boot application, and how to write integration tests using Testcontainers.

The Use Case

We will build a simple application that takes a user payload via a REST API, publishes it to a Kafka topic, and then uses a listener to read the message from that same topic.

Project Structure

The project has three main components: - User: A record representing our domain model. - UserResource: A REST controller that accepts POST requests and publishes the user to Kafka. - UserListener: A Kafka listener that listens to messages on the user-events topic.

The Record

We start by defining a User record to hold our data. Using records ensures that Jackson can automatically serialize it into JSON without any explicit annotations.

public record User(Long id, String name, String username, String email) {
}

The Controller

Our UserResource will accept a POST request with the user JSON and then publish it using KafkaTemplate.

@RestController
class UserController {

    private final KafkaTemplate<String, User> kafkaTemplate;

    UserController(KafkaTemplate<String, User> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    @PostMapping("/users")
    @ResponseStatus(HttpStatus.ACCEPTED)
    void publish(@RequestBody User user) {
        kafkaTemplate.send("user-events", String.valueOf(user.id()), user);
    }
}

The Listener

Our UserListener uses the @KafkaListener annotation to listen to the user-events topic.

@Component
class UserListener {

    private final List<User> receivedUsers = Collections.synchronizedList(new ArrayList<>());

    @KafkaListener(topics = "user-events", groupId = "user-group")
    void listen(User user) {
        receivedUsers.add(user);
    }

    public List<User> getReceivedUsers() {
        return List.copyOf(receivedUsers);
    }
}

Testing

We will write an integration test using Testcontainers to bring up an actual Kafka container. The setup for Kafka can be found in TestcontainersConfiguration.

@TestConfiguration(proxyBeanMethods = false)
public class TestcontainersConfiguration {

    @Bean
    @ServiceConnection
    KafkaContainer kafkaContainer() {
        return new KafkaContainer(DockerImageName.parse("apache/kafka-native:latest"));
    }

}

The @ServiceConnection annotation automatically configures the broker properties for us, so we don’t need to specify the connection string manually in our properties file.

@Import(TestcontainersConfiguration.class)
@AutoConfigureRestTestClient
@SpringBootTest(webEnvironment =  RANDOM_PORT)
class PublishUserTests {

    @Autowired
    private RestTestClient client;

    @Autowired
    private UserListener userListener;

    @BeforeEach
    void setUp() {
        userListener.getReceivedUsers().clear();
    }

    @Test
    @DisplayName("When user is published via REST, it should be received by the Kafka listener")
    void publish() {
        var user = new User(1L, "Rashidi Zin", "rashidi.zin", "[email protected]");

        client.post().uri("/users")
                .contentType(MediaType.APPLICATION_JSON)
                .body(user)
                .exchange()
                .expectStatus().isAccepted();

        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() ->
            assertThat(userListener.getReceivedUsers())
                    .hasSize(1)
                    .extracting("name")
                    .containsExactly("Rashidi Zin")
        );
    }
}

This setup ensures our tests are robust and use real infrastructure instead of mocks.