Kafka Testing Strategy — Spring EmbeddedKafka & Python testcontainers
-
Jason Yang - 26 Mar, 2026
- Views —
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 Test | How to Test |
|---|---|
| Business logic (event handlers) | Unit tests (pure functions, no mocks) |
| Producer / Consumer behavior | Integration tests (real Kafka required) |
| Serialization / Deserialization | Unit tests (Kafka not required) |
| Topic routing, offset commits | Integration 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-kafkadependency 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.
Related Posts
- Spring Boot Kafka Guide
- FastAPI Kafka Guide
- FastAPI Kafka Consumer
- Outbox Pattern