Learn Labs
3. Kafka Producers: Writing Messages to Kafka

3.3 The three ways to send

Everything that goes wrong at or after the broker is invisible.

Fire-and-forgetSynchronousAsynchronous + callback
producer.send(rec) — ignore the Futureproducer.send(rec).get() — block on the Futureproducer.send(rec, cb) — callback on response
Fastest. Silent message loss.Slowest by far. Poor performance.Fast and handles errors — ← the production choice

3.1 Fire-and-forget

ProducerRecord<String, String> record =
    new ProducerRecord<>("CustomerCountry", "Precision Products", "France");
try {
    producer.send(record);
} catch (Exception e) {
    e.printStackTrace();
}

"Most of the time it will arrive successfully, since Kafka is highly available and the producer will retry sending messages automatically. However, in case of nonretriable errors or timeout, messages will get lost and the application will not get any information or exceptions about this."

Subtle but important: even here you can still get an exception — but only for failures before the message reached the broker path:

  • SerializationException — failed to serialize
  • BufferExhaustedException / TimeoutException — the buffer is full
  • InterruptException — the sending thread was interrupted

Everything that goes wrong at or after the broker is invisible. "This method of sending messages can be used when dropping a message silently is acceptable. This is not typically the case in production applications."

3.2 Synchronous send

producer.send(record).get();      // ← the .get() is the whole difference

What it buys you: catches exceptions when Kafka responds to the produce request with an error, or when send retries were exhausted.

What it costs — and the numbers are worth memorizing:

"Depending on how busy the Kafka cluster is, brokers can take anywhere from 2 ms to a few seconds to respond to produce requests. If you send messages synchronously, the sending thread will spend this time waiting and doing nothing else, not even sending additional messages. This leads to very poor performance, and as a result, synchronous sends are usually not used in production applications (but are very common in code examples)."

That parenthetical is a warning about copy-paste: the pattern you see in every tutorial is the one you must not ship.

3.3 Retriable vs non-retriable errors

KafkaProducer has two types of errors:

RETRIABLE — resolvable by sending again:

  • connection error → the connection may get reestablished
  • “not leader for partition” → resolved when a new leader is elected and the client metadata is refreshed

NON-RETRIABLE — retrying cannot help:

  • “Message size too large” → KafkaProducer will not attempt a retry and returns the exception immediately

KafkaProducer can be configured to retry retriable errors automatically, "so the application code will get retriable exceptions only when the number of retries was exhausted and the error was not resolved."

3.4 Asynchronous send with a callback — the production pattern

The arithmetic that makes the case:

"Suppose the network round-trip time between our application and the Kafka cluster is 10 ms. If we wait for a reply after sending each message, sending 100 messages will take around 1 second. On the other hand, if we just send all our messages and not wait for any replies, then sending 100 messages will barely take any time at all."

"In most cases, we really don't need a reply — Kafka sends back the topic, partition, and offset of the record after it was written, which is usually not required by the sending app. On the other hand, we do need to know when we failed to send a message completely so we can throw an exception, log an error, or perhaps write the message to an 'errors' file for later analysis."

private class DemoProducerCallback implements Callback {
    @Override
    public void onCompletion(RecordMetadata recordMetadata, Exception e) {
        if (e != null) {
            e.printStackTrace();     // production: real error handling here
        }
    }
}

ProducerRecord<String, String> record =
    new ProducerRecord<>("CustomerCountry", "Biomedical Materials", "USA");

producer.send(record, new DemoProducerCallback());

Implement org.apache.kafka.clients.producer.Callback — a single method, onCompletion(). If Kafka returned an error, the exception argument is non-null.

⚠️ WARNING — the callback threading rule

Callbacks execute in the producer's main thread.

Benefit: guarantees that when you send two messages to the same partition one after another, their callbacks execute in the same order you sent them.

Cost: "the callback should be reasonably fast to avoid delaying the producer and preventing other messages from being sent. It is not recommended to perform a blocking operation within the callback. Instead, you should use another thread to perform any blocking operation concurrently."

This is a top-tier production landmine. A callback that writes failures to a database, calls an HTTP endpoint, or logs synchronously to a slow appender will throttle your entire producer. Symptom: producer throughput collapses and nobody can see why, because the code "looks async."


On this page