Topic 7.1
A Production Spring Boot Producer
In one line
Configure a KafkaTemplate for JSON events with a key per entity, headers for tracing and event type, acks=all, idempotence, zstd compression and batching, and handle send results asynchronously so failures are logged, counted and never silently dropped.
Think of it like this
A well-run shipping desk. Every parcel gets an address label (key), a tracking sticker (headers), is packed efficiently (batching and compression), and the desk records every delivery confirmation or failure (async result handling).
Key ideas
- 01
spring.kafka.producer.*properties map to producer configs;KafkaTemplate.send(topic, key, value)returns aCompletableFuture<SendResult>in Spring Kafka 3. - 02
JSON:
JsonSerializerwith type headers disabled (spring.json.add.type.headers=false) for cross-language consumers, or Avro/Protobuf with Schema Registry (Phase 8). - 03
Headers: add event type, schema version and correlation ID. Tracing: Micrometer observation (
template.setObservationEnabled(true)) propagates trace context in headers. - 04
Don't block request threads on
send().get()unless you must return the result to the caller; preferwhenCompletecallbacks and the outbox pattern for must-not-lose events. - 05
Other languages: librdkafka-based clients (confluent-kafka-go, confluent-kafka-python) use
linger.ms/batch.sizenames too, but their default partitioner isconsistent_random(CRC32), not Java's murmur2: setpartitioner=murmur2_randomif keys must land in the same partitions as Java producers.
Code & diagrams
spring:
kafka:
bootstrap-servers: kafka1:9092,kafka2:9092,kafka3:9092
producer:
client-id: order-service
acks: all
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
compression-type: zstd
batch-size: 65536
properties:
enable.idempotence: true
linger.ms: 10
delivery.timeout.ms: 120000
spring.json.add.type.headers: false@Component
@RequiredArgsConstructor
public class OrderEventPublisher {
private final KafkaTemplate<String, OrderEvent> template;
private final MeterRegistry metrics;
public CompletableFuture<SendResult<String, OrderEvent>> publish(OrderEvent evt) {
ProducerRecord<String, OrderEvent> record = new ProducerRecord<>("commerce.orders", evt.orderId(), evt);
record.headers()
.add("eventType", evt.type().name().getBytes(UTF_8))
.add("schemaVersion", "2".getBytes(UTF_8))
.add("eventId", evt.eventId().toString().getBytes(UTF_8));
return template.send(record).whenComplete((res, ex) -> {
if (ex != null) {
metrics.counter("orders.publish.failed").increment();
log.error("publish failed order={} eventId={}", evt.orderId(), evt.eventId(), ex);
} else {
RecordMetadata m = res.getRecordMetadata();
log.debug("published order={} partition={} offset={}", evt.orderId(), m.partition(), m.offset());
}
});
}
}When it breaks
Python and Java producers write the same keyed topic with default partitioners
What you see
The same key lands in different partitions depending on which service produced it; per-key ordering breaks.
Fix & prevent
Configure librdkafka clients with partitioner=murmur2_random (Java-compatible) or route all writes for a key through one service.
Explain it without notes
Why disable Spring's JSON type headers for shared topics?
Practice
Implement the publisher, send 1,000 events for 50 orders, and verify each order's events share a partition in order.
Trade-offs
- ↔
Async sends maximise throughput; synchronous confirmation simplifies API semantics at a latency cost.
Done when you can
I can configure and write a production Spring Kafka producer with keys, headers and failure handling.