Testing Kafka producers and consumers in Java with Embedded Kafka
Apache Kafka sits at the heart of many Australian event-driven platforms, from the real-time price feeds at ASX to the transaction pipelines powering buy-now-pay-later services modelled on Afterpay. When engineers in Melbourne and Sydney write code against Kafka topics, they need fast feedback loops, and spinning up a full broker just to verify a producer record is a luxury nobody has time for during a sprint. Embedded Kafka offers a compact alternative that fits neatly inside a JUnit run.
The library, originally lifted from the Confluent test harness and now maintained as spring-kafka-test, lets you start a real broker in-process. Because the protocol bytes are identical to a production cluster, your producers and consumers exercise the same serializer, partitioner and offset logic they will use in staging. For Australian teams delivering on tight AEST deadlines or supporting squads across APAC, that parity removes a whole class of "works on my machine" surprises that used to plague Friday afternoon deployments.
Why Embedded Kafka fits Australian delivery cycles
A lot of local teams operate under the ACSC Essential Eight maturity expectations, where automated tests are part of the standard delivery pipeline. Embedded Kafka lets you honour those controls without a separate Docker daemon for every developer laptop, which is helpful when half your team is on a flight between Brisbane and Perth and the rest are working from home in Adelaide. The broker boots inside the JVM, so tests stay reproducible regardless of the host operating system.
There is also a practical cultural angle. Australian engineering culture tends to favour pragmatic, low-ceremony tooling — fair dinkum tools that get the job done without layers of YAML. Embedded Kafka matches that ethos: a couple of annotations on your test class and the broker is up before your @BeforeEach method even returns. You can focus on the message contract, not the cluster plumbing.
Performance-wise, the embedded broker is not designed for load testing, but for verifying serialisation, partition routing and consumer acknowledgement, it is more than quick enough. A typical producer test finishes in under a second, which keeps your CI feedback cycle tight even on a busy Monday morning when the Sydney office network is still warming up.
Wiring spring-kafka-test into a build
Most Java shops around Martin Place and South Bank already use Spring Boot, so pulling in the test artefact is the path of least resistance. In Maven you add spring-kafka-test as a test scope dependency, and Gradle users reach for testImplementation. The library transitively brings in kafka-clients, kafka_2.13 and a minimal log4j binding, so you do not need to declare Confluent platform artefacts separately.
Once the dependency is in place, you annotate your test class with @EmbeddedKafka. The annotation accepts a partitions count, a topics array and a brokerProperties map. For most Australian teams, three partitions per topic is a sensible default that mirrors what you would run in a non-production environment hosted on AWS ap-southeast-2. You can also pin a KRaft-only configuration to skip the older Zookeeper path, which keeps the test aligned with modern Confluent Platform defaults. Add @SpringBootTest if you want the full application context, or stay lighter with @ExtendWith(SpringExtension.class) for a slimmer unit-style test.
Configuration choices matter here. Australian banks and fintechs often require message headers to carry a trace id matching an internal standard, so you will want to keep the broker properties identical to your staging cluster. Pinning auto.create.topics.enable=false forces the test to fail loud if your production code tries to publish to a topic that does not exist — a small but valuable safety net before the code reaches the regulator-facing environment.
Verifying producer behaviour with EmbeddedKafkaBroker
A producer test usually wants to confirm three things: that a message lands on the expected topic, that the payload serialises correctly, and that headers or keys are attached the way downstream consumers expect. Embedded Kafka exposes an EmbeddedKafkaBroker bean that you can inject directly, then call consumeFromAnEmbeddedTopic with a ConsumerRecord collector or a one-record poll.
@EmbeddedKafka(partitions = 1, topics = {"orders.placed"})
class OrderProducerTest {
@Autowired EmbeddedKafkaBroker broker;
@Autowired OrderProducer producer;
@Test
void publishesOrderEventWithSchemaIdHeader() {
producer.send("orders.placed", "ORD-42", new OrderPlaced(...));
ConsumerRecord<String, byte[]> record = broker.consumeFromAnEmbeddedTopic(
"orders.placed", new StringDeserializer(), new ByteArrayDeserializer());
assertThat(record.key()).isEqualTo("ORD-42");
assertThat(record.headers().lastHeader("schema-id")).isNotNull();
}
}
The trick is to consume with the same Deserializer configuration the production consumer will use. If your real consumer expects Avro with the Confluent Schema Registry wire format, the test consumer should too. Otherwise the assertion will pass on a JSON byte string and silently break the moment an Avro decoder wraps the payload.
A common pattern inside Australian squads is to wrap the consumer poll in a small Awaitility helper that retries for up to five seconds. That gives Kafka time to flush the producer buffer and avoids a flaky test caused by the JVM warm-up of the embedded broker — important when a developer runs the suite straight after a coffee break in the Surry Hills office.
Confirming consumers read what producers wrote
Consumer tests are usually the harder half. You need to start the listener, push a record, and then assert that the application reacted correctly. With Embedded Kafka you have two complementary options. The first is to use @KafkaListener against the embedded broker, post a real ProducerRecord, and verify a downstream side effect such as a database row or a mocked external call. The second is to bypass the listener and assert directly on what the consumer would have received, using a manual KafkaConsumer configured with the same group.id and auto.offset.reset settings.
The first approach gives you the strongest confidence because it exercises deserialisation, error handling and acknowledgement in one go. Australian teams building against the Consumer Data Right (CDR) standards, for example, often need to prove that an event lands in a downstream ledger system within a defined service-level window. A listener-based test lets you simulate that path end-to-end and assert on the timestamp, all without ever touching a real cluster.
Watch out for the group.id collision. If your tests share a default group name across runs, the embedded broker can hand back records from a previous test, and you will spend half the morning chasing intermittent failures. The fix is to give every test class its own randomly generated group id, which the Spring test support makes easy via a UUID appended to a base name.
Embedded Kafka against the alternatives
The comparison below sketches how Embedded Kafka compares with the two most common alternatives for Australian Java teams: Testcontainers Kafka and lightweight mock libraries such as MockConsumer.
| Approach | Fidelity | Setup speed | Best for |
|---|---|---|---|
Embedded Kafka (spring-kafka-test) |
High — real protocol, real broker binary | Fast, in-JVM | Producer/consumer logic, serialisation, partition routing |
| Testcontainers Kafka | Highest — same image as staging | Slower, requires Docker | Full integration, Schema Registry, Connect, ksqlDB |
Mock consumer/producer (MockConsumer, Mockito) |
Low — bypasses broker | Fastest | Pure logic around the Kafka client API |
For most everyday producer and consumer tests, Embedded Kafka hits the sweet spot. Reach for Testcontainers when you need Schema Registry, Kafka Connect or a MirrorMaker scenario, and reach for mocks only when you are unit-testing an adapter that does not care about Kafka semantics.
Pitfalls you will meet on the way
The most common stumble is asserting on payload bytes without sharing the deserialiser config. Another is forgetting that the embedded broker binds to a random free port — your producers must resolve bootstrap servers through EmbeddedKafkaBroker.getBrokersAsString() rather than hard-coding localhost:9092. A third is treating the embedded broker as a load-testing tool; it is not, and pushing thousands of records per second will skew timings and obscure real regressions.
Australian teams also need to think about time zones. Kafka timestamps are stored in UTC, but if you assert on Instant or ZonedDateTime in your test, make sure the production code does the same. A test that passes in Sydney but fails for a colleague in Perth usually points at a missing ZoneOffset.UTC normalisation rather than a flaky broker.
Finally, keep your topic names in sync with the rest of the estate. A producer test passing locally but failing in the shared CI environment hosted in ap-southeast-2 often traces back to a topic name that does not yet exist on the real cluster, or one that has been renamed as part of an internal naming standardisation push.
Practical recommendations for your test suite
- Pin the same
spring-kafkaandkafka-clientsversions in test scope that you ship in production to avoid subtle API drift. - Set
auto.create.topics.enable=falseon the embedded broker so missing topics surface as fast feedback rather than silent bugs in staging. - Always consume from the embedded broker with the same deserialiser pair the production consumer uses, including the Schema Registry client when applicable.
- Generate a unique
group.idper test class to prevent cross-test bleed when the suite runs in parallel. - Use
Awaitilitywith a generous timeout for consumer assertions to absorb JVM warm-up on developer machines. - Reserve Testcontainers for the handful of tests that need Schema Registry or Connect, and keep the rest on Embedded Kafka for speed.
- Add a
@Tag("integration")marker so you can split fast unit tests from slower broker-backed tests in your CI matrix.
Remember that Embedded Kafka is a development convenience, not a replacement for end-to-end tests. The right balance is a fast inner loop on your laptop, a slightly slower Testcontainers run in CI, and a small suite of production-like smoke tests against the real cluster — that combination catches most regressions before they ever reach a customer in Melbourne, Brisbane or anywhere else your service operates.