Producer Configuration
Learn ProducerBuilder, Acks.All, EnableIdempotence, and ProduceAsync. Understand what each config option actually does to reliability and throughput.
SPACED REPETITION Β· 15 practice questions
Make this lesson stick.
Try 3 questions now. No account needed. Sample answers aren't saved.
or sign in to practice all 15First Contact: ProducerConfig, ProducerBuilder, and the Life of One ProduceAsync Call
Your first Kafka producer will almost certainly work. That is the trap. Three lines, a record lands in a topic (a named stream of records), and it feels like you have learned the client. Then a deployment drops a few thousand records during shutdown, or a message appears on partition 2 when you assumed partition 0. This section builds the mental model on the smallest possible producer β one record β so the machinery stays visible while you learn it.
ProducerConfig is a typed bag of strings
ProducerConfig is a property bag: each .NET property maps to a lowercase librdkafka configuration key. BootstrapServers becomes bootstrap.servers, Acks becomes acks, EnableIdempotence becomes enable.idempotence. The typed class is a convenience over a dictionary: setting a property checks nothing. Validation happens at Build(), where librdkafka rejects unknown keys and illegal values β and even then nothing is checked against a broker (a Kafka server). Confluent.Kafka wraps librdkafka, the C library that runs the sockets, the buffering, and the background I/O threads, so every setting you write ends up as a librdkafka config string.
ProducerBuilder<TKey, TValue> takes that config and adds the one thing a property bag cannot hold: serializers, the converters that turn your TKey and TValue into bytes. For string, the builder uses the built-in UTF-8 serializer, so bytes come free. Build() then creates the native librdkafka producer handle β the object that owns the connection pool, the send queue, and the threads doing the network work.
using Confluent.Kafka; // ProducerConfig, ProducerBuilder, Message
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092"
};
using var producer = new ProducerBuilder<string, string>(config).Build();
var result = await producer.ProduceAsync(
"orders",
new Message<string, string>
{
Key = "customer-42",
Value = "order-1"
});
// Result type: DeliveryResult<string, string>
// Prints topic, partition and offset, e.g. orders [[0]] @7
// (Partition prints its own brackets, hence the double pair)
Console.WriteLine(result.TopicPartitionOffset);
result is a DeliveryResult<TKey, TValue>: the broker's receipt for one record. A partition is one ordered shard of a topic; an offset is a record's position within that shard's log. result.Partition and result.Offset are the two fields worth memorizing β and neither is assigned by your code.
BootstrapServers is a phone book, not a destination
BootstrapServers is the initial contact list the client uses to fetch cluster metadata β which brokers exist, and which broker is the leader (the broker currently responsible) for each partition. Metadata is fetched at startup and refreshed periodically (topic.metadata.refresh.interval.ms). After the first refresh, the producer talks directly to the partition leader; records never route through the bootstrap address.
Two consequences follow. A broker restart does not mean editing every client: as long as one reachable address stays in the list, the client refetches metadata and finds the new leader. And listing every broker is optional β one stable address that outlives leader changes is enough to bootstrap.
What ProduceAsync actually does
ProduceAsync does not perform a synchronous socket write. It serializes the message, enqueues it into librdkafka's internal buffer, and returns a Task<DeliveryResult<TKey, TValue>>. The await completes later, when the record is acknowledged or the client gives up β governed by Acks (how many replicas must have the record), retries, and timeouts (the full durability contract for those knobs belongs to "Acks, Retries, and the Real Meaning of 'Delivered'"). The queue holding unsent records is bounded by QueueBufferingMaxMessages and QueueBufferingMaxKbytes; when it is full, ProduceAsync fails at once with a local queue-full error instead of accepting unlimited work.
That single fact β enqueue now, resolve later β explains most producer behavior downstream.
Worked trace: one key, one partition
Suppose orders has 3 partitions. We send two records with the same key, then one with no key. The default partitioner in Confluent.Kafka is consistent_random: a non-null key is hashed to a partition; a null or empty key gets a random partition (the same one for a few milliseconds at a time, so batches can form). The hash is CRC32, whereas the Java client hashes keys with a different function, murmur2, so a default .NET producer and a default Java producer can send the same key to different partitions β a real hazard when a Java service and a .NET service must agree (murmur2_random is the partitioner for that case).
send: key="customer-42" value="order-1"
partitioner: hash("customer-42") -> partition 0
deliver: orders [[0]] @7 result.Partition=0, result.Offset=7
send: key="customer-42" value="order-2"
partitioner: hash("customer-42") -> partition 0 (same key, same partition)
deliver: orders [[0]] @8
send: key=null value="order-3"
partitioner: no key -> random partition
deliver: orders [[2]] @3 (could have landed anywhere)
Ordering is per partition, not per topic. Both customer-42 records land on partition 0 in enqueue order, so a reader of that partition sees them in order; order-3 shares no ordering relationship with them once it lands elsewhere. Choosing keys is the subject of the lesson "Topics, Partitions & Keys" β here you only need result.Partition and result.Offset.
Predict before you run
Commit to an answer first.
Topic
ordershas 3 partitions, each kept as copies (replicas) on several brokers, andmin.insync.replicas=1(one up-to-date copy is enough for anAcks=Allwrite to be accepted). The producer runs withAcks=All(wait for every up-to-date replica) andEnableIdempotence=false(no duplicate protection on retries). You callProduceAsync, and the partition leader crashes immediately afterward. Does the returned task complete successfully, throw, or hang?
Check your answer
It can complete successfully or throw, depending on timing β it does not hang forever, because a timeout bounds it.
A replica is a copy of a partition's log on a broker (the leader is one of them), and min.insync.replicas is the smallest number of in-sync replicas β replicas caught up with the leader, the leader included β the broker will accept an Acks.All write with. With it set to 1, a leader that is the only in-sync replica may acknowledge on its own. So the call may well succeed β and still be lossy: if that leader dies before a follower copies the record, the record is gone even though your await returned success. Acks.All is therefore only as strong as the topic's min.insync.replicas, which is a broker/topic setting, not a producer setting.
If no leader is elected before the deadline, the task faults with a ProduceException. The governing timeout is MessageTimeoutMs (message.timeout.ms, which librdkafka also accepts under the alias delivery.timeout.ms; 5 minutes by default): the total time a record may spend in the producer, including retries and time waiting for a leader. (request.timeout.ms is something else β how long the broker waits for replicas on one request β and is not the whole budget.)
The shutdown mistake
Here is the bug that bites almost everyone once:
// BUG: fire-and-forget sends, then immediate dispose
// (orders: your order objects, with CustomerId and Json; config as above)
using (var producer = new ProducerBuilder<string, string>(config).Build())
{
foreach (var order in orders)
{
producer.Produce( // callback-based overload; not awaited
"orders",
new Message<string, string> { Key = order.CustomerId, Value = order.Json });
}
} // Dispose here
Produce returns before the record is delivered. When the using block exits, Dispose() destroys the native handle without flushing: queued records are dropped, in-flight ones are never reported, and because no delivery handler was attached, nothing tells you.
Make shutdown explicit:
foreach (var order in orders)
{
producer.Produce(
"orders",
new Message<string, string> { Key = order.CustomerId, Value = order.Json },
deliveryReport => // always observe the outcome
{
if (deliveryReport.Error.IsError)
logger.LogError(deliveryReport.Error.Reason); // logger: your ILogger
});
}
// Bound the wait; a non-zero return means records are still queued or unacknowledged.
var remaining = producer.Flush(TimeSpan.FromSeconds(10));
if (remaining > 0)
throw new InvalidOperationException($"{remaining} records unflushed at shutdown");
Flush(TimeSpan) waits until everything queued has been delivered (or has failed) or the timeout elapses, and returns roughly how many records are still waiting to be sent or acknowledged; zero versus non-zero is the reliable signal. Treat a non-zero return as a shutdown failure, not a warning β those are the records you will be missing tomorrow.
Your turn
Modify the minimal example to send three records: two with key "customer-42" and one with key "customer-7". Before running it, write down which records you expect to share a partition (the partition number is the hash's business; only the grouping is predictable). Then diagnose this: a service disposes its producer immediately after a loop of callback Produce calls, and a downstream consumer occasionally sees gaps. Which single line would you add, and what does its return value tell you?
Acks, Retries, and the Real Meaning of 'Delivered'
ProduceAsync hands back a DeliveryResult, and that return value feels like a promise: the record was delivered. But delivered is a claim made by a broker β the server that stores records and answers client requests β and different brokers make that claim on very different evidence. A leader that wrote the record only to its own log and a full set of replicas that all wrote it return the same successful task. The gap between those two only becomes visible when a machine dies.
Choosing Acks is choosing which evidence you demand. The retry and timeout settings decide how long you are willing to wait for it.
The durability ladder: None, Leader, All
Acks is the producer setting that names how many replicas must have the record before the produce is acknowledged (with None, nothing is acknowledged at all); in Confluent.Kafka it maps to librdkafka's acks key. Three rungs:
Acks |
Waits for | Where loss can still happen |
|---|---|---|
Acks.None |
nothing | any broker-side failure or rejection β none is reported |
Acks.Leader |
leader's local log | leader dies before replicas copy it |
Acks.All |
every in-sync replica | only as strong as the topic's min.insync.replicas |
Three words to keep straight: a partition is one ordered slice of a topic; a replica is a copy of a partition stored on a broker; the leader is the replica that accepts writes for that partition while followers copy from it.
Acks.None is the strange rung. The client writes the record to the socket and reports success immediately; the returned DeliveryResult carries no broker-assigned offset, because no broker has spoken. Retries cannot react to a broker-side failure here β Acks.None produces no failure signal. Set Acks explicitly rather than inheriting a default: librdkafka defaults to all, the Java client defaulted to leader-only before Kafka 3.0, and this is not a setting to discover by accident.
Acks.Leader is a real improvement β the leader genuinely appended the record β but the acknowledgement is local. If that leader dies before followers copy the record, the record is gone.
Acks.All asks the leader to wait until every in-sync replica (ISR) β every copy currently caught up with the leader β has written the record. It is the strongest thing a producer can request.
Acks.All is a promise about the topic, not just the producer
Here is where producers get fake confidence. Acks.All does not mean all replicas. It means all replicas currently in sync. If replication has fallen behind, that set is smaller, and the promise shrinks with it. Two properties of the topic decide how much it is worth:
- replication.factor β how many copies of each partition exist at all.
- min.insync.replicas β the smallest ISR count the broker will accept an
Acks.Allwrite with; other writes are not checked against it.
β οΈ min.insync.replicas is a broker/topic setting. No producer config can set it, and no producer config can override it. With a replication factor of 3:
min.insync.replicas = 2
ISR {leader, followerA, followerB} -> accepted
followerB dies, ISR {leader, followerA} -> 2 >= 2, accepted
followerA dies, ISR {leader} -> 1 < 2, REJECTED
=> survives one broker failure, stays durable
min.insync.replicas = 1
ISR {leader, followerA, followerB} -> accepted
both followers die, ISR {leader} -> 1 >= 1, acked by one replica
leader dies before replicating -> record LOST
The second case is the trap: Acks.All is set, the produce returns success, and the data is still losable. The first case trades availability for durability; the second trades durability for availability.
When the ISR drops below min.insync.replicas, the broker rejects an Acks.All write with one of two error codes:
| Error code | What happened |
|---|---|
NotEnoughReplicas |
rejected before the leader appended |
NotEnoughReplicasAfterAppend |
leader appended, then the ISR shrank below the minimum |
librdkafka retries both: after NotEnoughReplicas the record is not in the log, after NotEnoughReplicasAfterAppend it may be. Your code sees the code (in a ProduceException<TKey, TValue> from ProduceAsync, or in the delivery report passed to a Produce callback) only if the retry count runs out; with the default, practically unlimited count the record expires instead and the error is Local_MsgTimedOut.
The retry and timeout budget
If Acks says what you are waiting for, the timeouts say how long you will wait, and the retries say how stubbornly.
| Confluent.Kafka | librdkafka | Bounds |
|---|---|---|
MessageSendMaxRetries |
message.send.max.retries (alias retries) |
retries after the first attempt |
RetryBackoffMs |
retry.backoff.ms |
first pause; doubles per retry up to retry.backoff.max.ms (1 s by default) |
RequestTimeoutMs |
request.timeout.ms |
how long the broker waits for replica acks on one request |
SocketTimeoutMs |
socket.timeout.ms |
how long the client waits for the response to one request |
MessageTimeoutMs |
message.timeout.ms |
the whole record, queue included |
The last one is the one people get wrong by name. MessageTimeoutMs is librdkafka's version of the JVM producer's delivery.timeout.ms (librdkafka accepts that name as an alias): it covers the record's entire life inside the producer β time parked in the local queue, every attempt, every backoff. The names differ across clients and so do the defaults (5 minutes here, 2 in the Java client); the meaning does not.
MessageTimeoutMs whole record: queue + retries + backoffs
|
+-- per attempt: the client waits min(SocketTimeoutMs, what is left of MessageTimeoutMs)
attempt 1 (broker waits up to RequestTimeoutMs for replicas, then answers)
backoff RetryBackoffMs
attempt 2
backoff 2 x RetryBackoffMs (doubling, capped at retry.backoff.max.ms)
...
Two things to hold on to. A single attempt cannot outlive the record's budget, because the client caps each request at whatever is left of MessageTimeoutMs. And RequestTimeoutMs is not a client-side clock, as the setting of the same name is in the Java client: it tells the broker how long to wait for replicas before answering with a timeout error (30 s by default, inside the client's 60 s wait). In the other direction, an enormous MessageTimeoutMs β the default is 5 minutes β means a shutdown flush can sit waiting for it. β οΈ That is not a hang you want during a rolling deployment; "Configuration Discipline: Validation, Environment Overrides, Shutdown, and a Combined Diagnostic" shows how to bound it. Pick MessageTimeoutMs from what your caller can tolerate, then keep the per-request timeouts well below it.
Idempotence: retries that don't duplicate or reorder
Retries create an ambiguity. Suppose the leader appends your record and then dies before the acknowledgement reaches you. The producer cannot tell a record that never arrived from a record that arrived with a lost ack. Its instinct is to retry β which can append the record twice.
EnableIdempotence removes the ambiguity. With it on, the producer obtains a producer ID from the broker and stamps each record with a per-partition sequence number. The broker remembers which sequences it has already accepted, so a retried record it has seen before is discarded rather than appended, and records are appended in sequence order β which preserves per-partition ordering across retries.
The fact to carry out of this section: idempotence is the bridge between retry-on-failure and retry-without-duplicates-or-reordering. It requires Acks.All and a positive retry count, and it requires MaxInFlight β the cap on requests sent but not yet acknowledged per broker connection β to be 5 or lower; how that cap interacts with concurrency belongs to 'ProduceAsync in the Hot Path'.
β οΈ Idempotence is per producer session, and it covers only the client's own retries. Restart the producer and it gets a new producer ID; send a record again from your own code β after a restart, or after a ProduceAsync that failed β and the broker sees a new record, which can duplicate the first. Producer-side idempotence is not end-to-end exactly-once β which is why "effectively once" in practice usually means at-least-once delivery (every record arrives, possibly more than once) plus an idempotent consumer.
Worked trace: one call, two contracts
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
Acks = Acks.All, // durability ladder: None < Leader < All
MessageSendMaxRetries = 3, // 3 retries on top of the first attempt
RetryBackoffMs = 200, // first backoff; doubles on each retry
MessageTimeoutMs = 20_000 // 20s total budget for one record
};
// topic: replication.factor = 2, min.insync.replicas = 2
// ISR when the call is made: { leader } (the follower is already out)
Scene: the follower's machine died a minute ago, and it is long gone from the ISR (a dead broker is removed within seconds, not milliseconds). Times are approximate; each backoff varies randomly by up to 20%.
t=0ms ProduceAsync sends to leader
tβ10ms ISR is {leader}: 1 < min.insync.replicas=2
leader rejects with NotEnoughReplicas, nothing appended (attempt 1)
backoff β200 ms
tβ0.2s attempt 2 -> NotEnoughReplicas
backoff β400 ms
tβ0.6s attempt 3 -> NotEnoughReplicas
backoff β800 ms
tβ1.4s attempt 4 -> NotEnoughReplicas
retries exhausted (3 retries + the first attempt)
-> ProduceException thrown to the awaiting caller
Read the trace closely: the 20-second MessageTimeoutMs never became the binding constraint, because retry exhaustion got there first, in about a second and a half. Had the retry count been larger, the deadline would take over instead β retries about once a second (the backoff stops growing at retry.backoff.max.ms) until t=20s, then a ProduceException carrying Local_MsgTimedOut. And if the follower rejoins the ISR before either limit hits, a retry succeeds and the caller sees an ordinary DeliveryResult. Three exits, one call; the bad exit honestly reports that the record is not in the log, because NotEnoughReplicas rejects before the append.
Now change exactly one line β Acks = Acks.Leader β and leave everything else identical:
t=0ms ProduceAsync sends to leader
tβ10ms leader appended locally -> ack returned
ProduceAsync completes SUCCESSFULLY
Same code, same dead broker, opposite result: min.insync.replicas is enforced only for Acks.All. And that success is thin: the only other replica β the follower β is on the machine that is down, so the record exists in exactly one place. If the leader goes down next, there is no second copy, and no retry in your service can recover what was never replicated.
Decide: three producers
For each row, choose Acks, the min.insync.replicas you would demand of the topic, a retry count, whether EnableIdempotence must be on, and the failure mode you are accepting β loss, duplicate, or latency.
| Producer | Publishes | Accepts |
|---|---|---|
| Payments | charge-authorization events for a ledger | ? |
| Clicks | page-view counts for a dashboard | ? |
| Audit | compliance records that must not vanish | ? |
Then answer the twist: the audit topic was created with replication.factor=3 and min.insync.replicas=1, and the requirement is no loss during a single broker failure. What is the one configuration change that satisfies it?
Check your answer
- Payments:
Acks.All,min.insync.replicas>=2, idempotence on, generous retries. Accept latency, not loss β and the ledger still needs its own idempotency key, because producer config cannot see the consumer's world. - Clicks:
Acks.Leader(whichmin.insync.replicasdoes not govern), low retry count, idempotence off (it requiresAcks.All). Accept loss for throughput and availability; a missing page-view count is noise. - Audit:
Acks.All,min.insync.replicas>=2, idempotence on, generous retries. Accept duplicates (deduplicate downstream) but never loss.
The twist has one honest answer: change the topic, not the producer. Set min.insync.replicas=2 on the audit topic. Because min.insync.replicas is a topic setting, a producer-only fix cannot work: at min.insync.replicas=1, Acks.All is satisfied by a leader that is the only in-sync replica, and that leader dying before a follower copies the record loses it. The replication factor is what makes the change affordable: replication.factor=2 with min.insync.replicas=2 would leave no slack β one broker down and every write is rejected β whereas replication.factor=3, min.insync.replicas=2 survives one broker failure while still accepting writes.
β οΈ Finally, the pitfall these settings invite. Acks.All plus retries without EnableIdempotence can append the same record twice: the first attempt succeeded, the ack was lost, and the retry wrote a copy. With more than one request in flight at a time (max.in.flight.requests.per.connection > 1, which is the default), a retried earlier record can also land after a later one, reordering the partition. These settings buy at-least-once delivery with a reorder-free option β not exactly-once (each record taking effect one time, end to end). Producer configuration alone never gives you exactly-once.
ProduceAsync in the Hot Path: Ordering, Concurrency, and In-Flight Limits
The moment you stop sending one message and start sending a stream, a new question appears: what does the log actually look like when several sends are outstanding at once? This section is about the three-way tension between throughput, per-partition order, and the in-flight window β and about the timeouts that decide whether a record you believe you sent ever landed.
Awaiting versus pipelining
The await in await producer.ProduceAsync(...) completes when librdkafka reports an acknowledgement. In a loop, that makes everything sequential:
// orders: your order objects; Ser(...): your serializer (both elided)
foreach (var order in orders)
{
// Each iteration waits for a full broker round trip before queuing the next.
await producer.ProduceAsync("orders",
new Message<string, string> { Key = order.AccountId, Value = Ser(order) });
}
Every iteration pays enqueue β network β append β ack before the next record is even queued. Ordering is trivially safe; throughput is a fraction of the client's capacity. This is the sequential producer: simple, slow, ordered.
To pipeline, stop awaiting each send before issuing the next:
// Shape 1: issue all sends, then await together.
var pending = orders.Select(o => producer.ProduceAsync("orders",
new Message<string, string> { Key = o.AccountId, Value = Ser(o) })).ToList();
var results = await Task.WhenAll(pending); // β οΈ unbounded β see the pitfall below
// Shape 2: fire-and-forget with a delivery callback (Log: your logger).
producer.Produce("orders",
new Message<string, string> { Key = order.AccountId, Value = Ser(order) },
report =>
{
if (report.Error.IsError)
Log.Error("produce failed: {reason}", report.Error.Reason);
});
Both keep multiple requests in flight. Both move responsibility onto you: with Produce, the callback is the only place a failure surfaces β ignore it and failed records vanish silently. Any path that queues without awaiting must call producer.Flush(TimeSpan.FromSeconds(10)) before disposal; the disposal trap is covered in "First Contact: ProducerConfig, ProducerBuilder, and the Life of One ProduceAsync Call".
Order lives in a partition, not a topic
Kafka appends records to partitions β the ordered, append-only shards that a topic is split into. Within one partition, offsets increase in append order, and that is the only order Kafka promises; across partitions of the same topic, there is none. The partitioner is the client-side function that chooses the partition: with a non-null key, the default implementation hashes the key so identical keys always map to the same partition. Key choice is ordering choice. Which keys you pick, and the traps in doing so, belong to the lesson "Topics, Partitions & Keys"; here we need only the consequence.
The in-flight window is where order breaks
An in-flight request is a produce request librdkafka has sent but not yet had acknowledged. max.in.flight.requests.per.connection β MaxInFlight in Confluent.Kafka .NET β caps how many can be outstanding per broker connection. The .NET default is effectively unbounded (librdkafka's default is 1000000), so unless you set it, you pipeline as hard as the client allows.
Prediction before the trace: three records with key acct-7 are sent far enough apart that, with LingerMs = 0 (no waiting to collect a batch), each leaves in its own produce request. The settings are Acks.All, retries on, MaxInFlight = 5, EnableIdempotence = false. Request 2 gets a retriable error from the broker and is retried; requests 1 and 3 succeed. What order does the log show?
All three requests are in flight at once (records that were queued together would instead share one batch and could not overtake each other):
- record 1 β acked, appended;
- record 2 β the broker answers with a retriable error (for example
NotEnoughReplicasat a moment when the ISR was short); librdkafka requeues the record and retries after the backoff; - record 3 β acked, appended before the retry lands.
Log order: 1, 3, 2.
intended: 1 β 2 β 3
log: 1 β 3 β 2 (retried 2 appended after 3)
Nothing was lost; the acknowledgement for 2 eventually succeeded. Order was what broke. Two fixes trade differently:
| Fix | Order preserved | Pipelining |
|---|---|---|
MaxInFlight = 1 |
Yes β next request waits for ack | None: one request at a time |
EnableIdempotence = true, MaxInFlight β€ 5 |
Yes β broker sequences records | Up to 5 requests |
MaxInFlight = 1 is blunt: only one request per broker connection is outstanding at a time, so pipelining is gone (each request still carries a whole batch, which keeps it ahead of the sequential loop). Idempotence is surgical. With it on, the producer attaches a producer ID and per-partition sequence numbers, and the broker accepts appends in sequence; librdkafka requires MaxInFlight to be 5 or lower when idempotence is enabled. This is the same machinery described in "Acks, Retries, and the Real Meaning of 'Delivered'", where it explained duplicate suppression across retries β here it is what stops a retried record from jumping the queue. β οΈ The guarantee is per producer session. Restart the producer and a new session begins with a new producer ID; send records again from your own code and they are brand-new records with new sequence numbers. Either way the broker cannot link them to the originals; that gap is picked up in "Configuration Discipline: Validation, Environment Overrides, Shutdown, and a Combined Diagnostic".
Three clocks, one record
When a send hangs or is cancelled, you must know which clock fired:
CancellationTokenonProduceAsyncβ cancels your wait, not the write. The record may still be enqueued and delivered, and a late acknowledgement may still arrive. There is no "unsend."- The per-request clock, which has two sides β
RequestTimeoutMs(request.timeout.ms) is the broker's: how long the leader waits for replica acknowledgements before answering one request with a timeout error.SocketTimeoutMs(socket.timeout.ms) is the client's: how long it waits for that answer, capped by what is left of the record's budget. When either expires the request may be retried; it is not the end of the record's life. MessageTimeoutMs(librdkafka'smessage.timeout.ms, the counterpart of the Java client'sdelivery.timeout.ms) β the producer-level deadline for the record itself, from enqueue through all retries to final success or failure.
Set MessageTimeoutMs when you need a record resolved within a bounded time, and read a cancelled ProduceAsync as "unknown, possibly written" rather than "not sent." Sending it again can duplicate it; treating it as done can lose it.
Guided task: ordered and pipelined
This loop sends order events for many accounts. Every event for an account must land in production order, but you refuse to make the whole service sequential:
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
Acks = Acks.All,
EnableIdempotence = false,
// MaxInFlight left at its default
};
using var producer = new ProducerBuilder<string, string>(config).Build();
await Parallel.ForEachAsync(orderEvents, async (evt, ct) =>
{
// Ser(evt) is your serializer; System.Text.Json is fine.
await producer.ProduceAsync("orders",
new Message<string, string> { Key = evt.AccountId, Value = Ser(evt) }, ct);
});
Change the configuration β and the loop, if needed β so all messages for one account appear in the log in the order produced, without sending one record at a time across the whole service.
Check your answer
Order can break in two places, so there are two changes:
- Retries.
EnableIdempotence = true(withMaxInFlight β€ 5; librdkafka sets 5 itself when you leave it unset). Order survives retries while up to five requests stay in flight. Idempotence requiresAcks.All, which the config already sets. - Submission.
Parallel.ForEachAsyncover single events lets two events of the same account race into the producer's queue, and no producer setting can repair an order that was already wrong at enqueue. Group the events byAccountId, run the accounts in parallel, and send each account's events one after another:
await Parallel.ForEachAsync(orderEvents.GroupBy(e => e.AccountId), async (account, ct) =>
{
foreach (var evt in account) // one account: strictly in order
{
await producer.ProduceAsync("orders",
new Message<string, string> { Key = evt.AccountId, Value = Ser(evt) }, ct);
}
});
MaxInFlight = 1 would also protect order across retries, but it is the option the task rules out: it limits every broker connection to one outstanding request, for the whole service.
Record in the service's docs which guarantee the rest of the system relies on β the configuration half of the fix is invisible at the call site.
Pitfall: unbounded concurrency
Starting every ProduceAsync before awaiting any looks like "just await them all," but it enqueues every record at once. The producer's local queue is bounded (QueueBufferingMaxMessages, 100000 records by default, and QueueBufferingMaxKbytes in Confluent.Kafka), and when it is full, ProduceAsync fails immediately with a local queue-full error (ErrorCode.Local_QueueFull) β an error that looks like broker trouble but originates in your process. Bound the fan-out with a SemaphoreSlim, a Channel, or Parallel.ForEachAsync's MaxDegreeOfParallelism, and let backpressure β the slowdown passed back up the pipeline β reach whatever produces the events. Batching and buffer sizing are treated in "Throughput vs. Latency: Batching, Compression, and Buffer Memory".
Before moving on: can you say (a) which producer paths await each send and which pipeline, (b) whether MaxInFlight is set in your config or left at the default, (c) what happens to order on a retry under your settings, and (d) which of the three clocks bounds a record's life in your service?
Throughput vs. Latency: Batching, Compression, and Buffer Memory
An orders service adds one line β LingerMs = 50 β to "improve throughput." Produce p99 latency (the value 99% of sends stay under) jumps from 8 ms to 53 ms and breaks a downstream target of sub-20 ms acknowledgements. Nothing about acks, retries, or serialization changed. Batching knobs are not switches you flip; they are a budget you spend, buying broker and network efficiency with client-side milliseconds at a rate you should calculate rather than guess.
Batching: LingerMs Sets the Deadline, BatchSize Sets the Cap
linger.ms (Confluent.Kafka: LingerMs) tells the producer to wait up to that many milliseconds for more records to accumulate in a batch β one grouped set of records sent in a single produce request β before sending. batch.size (BatchSize, in bytes) is the target maximum size of that batch. Two moving parts, two jobs: LingerMs is the deadline, BatchSize is the cap, and whichever triggers first ends the wait. (librdkafka also caps a batch at batch.num.messages, 10000 records by default, which binds before BatchSize only for very small records.)
record arrives
β
append to the partition's batch
β
send when: BatchSize reached OR LingerMs expired
β
compress the batch (if enabled)
β
one produce request to the partition leader
That distinction matters because learners tune the wrong knob. At low volume a 1 MB BatchSize almost never fills β the deadline does all the work. One stack detail worth internalizing: Confluent.Kafka inherits librdkafka's defaults, which are not the Java client's. librdkafka's linger.ms defaults to 5; the Java client defaulted to 0 before Kafka 4.0 and to 5 since. batch.size still differs: 1000000 bytes in librdkafka, 16384 in the Java client. Tuning notes copied verbatim from a Java blog post are tuning a different default. And linger.ms = 0 does not mean "one record per request" β it means "do not wait for more," while the producer still ships whatever is already queued as one batch.
The Accumulation Math
Say 10,000 messages/sec at 500 bytes each β 5 MB/s of payload, all headed for one partition (batches are per partition). With LingerMs = 5 and BatchSize = 1 MB, a full batch needs 1 MB Γ· 5 MB/s = 200 ms of arrivals. Because the 5 ms deadline fires long before the cap, the batch actually ships with far less than 1 MB and latency stays near 5 ms plus network time. Raise LingerMs = 200 and the same stream fills a real 1 MB batch β you finally get the compression and request-count wins β but each record can now wait up to 200 ms, with the typical wait roughly half that depending on arrival gaps. Compute both ends before shipping a change.
Compression: It Rides on Batches
compression.type (CompressionType) selects None, Gzip, Snappy, Lz4, or Zstd. Compression is applied per batch, not per record, so batching is a prerequisite: compressing single-record batches burns CPU and saves almost nothing. Turn compression on only after LingerMs and BatchSize are doing their job.
| Codec | Ratio | CPU cost | Typical use |
|---|---|---|---|
| none | 1Γ | lowest | latency-critical, small payloads |
| lz4 | modest | low | speed-first throughput |
| zstd | strong | moderate | large volumes |
| gzip | strong | higher | compatibility-first |
| snappy | modest | low | legacy ecosystems |
Treat the table as a starting heuristic, not a ranking β real ratios depend on payload shape, and the only authority on CPU cost is your own producer's CPU graph. Compressing 50 MB/s on a container with two shared cores moves the bottleneck from the network into the process, and now every handler on that pod slows down.
Buffer Memory: Where Backpressure Lives
buffer.memory is the JVM client's name for the total bytes a producer may hold for unsent records. librdkafka has no setting by that name; its queue is bounded twice, by QueueBufferingMaxMessages (the count of queued records, 100000 by default) and QueueBufferingMaxKbytes (their total size in kilobytes, 1048576 by default, which is 1 GB). Whichever limit is reached first makes the queue full β with 500-byte records that is the record count, at about 50 MB.
When the queue is full, Produce and ProduceAsync fail immediately β a ProduceException carrying ErrorCode.Local_QueueFull, a local error, not a broker error. Nothing blocks: the Java client's max.block.ms wait has no counterpart here. That is backpressure arriving at your caller, and the right response is to slow the source (pause intake, wait, try again). Raise the limits only to ride out a stall you have sized, as below; raised as a reflex, they convert an immediate, visible failure into a delayed, larger one: the queue soaks up records a sick broker isn't draining, and when it finally fills you have far more unsent data and a memory-pressure failure instead of a throughput problem.
Run the headroom numbers with 500-byte records. At 10,000 messages/sec (5 MB/s), a 2-second broker stall queues 20,000 records, about 10 MB β far inside both defaults. At 100,000 messages/sec (50 MB/s), the same stall would queue 200,000 records, about 100 MB; the default limit of 100000 records is reached after 1 second, long before the 1 GB size limit matters. Riding it out means raising QueueBufferingMaxMessages (to 300000, say) and accepting up to 150 MB of queued payload in your process. Size the queue for your worst stall, or accept that produce calls will fail with queue-full and handle that; decide deliberately.
Size Limits Fail at Publish, Not at Build
Producer-side MessageMaxBytes (message.max.bytes, 1000000 by default) caps the size of what the client will send; the broker's message.max.bytes and the topic's max.message.bytes cap what the broker accepts. These are separate limits, and Build() validates none of them against the broker: a producer whose records exceed the topic's maximum builds fine and then fails every oversized publish. If a payload format doubles your record size, check all three numbers before blaming the network.
Practice: Tune Two Producers
Choose LingerMs, BatchSize, CompressionType, the queue limits, and Acks for each. Then name the one value you would not tighten further without measurement β and what you would measure.
- A β Payment authorizations. Moderate volume, small payloads, p99 produce latency budget under 20 ms, zero tolerance for loss on a single broker failure.
- B β Clickstream. High volume, repetitive payloads, 200 ms extra latency acceptable, a few seconds of lost metrics tolerable.
Check your answer
A is latency-bound, so LingerMs stays near 0 and BatchSize stays default β a large cap is harmless because the deadline fires first. CompressionType stays None (small payloads gain little) or Lz4 if you measure a win. Default queue limits are ample at moderate volume. Acks = All, paired with a topic whose min.insync.replicas supports it β that pairing is developed in "Acks, Retries, and the Real Meaning of 'Delivered'." The value not to touch without measurement is LingerMs: every millisecond you add comes straight out of the 20 ms budget, so measure p99 before and after any change.
B is throughput-bound, so LingerMs rises (20β100 ms), BatchSize can stay at librdkafka's 1 MB default, and CompressionType becomes Lz4. Size the queue limits from the arithmetic above β at tens of MB/s, a multi-second stall needs hundreds of MB and a QueueBufferingMaxMessages raised to match, or you accept queue-full failures and drop. Durability can be looser than for payments. The value not to tighten further is compression: do not move from Lz4 to Zstd (or raise the level) until you have producer CPU numbers, because CPU is the resource most likely to become the new bottleneck.
The general judgment: identify the scarce resource β network, latency budget, CPU, or memory β and spend the others to protect it.
β οΈ Three pitfalls survive code review. linger.ms = 0 does not guarantee low latency; it removes deliberate waiting while still batching whatever is already queued. Enabling compression without measuring CPU relocates the bottleneck into your process. And raising the queue limits without sizing them masks a broker throughput problem until memory pressure turns it into a harder, later failure β the pattern that bounded concurrency in "ProduceAsync in the Hot Path" is designed to prevent.
Configuration Discipline: Validation, Environment Overrides, Shutdown, and a Combined Diagnostic
A producer that passes every dev test can still lose records in production, and the cause is usually not one exotic setting β it is an assembly problem. The effective config was stitched together from appsettings.json, a deployment values file, and a code default, and nobody ever looked at the result. Treat producer configuration as a reviewed artifact: one validated base, overrides that cannot vanish silently, a bounded shutdown, and a diagnostic that reads an entire config at once.
One typed base, validated before Build()
Confluent.Kafka's ProducerConfig is a property bag of librdkafka configuration keys (introduced in "First Contact: ProducerConfig, ProducerBuilder, and the Life of One ProduceAsync Call"). Build it in exactly one place β a durable base plus environment overrides β then validate the resolved object, not the source file.
public static ProducerConfig BuildProducerConfig(
IConfiguration config, bool isProduction)
{
// 1. Base for durability-critical producers, in every environment.
var cfg = new ProducerConfig
{
BootstrapServers = config["Kafka:BootstrapServers"],
Acks = Acks.All, // wait for all in-sync replicas
EnableIdempotence = true, // per-session dedupe + ordering
MaxInFlight = 5, // legal alongside idempotence
MessageSendMaxRetries = 10, // message.send.max.retries
RetryBackoffMs = 200, // first backoff; doubles up to 1 s
RequestTimeoutMs = 5_000, // broker waits for replicas
SocketTimeoutMs = 10_000, // client waits for a response
LingerMs = 10, // small batch window
CompressionType = CompressionType.Lz4,
MessageTimeoutMs = 30_000, // total budget incl. retries
QueueBufferingMaxMessages = 100_000, // queue cap in records (default)
QueueBufferingMaxKbytes = 65_536, // queue cap in size: 64 MB, not 1 GB
};
// 2. Environment overrides win. Keys here are raw librdkafka keys.
foreach (var entry in config.GetSection("Kafka:Producer").GetChildren())
{
cfg.Set(entry.Key, entry.Value);
}
// 3. Fail fast on combinations that break durability or fail obscurely at Build().
if (cfg.EnableIdempotence == true && cfg.Acks != Acks.All)
{
throw new InvalidOperationException(
$"EnableIdempotence=true requires Acks=All; resolved Acks={cfg.Acks}.");
}
if (isProduction && cfg.Acks == Acks.None)
{
throw new InvalidOperationException(
"Acks=None is not permitted in production: fire-and-forget writes.");
}
if (isProduction && cfg.EnableIdempotence != true)
{
throw new InvalidOperationException(
"EnableIdempotence must be true in production.");
}
return cfg;
}
Three things worth naming. cfg.Set(key, value) accepts the raw librdkafka key, so an override section can say "linger.ms": "25" and bypass the typed properties; LingerMs and MessageTimeoutMs are .NET-side conveniences over those same keys. One trap: the typed getters read one spelling only (Acks reads acks, MaxInFlight reads max.in.flight), so an override written with an alias such as request.required.acks slips past the checks in step 3. Second, librdkafka auto-adjusts acks, max.in.flight.requests.per.connection, and retries when enable.idempotence=true β but only for settings you did not set explicitly; explicit incompatible values fail at Build(). Your startup check exists to surface that failure earlier, in your own words, with the resolved values attached. Third, SASL credentials and TLS key passwords come from environment variables or a secret store, never from a source-controlled file.
Quick check β does validation catch it? A production deployment resolves to Acks=None, EnableIdempotence=true, BootstrapServers=kafka:9092, LingerMs=5. Which rules does it break, which exception do you actually see, and what would the service have done without the checks?
Check your answer
It breaks two rules: the idempotence/acks rule (idempotence requires All) and the production Acks=None rule. You see only the first, because the validator throws at the first failed check; fix that one and the next run tells you about the other. Without the checks, librdkafka rejects the pair at Build() with a message that names neither the file nor the override that caused it β or, worse, if someone later dropped idempotence while leaving Acks=None, every write would be fire-and-forget and a broker failover could drop records with no error returned to your caller.
Log the effective config, then probe the broker
The resolved object can be enumerated, which makes redacted logging straightforward β provided the list of sensitive keys covers the security settings you actually use:
private static readonly string[] SensitiveKeys =
{ "sasl.password", "sasl.oauthbearer.client.secret",
"ssl.key.password", "ssl.key.pem", "ssl.keystore.password" };
static string Redact(string key, string value) =>
SensitiveKeys.Contains(key, StringComparer.OrdinalIgnoreCase) ? "***" : value;
foreach (var kv in cfg) // the ProducerConfig enumerates its key/value pairs
{
_logger.LogInformation("Kafka producer config: {Key}={Value}",
kv.Key, Redact(kv.Key, kv.Value));
}
Logging the effective values is what catches a silent misconfiguration such as Acks=None reaching production, or an idempotence flag disabled by an override nobody reviewed. Then prove the producer can actually write β a metadata fetch only tells you a broker answered:
public static async Task VerifyKafkaWritableAsync(
IProducer<string, string> producer, string canaryTopic, ILogger logger)
{
var result = await producer.ProduceAsync(
canaryTopic,
new Message<string, string> { Key = "startup", Value = "probe" });
logger.LogInformation(
"Kafka write probe OK: {Topic} [{Partition}] @ offset {Offset}",
result.Topic, result.Partition.Value, result.Offset.Value);
}
Publishing a probe verifies connectivity and write authorization in the same call. A configuration review that stops at the producer is incomplete: check the topic side too, because a producer with Acks.All on a topic with min.insync.replicas=1 is durable only until that lone replica dies β the interaction covered in "Acks, Retries, and the Real Meaning of 'Delivered'".
Shutdown discipline: bounded flush before dispose
Shutdown signal
β
1. stop intake (no new work enters the pipeline)
β
2. await in-flight ProduceAsync tasks (bounded too)
β
3. Flush(TimeSpan) bounded, e.g. 10 seconds
β
4. remaining > 0 ? log as shutdown failure : continue
β
5. Dispose() does not flush; what is still queued is dropped
// _inFlight: the ProduceAsync tasks not yet completed (tracking elided)
public async Task StopAsync()
{
_acceptingWork = false; // 1. stop intake
try // 2. await in-flight sends
{
await Task.WhenAll(_inFlight).WaitAsync(TimeSpan.FromSeconds(10));
}
catch (Exception ex) // a failed or slow send
{ // must not skip the flush
_logger.LogError(ex, "Kafka shutdown: in-flight sends failed or timed out");
}
var remaining = _producer.Flush(TimeSpan.FromSeconds(10)); // 3. bounded
if (remaining > 0) // 4. evidence
{
_logger.LogError(
"Kafka shutdown flush left {Count} records unacknowledged", remaining);
}
_producer.Dispose(); // 5. dispose last
}
The order is the contract. Dispose() takes no timeout parameter and does not flush: whatever is still queued when it runs is dropped, with no exception and no delivery report. Flush is the only call that waits for the queue to drain, and the TimeSpan overload gives you a clock and a return value β roughly how many records are still waiting to be sent or acknowledged. Zero versus non-zero is your evidence that the drain succeeded or did not. Pick bounds whose sum fits inside your orchestrator's termination grace period; a drain longer than the grace period is a drain the platform kills, which is the same as no drain. "First Contact: ProducerConfig, ProducerBuilder, and the Life of One ProduceAsync Call" showed the dispose-without-flush loss for a one-shot script; in a long-running service it becomes a lifecycle requirement.
Combined diagnostic: seven settings, one page at 2 a.m.
A service publishes order events with Acks.All, EnableIdempotence=true, max.in.flight=5, linger.ms=0, compression.type=none, queue.buffering.max.kbytes=16384 (16 MB), and delivery.timeout.ms=120000. During a broker leader failover, throughput collapses: many ProduceAsync calls fail at once with a queue-full error, others throw after two minutes. The application retries the failures, so every message is eventually delivered, and a downstream consumer that is not idempotent sees a few duplicates. Before reading the diagnosis, decide which single setting explains the two-minute throw.
| Setting | Value | Verdict |
|---|---|---|
| acks | all | Correct for durability |
| enable.idempotence | true | Correct |
| max.in.flight | 5 | Legal with idempotence β not the cause |
| linger.ms | 0 | No deliberate wait; still batches what is queued β not the cause |
| compression.type | none | Largest requests while catching up |
| queue.buffering.max.kbytes | 16384 | 16 MB: fills fast, then sends fail |
| delivery.timeout.ms | 120000 | The two-minute wait |
Why throughput collapsed. While partitions have no leader, nothing sent to them is acknowledged, so records pile up in the local queue. linger.ms=0 is not the culprit: it only removes the deliberate wait, and a backlog is exactly when the producer batches most, because it ships whatever is already queued as one batch (see "Throughput vs. Latency: Batching, Compression, and Buffer Memory"). What hurt was the 16 MB queue limit: with nothing draining it, the queue fills quickly, every further ProduceAsync fails immediately with Local_QueueFull, and the callers have to back off and try again β backpressure, arriving as errors. With compression.type=none, the catch-up once a new leader appears also sends the largest possible requests.
Why two minutes. delivery.timeout.ms (an alias of message.timeout.ms, typed as MessageTimeoutMs in Confluent.Kafka) bounds the total time a record may spend in the producer including retries. When no leader can acknowledge, the producer keeps retrying until that budget expires and then the ProduceAsync task faults. 120 s is also longer than most deployment grace periods, so a rollout during a failover risks being killed mid-flush.
Why duplicates despite idempotence. Idempotence deduplicates the client's own retries within a single producer session β a producer ID plus per-partition sequence numbers β and cannot deduplicate across sessions: a restarted producer instance receives a new ID, so the broker has no way to recognize the earlier copy. The second path is application-level retry: during a failover the record is often appended and only the acknowledgement is lost, so your retry after a timed-out ProduceAsync is a new record with a new sequence number, and it appends a second copy. Both are outside the producer's control.
Revised configuration.
LingerMs = 20, // optional, not the fix: larger batches day to day
CompressionType = CompressionType.Lz4, // cheap once batching is on
QueueBufferingMaxKbytes = 131_072, // sized for the failover backlog: 128 MB, not 16
QueueBufferingMaxMessages = 300_000, // or the default 100000 records binds first
MessageTimeoutMs = 30_000, // fail within 30 s, inside a typical grace period
// Acks=All, EnableIdempotence=true, MaxInFlight=5 stay as they are
One consumer-side guard. Deduplicate on a business key β a processed-message table, a unique constraint on order ID, or an upsert. That guard belongs to the consumer; no producer setting supplies it. Also review what sits next to these settings: transactional.id (the setting that turns on Kafka transactions) converts this into a transactional producer with a different configuration family and should not be switched on casually, and topic settings such as min.insync.replicas, max.message.bytes and topic-level compression can invalidate producer expectations.
Your turn: an audit-event producer review
Write a one-page producer configuration review for a new service that publishes audit events at moderate volume (a few hundred per second), must not lose events during a single broker failure, and must keep p99 produce latency under 50 ms. For each value state the reason, the trade-off you accept, and name the one metric you would alert on. Also state what you must verify on the topic side. Attempt it before opening the solution.
Check your answer
Acks = Acks.All, // no loss on a single broker failure
EnableIdempotence = true, // dedupe + ordering within a session
MaxInFlight = 5, // pipelining, order still preserved
LingerMs = 5, // the default; 5 ms fits a 50 ms p99
CompressionType = CompressionType.None, // 1-2 records per batch: nothing to gain
MessageSendMaxRetries = 5,
RetryBackoffMs = 100, // then 200, 400, 800, 1000: 2.5 s of pauses
RequestTimeoutMs = 3_000, // broker waits for replicas
SocketTimeoutMs = 5_000, // client waits for a response
MessageTimeoutMs = 20_000, // whole record; under a 30 s grace period
// queue limits: the defaults (100000 records, 1 GB) are ample at this volume
Topic side: replication.factor=3 and min.insync.replicas=2. Without that, Acks.All can be satisfied by a single in-sync replica and the no-loss requirement is unmet regardless of producer settings β a producer-only change cannot fix a topic with one in-sync replica.
Trade-offs: Acks.All with min.insync.replicas=2 favors durability over availability, so produces fail rather than risk loss when the ISR shrinks below two; linger.ms=5 adds a few milliseconds to p99; compression stays off because a few hundred events per second put only one or two records in a 5 ms batch, so there is nothing for a codec to work on. The alert: p99 produce latency, because the requirement is latency-bounded, with delivery-timeout errors as the secondary signal.
Self-assess against this checklist: Acks and min.insync.replicas aligned; idempotence scope understood (session-bound, not a cross-restart guarantee); retry and timeout budget stated and ordered (RequestTimeoutMs below SocketTimeoutMs, both well below MessageTimeoutMs, and MessageTimeoutMs inside the shutdown grace period); batching and compression budget stated as added latency; shutdown flush bounded and shorter than the grace period; effective config logged with secrets redacted and validated at startup.
Reading a correct answer is not the same as producing one β write the review for your own service's numbers before you consider the producer chapter finished. Producer configuration is the last mile of the Confluent.Kafka baseline, and it feeds directly into the Akka.Streams.Kafka pipelines and consumer-side idempotency that close the roadmap: abstractions do not delete distributed systems, they only put a nicer suit on them.