KafkaTemplate.send() returns a CompletableFuture<SendResult> — what does that mean for correctness, and how do you handle success and failure?
answer
- buffer + background I/O thread = async
- future completes per acks setting
- ignore future -> silent failures
- whenComplete / exceptionally, non-blocking callback
- get() blocks; flush() forces out
basics
~20 ssend() is asynchronous — it returns a future that completes when Kafka acknowledges the write. You attach whenComplete/thenAccept callbacks to handle the SendResult on success or the exception on failure, instead of assuming send succeeded.
solid answer
~40 ssend() enqueues the record into the producer's buffer and immediately returns a CompletableFuture<SendResult<K,V>>. The future completes on the producer's I/O thread when the broker acks (per the acks config) or fails. On success SendResult carries RecordMetadata (partition, offset, timestamp). Because the call doesn't throw on delivery failure, you must observe the future — typically future.whenComplete((result, ex) -> ...) or thenAccept/exceptionally — to log/retry/compensate; otherwise errors are swallowed. If you truly need synchronous semantics you call future.get(), but that blocks the calling thread and couples it to network latency. Note the future is completed on a Kafka thread, so callbacks must be fast and non-blocking. Serialization errors and buffer-full conditions can still throw synchronously from send() itself.
code
java · 16 linespublic CompletableFuture<SendResult<String, String>> publish(String key, String value) {
CompletableFuture<SendResult<String, String>> future =
kafkaTemplate.send("orders", key, value);
future.whenComplete((result, ex) -> {
if (ex != null) {
// delivery failed -> log, retry, or write to a dead-letter store
log.error("Failed to send key={}", key, ex);
} else {
RecordMetadata md = result.getRecordMetadata();
log.info("Sent key={} to {}-{}@offset {}",
key, md.topic(), md.partition(), md.offset());
}
});
return future; // let caller compose further
}go deeper
Know send() is async and returns a future that completes later.
Explain success (SendResult/RecordMetadata) vs failure handling and why ignoring the future loses errors.
Discuss acks-driven completion, callback runs on the I/O thread (don't block), get() vs flush() trade-offs, sync exceptions.
Reason about ordering with retries/idempotence, throughput vs durability trade-offs, and integrating the future into reactive/transactional flows.
**The async model.** Kafka's producer is built for throughput: records are accumulated into per-partition batches in an in-memory `RecordAccumulator` and a background *Sender* (I/O) thread ships full/expired batches to brokers. So `KafkaTemplate.send(...)` doesn't wait — it serializes the record, appends it to the batch buffer, and returns a `CompletableFuture<SendResult<K,V>>` (in Spring for Apache Kafka 3.x; earlier versions returned a `ListenableFuture`). **When the future completes.** The future is completed by the producer's I/O thread once the broker responds, and *how many* brokers must respond is governed by the `acks` producer setting: - `acks=0` — fire-and-forget, completes as soon as it's on the wire (weakest, can lose data). - `acks=1` — leader has written it. - `acks=all` (a.k.a. `-1`) — leader + all in-sync replicas; strongest durability. **On success** you get a `SendResult<K,V>` with the original `ProducerRecord` and the `RecordMetadata` (assigned partition, offset, timestamp). **On failure** the future completes exceptionally with e.g. `TimeoutException`, `RecordTooLargeException`, or broker errors. Crucially, `send()` itself normally does **not** throw for delivery problems — so if you ignore the future, failures vanish silently. This is the #1 correctness gotcha. **Handling patterns.** ```java kafkaTemplate.send("orders", key, value) .whenComplete((SendResult<String,String> res, Throwable ex) -> { if (ex != null) log.error("send failed for {}", key, ex); else log.info("sent to {}-{}@{}", res.getRecordMetadata().topic(), res.getRecordMetadata().partition(), res.getRecordMetadata().offset()); }); ``` Use `thenAccept` for success-only, `exceptionally` for error-only, or `whenComplete` for both. **Do not block inside the callback** — it runs on the single producer I/O thread and blocking it stalls all other sends. **Synchronous send.** If you need to be sure before continuing (e.g. inside a request you must not ack until the message is durable): ```java SendResult<String,String> res = kafkaTemplate.send(...).get(10, TimeUnit.SECONDS); ``` `get()` blocks the caller until completion and rethrows failures wrapped in `ExecutionException`. It trades throughput for certainty and ties your thread to network latency — use sparingly. **Synchronous exceptions that CAN come straight out of send().** Serializer failures (bad object → `SerializationException`) and a full buffer when `max.block.ms` expires (`TimeoutException` / `BufferExhaustedException`) surface from the call itself, not the future. So robust code both try/catches send() *and* observes the future. **flush().** `kafkaTemplate.flush()` forces all buffered records out and blocks until they're sent — useful at shutdown or in tests. `send()` alone never guarantees the record has left the JVM. **Ordering caveat.** With retries enabled and `max.in.flight.requests.per.connection > 1`, a retried record can be reordered relative to later ones — enabling idempotence (`enable.idempotence=true`) preserves ordering while allowing in-flight >1.
- What is the risk of calling future.get() right after send() in a request-handling thread?It turns an async producer into a blocking call: the request thread waits for the broker ack, so throughput drops and latency spikes tie up your thread pool. It also loses the batching benefit if you do it per-message. Use it only when durability-before-response is required.
- Can send() ever throw an exception directly rather than via the future?Yes — serialization errors (SerializationException) and buffer-exhaustion timeouts (max.block.ms exceeded) are thrown synchronously from send(). Delivery/broker errors instead complete the future exceptionally.
saying these in an interview costs you the question
- Assuming a successful return from send() means the message reached the broker
- Never observing the returned future, so delivery failures are lost
- Blocking inside the whenComplete callback (it runs on the producer I/O thread)
- Calling get() per message and destroying throughput without realizing it