skip to content

Explain the producer's send() async model and the roles of flush() and close(). What can go wrong if you skip them?

level: seniorimportance: must knowfreq 65%

answer

  1. send() = buffer + return Future, async
  2. Sender thread does the real send
  3. flush() = block till buffered done, no close
  4. close() = implicit flush + release
  5. skip -> buffered records lost on exit

basics

~20 s

send() is asynchronous: it buffers the record and returns a Future immediately; a background Sender thread actually transmits batches. flush() blocks until all buffered records have been sent and acknowledged. close() flushes then releases resources. Skip them and you can lose unsent buffered records on exit.

solid answer

~50 s

KafkaProducer.send() doesn't send synchronously — it serializes the record, appends it to an in-memory buffer batched per partition, and returns a Future<RecordMetadata> immediately. A background Sender thread drains batches to brokers based on batch.size, linger.ms, and buffer.memory. So a record is durable only once its callback/Future completes successfully. flush() blocks the calling thread until every record buffered so far has completed (sent and acknowledged per acks), without closing the producer — useful at a checkpoint. close() does an implicit flush (waits for in-flight records) then shuts down the Sender thread and releases connections; close(Duration) bounds the wait and will drop still-pending records if the timeout elapses. The danger: if you let the JVM exit (or close() with a tiny timeout) while records sit in the buffer, those records are silently lost — send() returning is NOT a delivery guarantee.

code

java · 7 lines
java
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
    for (Order o : orders) {
        producer.send(new ProducerRecord<>("orders", o.id(), o.json()),
            (md, ex) -> { if (ex != null) log.error("send failed", ex); });
    }
    producer.flush(); // durability checkpoint: block until all complete
} // try-with-resources -> close() does an implicit flush + releases resources

go deeper

for a junior

Know send() is async (returns a Future) and that you should flush/close before exiting.

for a middle

Explain the Sender thread, that close() implies flush, and that buffered records are lost if you skip them.

for a senior

Tie delivery to acks + callback/Future, discuss close(Duration) trade-offs and linger/batch buffering.

for a principal

Define org-wide producer lifecycle and durability contracts (shared producer beans, mandatory result-checking, bounded close on shutdown, idempotence/acks=all).

## send() is asynchronous A frequent misconception is that `producer.send(record)` transmits the record. It doesn't. `send()`: 1. serializes key/value, 2. computes the partition (via the partitioner), 3. appends the record to an in-memory `RecordAccumulator` batch for that partition, 4. returns a `Future<RecordMetadata>` **immediately**. A separate **Sender** background thread drains accumulated batches and sends them to brokers. So after `send()` returns, the record is merely *buffered* — not yet on any broker, not yet durable. ## When is a record actually delivered? Only when its result completes successfully — either the `Future.get()` returns, or the optional callback `send(record, (metadata, exception) -> ...)` fires with `exception == null`. Delivery is governed by: - **`acks`** (0 / 1 / all): how many replicas must acknowledge. - **`batch.size`** and **`linger.ms`**: how batches fill / how long the Sender waits to batch before sending (linger trades a little latency for throughput). - **`buffer.memory`**: total buffer; if it fills, `send()` blocks up to `max.block.ms` then throws. - **`retries`/`delivery.timeout.ms`**: how long the client retries before the Future fails. ## flush() `producer.flush()` **blocks the caller until all records buffered up to this point have completed** (succeeded or failed) — it forces the Sender to send everything now and waits. It does **not** close the producer; you keep using it after. Use it at a logical checkpoint (e.g. after a batch import, before reporting success) when you need a hard 'everything is durable now' barrier without losing the producer. ## close() `producer.close()`: 1. performs an **implicit flush** — waits for in-flight and buffered records to complete, 2. stops the Sender thread, 3. closes network connections and frees the buffer. `close(Duration timeout)` bounds step 1: if the timeout elapses with records still pending, those records are **failed/abandoned** and the producer force-closes. `close()` with no argument waits up to a long default (effectively Long.MAX). Calling `close()` from inside a producer callback uses a zero-timeout variant to avoid deadlock. ## What goes wrong if you skip them - **Lost records on exit**: if the process exits (or you `System.exit`, or use `close(Duration.ZERO)`) while batches sit in the accumulator — perhaps still lingering for `linger.ms`, or queued — those buffered records never reach a broker and are **silently lost**. `send()` returning earlier gave no guarantee. - **Resource leaks**: not closing leaks the Sender thread, network sockets, and metrics. Always close in a `finally` or use try-with-resources (`KafkaProducer` implements `AutoCloseable`). - **False 'success'**: fire-and-forget `send()` that never checks the Future/callback hides delivery failures (e.g. record too large, auth error). For at-least-once you must observe the result *and* ensure flush/close before treating data as sent. ## Practical rules - Reuse one long-lived producer; `close()` once on shutdown. - Use try-with-resources for short-lived producers (auto-flush+close). - Call `flush()` when you need a durability checkpoint mid-life. - For delivery guarantees, check the callback/Future — don't rely on send() returning.

  • Does send() returning mean the record is safely in Kafka?
    No. send() only buffers the record and returns a Future; the background Sender transmits it later. The record is durable only when the Future/callback completes successfully per the acks setting. Relying on send() returning is the classic data-loss bug.
  • What's the difference between flush() and close()?
    flush() blocks until all currently buffered records complete but keeps the producer usable. close() does an implicit flush and then tears down the Sender thread and connections, after which the producer can't be used. close(Duration) bounds the wait and can drop pending records.
  • How can records be lost even though every send() call returned without throwing?
    send() returns after buffering, not delivery. If the JVM exits or you close with too short a timeout while batches are still in the accumulator (e.g. lingering for linger.ms), those buffered records never reach a broker and are silently lost. flush()/close() with adequate timeout prevents this.

saying these in an interview costs you the question

  • Saying send() transmits synchronously / guarantees delivery
  • Thinking flush() closes the producer
  • Closing with Duration.ZERO and assuming all records were sent
  • Fire-and-forget send() without ever checking the callback/Future
  • Creating and closing a producer per message

context