Software Engineer's Blog

Kafka Testing Strategy — Spring EmbeddedKafka & Python testcontainers

Kafka Testing Strategy — Spring EmbeddedKafka & Python testcontainers

This post focuses on testing in the Kafka series.
It is recommended to read Spring Boot Implementation and FastAPI Implementation first.

One of the most common mistakes when testing Kafka-related code is mocking Kafka itself.

# ❌ Do not do this
mock_producer = MagicMock()
mock_producer.send.return_value = AsyncMock()

If you mock Kafka, issues like serialization errors, topic misconfiguration, and offset commit bugs will never surface in tests.
To properly validate real-world behavior, you need to run tests against an actual Kafka instance.

1. Testing Strategy Principles

What to TestHow to Test
Business logic (event handlers)Unit tests (pure functions, no mocks)
Producer / Consumer behaviorIntegration tests (real Kafka required)
Serialization / DeserializationUnit tests (Kafka not required)
Topic routing, offset commitsIntegration tests (real Kafka required)

2. Spring Boot — @EmbeddedKafka

Spring Kafka provides an embedded Kafka broker for testing.
It runs a real Kafka broker inside the JVM, without needing Docker.

EmbeddedKafka vs testcontainers — which one? EmbeddedKafka is fast and needs no Docker, but it runs Kafka’s own libraries in-process, so its version tracks the spring-kafka dependency rather than the broker you run in production. testcontainers (below) starts the actual Kafka image — closer to prod, at the cost of Docker and a slower first run. Rule of thumb: EmbeddedKafka for quick JVM-side tests, testcontainers when version parity with production matters or when the service isn’t on the JVM (like the Python examples here).

Add Dependency (pom.xml)

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka-test</artifactId>
    <scope>test</scope>
</dependency>

Producer Integration Test

@SpringBootTest
@EmbeddedKafka(
    partitions = 1,
    topics = {
        MessagingConstants.TOPIC_MY_EVENT_CREATED,
        MessagingConstants.TOPIC_MY_EVENT_DELETED,
    }
)
@TestPropertySource(properties = {
    "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}",
    "spring.kafka.consumer.auto-offset-reset=earliest",
})
class KafkaMyEventPublisherTest {

    @Autowired
    private KafkaMyEventPublisher publisher;

    @Autowired
    private EmbeddedKafkaBroker embeddedKafka;

    @Test
    void publishEvent_success_publishes_message_to_topic() throws Exception {
        // given
        String userId = "user-123";
        MyEventPayload payload = new MyEventPayload("ITEM_CREATED", "item-456");

        // when
        boolean result = publisher.publishEvent("MY_EVENT_CREATED", userId, payload);

        // then
        assertThat(result).isTrue();

        // Consume from real Kafka to verify
        ConsumerRecord<String, String> record = KafkaTestUtils.getSingleRecord(
            createTestConsumer(MessagingConstants.TOPIC_MY_EVENT_CREATED),
            MessagingConstants.TOPIC_MY_EVENT_CREATED
        );

        assertThat(record.key()).isEqualTo(userId);
        assertThat(record.value()).contains("ITEM_CREATED");
    }

    private Consumer<String, String> createTestConsumer(String topic) {
        Map<String, Object> consumerProps = KafkaTestUtils.consumerProps(
            "test-group", "true", embeddedKafka
        );
        Consumer<String, String> consumer =
            new DefaultKafkaConsumerFactory<String, String>(consumerProps).createConsumer();
        consumer.subscribe(Collections.singletonList(topic));
        return consumer;
    }
}

Consumer Integration Test

@SpringBootTest
@EmbeddedKafka(partitions = 1, topics = {MessagingConstants.TOPIC_MY_EVENT_CREATED})
@TestPropertySource(properties = {
    "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}",
    "spring.kafka.consumer.auto-offset-reset=earliest",
})
class MyEventKafkaConsumerTest {

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @MockBean
    private MyBusinessService myBusinessService;

    @Test
    void message_received_triggers_business_service() throws Exception {
        // given
        String payload = """
            {"eventType":"MY_EVENT_CREATED","userId":"user-123","idempotencyKey":"uuid-1"}
            """;

        // when: publish test message
        kafkaTemplate.send(MessagingConstants.TOPIC_MY_EVENT_CREATED, "user-123", payload);

        // then: wait until consumer processes it (max 5 seconds)
        verify(myBusinessService, timeout(5000).times(1))
            .processCreatedEvent(eq("user-123"), any());
    }

    @Test
    void failure_in_business_logic_routes_to_dlt() {
        // given
        doThrow(new RuntimeException("DB error"))
            .when(myBusinessService).processCreatedEvent(any(), any());

        String payload = """
            {"eventType":"MY_EVENT_CREATED","userId":"user-456","idempotencyKey":"uuid-2"}
            """;

        // when
        kafkaTemplate.send(MessagingConstants.TOPIC_MY_EVENT_CREATED, "user-456", payload);

        // then: verify message is routed to DLT
        ConsumerRecord<String, String> dltRecord = KafkaTestUtils.getSingleRecord(
            createTestConsumer(MessagingConstants.TOPIC_MY_EVENT_DLT),
            MessagingConstants.TOPIC_MY_EVENT_DLT,
            15000
        );
        assertThat(dltRecord).isNotNull();
    }
}

Important Notes

// ❌ WRONG - shared broker across tests can cause flaky tests
@EmbeddedKafka

// ✅ CORRECT - isolate with unique group-id per test
@TestPropertySource(properties = {
    "spring.kafka.consumer.group-id=test-group-${random.uuid}"
})

3. Python — testcontainers

In Python, we use the testcontainers library to spin up a real Kafka container during tests.

