Kafka的min.insync.replicas配置未生效,如何复现预期异常?
min.insync.replicas Behavior in Kafka 0.10.0.0 Great question—let's break down exactly what's happening with your setup and how to get the expected NotEnoughReplicasException to surface properly.
Why Your Console Producer Didn't Throw the Exception
The default behavior of Kafka's console producer (in 0.10.0.0) uses acks=1, not acks=all. When acks=1, the producer only waits for the leader broker to acknowledge the write—it doesn't care about the number of in-sync replicas (ISR). That's why your message went through successfully even when the ISR shrank to 1 (only the leader was alive).
To fix this, you need to explicitly set acks=all when starting the console producer, which forces it to wait for all replicas in the ISR to confirm the write. This is the trigger that makes the broker check if the ISR size meets your min.insync.replicas requirement.
Why Your Java Producer Didn't Catch the Exception
There are two likely culprits here:
- You weren't waiting for the send result: If you used asynchronous sends (calling
send()without getting theFutureresult), the exception would only be logged internally by the producer or passed to a callback—your main thread wouldn't see it. - Retries were masking the exception: If your producer had
retries > 0(the default might be non-zero in 0.10.0.0), it would automatically retry the send, delaying or hiding the exception.
Step-by-Step to Reproduce the Expected Exception
Let's adjust your experiment to trigger the exception consistently:
1. Prep (Same as Your Original Setup)
- Keep your 3 brokers (0,1,2) running initially.
- Create the topic with the correct configs:
sudo ./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 3 --partitions 1 --topic insync --config min.insync.replicas=2 - Stop brokers 0 and 2, confirm the ISR for
insyncis only[1]via:sudo ./kafka-topics.sh --describe --zookeeper localhost:2181 --topic insync
2. Console Producer: Trigger the Exception
Start the console producer with acks=all:
sudo ./kafka-console-producer.sh --broker-list localhost:9092 --topic insync --property acks=all
When you type a message and hit enter, you'll immediately see the NotEnoughReplicasException in the console output.
3. Java Producer: Catch the Exception
Adjust your producer config and code to ensure you capture the exception:
Producer Config
Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("acks", "all"); // Critical: Wait for all ISR replicas to confirm props.put("retries", 0); // Disable retries to see the exception immediately props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Option 1: Synchronous Send (Catch Exception Directly)
KafkaProducer<String, String> producer = new KafkaProducer<>(props); ProducerRecord<String, String> record = new ProducerRecord<>("insync", "test", "message"); try { // Wait for the send to complete and retrieve any exception producer.send(record).get(); } catch (InterruptedException | ExecutionException e) { // Unwrap the underlying Kafka exception if (e.getCause() instanceof NotEnoughReplicasException) { System.err.println("Caught expected NotEnoughReplicasException!"); e.getCause().printStackTrace(); } else { e.printStackTrace(); } } finally { producer.close(); }
Option 2: Asynchronous Send with Callback
producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception != null) { if (exception instanceof NotEnoughReplicasException) { System.err.println("Caught expected NotEnoughReplicasException in callback!"); } exception.printStackTrace(); } } }); // Keep the producer running long enough to receive the callback Thread.sleep(5000); producer.close();
Key Notes for Kafka 0.10.0.0
min.insync.replicasonly enforces when the producer usesacks=all(or-1, which is equivalent). Without this, the broker doesn't check the ISR size against the config.- Console producers don't inherit topic-level configs for
acks—you have to set it explicitly via the--propertyflag. - Asynchronous sends won't throw exceptions in your main thread unless you explicitly check the
Futureresult or use a callback.
内容的提问来源于stack exchange,提问作者Arun Gowda

