Why this matters
- VaultCommerce uses async send with callbacks for order events; sync .get() only in the outbox relay where failure must be immediate.
- Callbacks run on the producer's I/O thread — keep them fast or offload heavy work.
- RecordMetadata carries partition, offset, and timestamp for tracing.
- Understanding sync vs async send prevents thread pool exhaustion under load.
Async send with callbacks
The happy path: producer.send(record, (metadata, exception) -> { ... }). On success, log metadata.partition() and metadata.offset() for correlation with consumer processing. VaultCommerce attaches traceId in headers and logs it alongside broker metadata.
Key points
Future<RecordMetadata>— handle to in-flight send; completes on broker ack- Callback — invoked on success or failure; must not block
- RecordMetadata — partition, offset, timestamp, serialized sizes
- flush() — force send of all buffered records before shutdown
- max.block.ms — how long send() blocks when metadata or buffer is unavailable
Blocking send for critical paths
producer.send(record).get(5, TimeUnit.SECONDS) blocks until broker ack or timeout. Use sparingly — VaultCommerce's outbox poller uses blocking send so failed publishes increment retry counters before the outbox row is marked failed.
VaultCommerce rollout checklist
Before promoting changes that touch vaultcommerce.orders.placed.v1 and vaultcommerce.payments.captured.v1, run the staging KRaft cluster (Kafka 3.7+, Schema Registry 7.x) through a 10k events/min soak test. Compare producer request latency p99 and consumer lag per group against the pre-deploy baseline. send() is async by default; .get() blocks for ack. Document the change in the internal topic registry, attach Grafana screenshots to the change ticket, and keep an engineer on lag dashboards for 30 minutes after production rollout — roll back the service release before altering broker-level settings if lag or under-replicated partitions spike.
// VaultCommerce — async publish with structured callback
public void publishOrderPlaced(OrderPlaced event) {
ProducerRecord<String, OrderPlaced> record =
new ProducerRecord<>("vaultcommerce.orders.placed.v1", event.orderId(), event);
record.headers().add("traceId", traceId.getBytes(StandardCharsets.UTF_8));
producer.send(record, (meta, ex) -> {
if (ex != null) {
metrics.increment("kafka.publish.failure", "topic", meta == null ? "unknown" : meta.topic());
log.error("Publish failed orderId={}", event.orderId(), ex);
return;
}
log.info("Published orderId={} partition={} offset={}",
event.orderId(), meta.partition(), meta.offset());
});
}
Quick recall
Everything you need if you only revisit this box.
- send() is async by default; .get() blocks for ack.
- Callbacks must be lightweight.
- RecordMetadata gives partition/offset for tracing.
Test yourself
Answer these before moving on — recall is what makes it stick.