Add Dependency

[tool.poetry.group.test.dependencies]
testcontainers = {version = "^4.0", extras = ["kafka"]}
pytest-asyncio = "^0.23"

Kafka Container Fixture

# tests/conftest.py
import pytest
import pytest_asyncio
from testcontainers.kafka import KafkaContainer
from my_service.core.config import KafkaSettings


@pytest.fixture(scope="session")
def kafka_container():
    """Kafka container shared across the test session."""
    with KafkaContainer("apache/kafka:latest") as kafka:
        yield kafka


@pytest.fixture
def kafka_settings(kafka_container) -> KafkaSettings:
    """KafkaSettings for tests (uses real container address)."""
    return KafkaSettings(
        bootstrap_servers=kafka_container.get_bootstrap_server(),
        security_protocol="PLAINTEXT",
    )

Producer Integration Test

import json
import pytest
import pytest_asyncio
from aiokafka import AIOKafkaConsumer
from my_service.services.kafka_producer import KafkaProducer
from my_service.services.kafka_publisher import KafkaMyPublisher
from my_service.schemas.events import MyEventType, MyItemEvent


@pytest_asyncio.fixture
async def publisher(kafka_settings):
    producer = KafkaProducer(kafka_settings)
    await producer.start()
    pub = KafkaMyPublisher(producer)
    yield pub
    await producer.stop()


@pytest.mark.asyncio
async def test_publish_item_event_sends_message_to_correct_topic(
    publisher, kafka_settings
):
    # given
    event = MyItemEvent(
        event_type=MyEventType.ITEM_CREATED,
        user_id="user-123",
        item_id="item-456",
        item_type="document",
    )

    # when
    await publisher.publish_item_event(event)

    # then: verify by consuming from real Kafka
    consumer = AIOKafkaConsumer(
        "myapp.my.item.created",
        bootstrap_servers=kafka_settings.bootstrap_servers,
        group_id="test-verify-group",
        auto_offset_reset="earliest",
    )
    await consumer.start()
    try:
        msg = await asyncio.wait_for(consumer.__anext__(), timeout=5.0)
        payload = json.loads(msg.value)
        assert payload["item_id"] == "item-456"
        assert payload["event_type"] == "ITEM_CREATED"
        assert msg.key == b"user-123"
    finally:
        await consumer.stop()

Consumer Integration Test

import asyncio
import json
import pytest
from unittest.mock import AsyncMock
from aiokafka import AIOKafkaProducer
from my_service.services.kafka_consumer import KafkaConsumer
from my_service.core.messaging_constants import MessagingConstants


@pytest.mark.asyncio
async def test_consumer_calls_handler_on_message(kafka_settings):
    mock_handler = AsyncMock()

    consumer = KafkaConsumer(
        settings=kafka_settings,
        topics=[MessagingConstants.TOPIC_MY_ITEM_CREATED],
        group_id="test-consumer-group",
        handler=mock_handler,
    )
    await consumer.start()

    producer = AIOKafkaProducer(
        bootstrap_servers=kafka_settings.bootstrap_servers
    )
    await producer.start()
    try:
        payload = json.dumps({"event_type": "ITEM_CREATED", "item_id": "item-1"}).encode()
        await producer.send_and_wait(
            MessagingConstants.TOPIC_MY_ITEM_CREATED,
            value=payload,
            key=b"user-1",
        )
    finally:
        await producer.stop()

    await asyncio.wait_for(_wait_for_mock_call(mock_handler), timeout=5.0)
    mock_handler.assert_called_once()

    await consumer.stop()


async def _wait_for_mock_call(mock, interval=0.1):
    while not mock.called:
        await asyncio.sleep(interval)

Serialization Unit Test (Kafka not required)

from my_service.schemas.events import MyEventType, MyItemEvent
import json


def test_my_item_event_serialization():
    event = MyItemEvent(
        event_type=MyEventType.ITEM_CREATED,
        user_id="user-123",
        item_id="item-456",
        item_type="document",
    )

    json_str = event.model_dump_json()
    payload = json.loads(json_str)

    assert payload["event_type"] == "ITEM_CREATED"
    assert payload["user_id"] == "user-123"
    assert "event_id" in payload
    assert "occurred_at" in payload


def test_my_item_event_deserialization():
    payload = {
        "event_id": "some-uuid",
        "event_version": "1.0",
        "occurred_at": "2024-01-01T00:00:00Z",
        "source_service": "my-service",
        "event_type": "ITEM_CREATED",
        "user_id": "user-123",
        "item_id": "item-456",
        "item_type": "document",
    }

    event = MyItemEvent.model_validate(payload)
    assert event.item_id == "item-456"

4. Common Mistakes

The Risk of Mocking Kafka

# ❌ This test does NOT verify:
# - JSON serialization correctness
# - partition key configuration
# - topic name accuracy
mock_producer = MagicMock()
await publisher.publish_item_event(event)
mock_producer.send.assert_called_once()
# ✅ Using real Kafka verifies everything
await publisher.publish_item_event(event)
msg = await consumer.getone()
assert json.loads(msg.value)["item_id"] == event.item_id

Test Isolation

# ❌ shared offsets across tests
group_id = "test-group"

# ✅ isolate each test with unique group-id
import uuid
group_id = f"test-group-{uuid.uuid4()}"

5. Test Performance Optimization

testcontainers can be slow on the first run because it downloads Docker images.
Using scope="session" allows the container to be reused across the test session, which significantly improves performance.

@pytest.fixture(scope="session")
def kafka_container():
    with KafkaContainer("apache/kafka:latest") as kafka:
        yield kafka

In CI environments, you can also use Docker-in-Docker or a dedicated Kafka server.