Command Palette

Search for a command to run...

Hectal
PHASE 7Intermediate ~7 min· topic 1 of 4

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.

0/4 · 0%

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

  1. 01

    spring.kafka.producer.* properties map to producer configs; KafkaTemplate.send(topic, key, value) returns a CompletableFuture<SendResult> in Spring Kafka 3.

  2. 02

    JSON: JsonSerializer with type headers disabled (spring.json.add.type.headers=false) for cross-language consumers, or Avro/Protobuf with Schema Registry (Phase 8).

  3. 03

    Headers: add event type, schema version and correlation ID. Tracing: Micrometer observation (template.setObservationEnabled(true)) propagates trace context in headers.

  4. 04

    Don't block request threads on send().get() unless you must return the result to the caller; prefer whenComplete callbacks and the outbox pattern for must-not-lose events.

  5. 05

    Other languages: librdkafka-based clients (confluent-kafka-go, confluent-kafka-python) use linger.ms/batch.size names too, but their default partitioner is consistent_random (CRC32), not Java's murmur2: set partitioner=murmur2_random if keys must land in the same partitions as Java producers.

Code & diagrams

application.ymlyaml
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
OrderEventPublisher.javajava
@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

01

Why disable Spring's JSON type headers for shared topics?

Practice

01

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.