From 78985b47a1662927336bf6c33c9d4d2e15b7a274 Mon Sep 17 00:00:00 2001 From: xuexiaoyue Date: Sun, 3 Apr 2022 15:24:10 +0800 Subject: [PATCH 1/2] fix comparator in TransactionManager --- .../internals/TransactionManager.java | 10 +++- .../internals/TransactionManagerTest.java | 55 +++++++++++++++++++ 2 files changed, 63 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java index 4362ad68af3dd..df457c3471b50 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java @@ -184,16 +184,22 @@ private static class TopicPartitionEntry { // responses which are due to the retention period elapsing, and those which are due to actual lost data. private long lastAckedOffset; + private final Comparator producerBatchComparator = (b1, b2) -> { + if (b1.baseSequence() < b2.baseSequence()) return -1; + else if (b1.baseSequence() > b2.baseSequence()) return 1; + else return b1.equals(b2) ? 0 : 1; + }; + TopicPartitionEntry() { this.producerIdAndEpoch = ProducerIdAndEpoch.NONE; this.nextSequence = 0; this.lastAckedSequence = NO_LAST_ACKED_SEQUENCE_NUMBER; this.lastAckedOffset = ProduceResponse.INVALID_OFFSET; - this.inflightBatchesBySequence = new TreeSet<>(Comparator.comparingInt(ProducerBatch::baseSequence)); + this.inflightBatchesBySequence = new TreeSet<>(producerBatchComparator); } void resetSequenceNumbers(Consumer resetSequence) { - TreeSet newInflights = new TreeSet<>(Comparator.comparingInt(ProducerBatch::baseSequence)); + TreeSet newInflights = new TreeSet<>(producerBatchComparator); for (ProducerBatch inflightBatch : inflightBatchesBySequence) { resetSequence.accept(inflightBatch); newInflights.add(inflightBatch); diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java index 4227db5e61e62..6a7c314d8f267 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java @@ -674,6 +674,61 @@ public void testBatchCompletedAfterProducerReset() { assertNull(transactionManager.nextBatchBySequence(tp0)); } + @Test + public void testDuplicateSequenceAfterProducerReset() throws Exception { + initializeTransactionManager(Optional.empty()); + initializeIdempotentProducerId(producerId, epoch); + + Metrics metrics = new Metrics(time); + + RecordAccumulator accumulator = new RecordAccumulator(logContext, 16 * 1024, CompressionType.NONE, 0, 0L, + 15000, metrics, "", time, apiVersions, transactionManager, + new BufferPool(1024 * 1024, 16 * 1024, metrics, time, "")); + + Sender sender = new Sender(logContext, this.client, this.metadata, accumulator, false, + MAX_REQUEST_SIZE, ACKS_ALL, MAX_RETRIES, new SenderMetricsRegistry(metrics), this.time, 10000, + 0, transactionManager, apiVersions); + + assertEquals(0, transactionManager.sequenceNumber(tp0).intValue()); + + Future responseFuture1 = accumulator.append(tp0, time.milliseconds(), "1".getBytes(), "1".getBytes(), Record.EMPTY_HEADERS, + null, MAX_BLOCK_TIMEOUT, false, time.milliseconds()).future; + sender.runOnce(); + assertEquals(1, transactionManager.sequenceNumber(tp0).intValue()); + + time.sleep(10000); // request time out + sender.runOnce(); + assertEquals(0, client.inFlightRequestCount()); + assertTrue(transactionManager.hasInflightBatches(tp0)); + assertEquals(1, transactionManager.sequenceNumber(tp0).intValue()); + sender.runOnce(); // retry + assertEquals(1, client.inFlightRequestCount()); + assertTrue(transactionManager.hasInflightBatches(tp0)); + assertEquals(1, transactionManager.sequenceNumber(tp0).intValue()); + + time.sleep(5000); // delivery time out + sender.runOnce(); // expired in accumulator + assertFalse(transactionManager.hasInFlightRequest()); + assertEquals(1, client.inFlightRequestCount()); // not reaching request timeout, so still in flight + + sender.runOnce(); // bump the epoch + assertEquals(epoch + 1, transactionManager.producerIdAndEpoch().epoch); + assertEquals(0, transactionManager.sequenceNumber(tp0).intValue()); + + Future responseFuture2 = accumulator.append(tp0, time.milliseconds(), "2".getBytes(), "2".getBytes(), Record.EMPTY_HEADERS, + null, MAX_BLOCK_TIMEOUT, false, time.milliseconds()).future; + sender.runOnce(); + sender.runOnce(); + assertEquals(0, transactionManager.firstInFlightSequence(tp0)); + assertEquals(1, transactionManager.sequenceNumber(tp0).intValue()); + + time.sleep(5000); // request time out again + sender.runOnce(); + assertTrue(transactionManager.hasInflightBatches(tp0)); // the latter batch failed and retried + assertTrue(responseFuture1.isDone()); + assertFalse(responseFuture2.isDone()); + } + private ProducerBatch writeIdempotentBatchWithValue(TransactionManager manager, TopicPartition tp, String value) { From b4d2e86469d0ee8dfb127d80d2f998b9d80eef5e Mon Sep 17 00:00:00 2001 From: xuexiaoyue Date: Tue, 5 Apr 2022 10:36:41 +0800 Subject: [PATCH 2/2] fix review comments --- .../internals/TransactionManager.java | 6 +++--- .../internals/TransactionManagerTest.java | 19 +++++++++++++------ 2 files changed, 16 insertions(+), 9 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java index df457c3471b50..de93bcaedf1eb 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java @@ -184,7 +184,7 @@ private static class TopicPartitionEntry { // responses which are due to the retention period elapsing, and those which are due to actual lost data. private long lastAckedOffset; - private final Comparator producerBatchComparator = (b1, b2) -> { + private static final Comparator PRODUCER_BATCH_COMPARATOR = (b1, b2) -> { if (b1.baseSequence() < b2.baseSequence()) return -1; else if (b1.baseSequence() > b2.baseSequence()) return 1; else return b1.equals(b2) ? 0 : 1; @@ -195,11 +195,11 @@ private static class TopicPartitionEntry { this.nextSequence = 0; this.lastAckedSequence = NO_LAST_ACKED_SEQUENCE_NUMBER; this.lastAckedOffset = ProduceResponse.INVALID_OFFSET; - this.inflightBatchesBySequence = new TreeSet<>(producerBatchComparator); + this.inflightBatchesBySequence = new TreeSet<>(PRODUCER_BATCH_COMPARATOR); } void resetSequenceNumbers(Consumer resetSequence) { - TreeSet newInflights = new TreeSet<>(producerBatchComparator); + TreeSet newInflights = new TreeSet<>(PRODUCER_BATCH_COMPARATOR); for (ProducerBatch inflightBatch : inflightBatchesBySequence) { resetSequence.accept(inflightBatch); newInflights.add(inflightBatch); diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java index 6a7c314d8f267..64be3aeaf47b2 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java @@ -680,13 +680,15 @@ public void testDuplicateSequenceAfterProducerReset() throws Exception { initializeIdempotentProducerId(producerId, epoch); Metrics metrics = new Metrics(time); + final int requestTimeout = 10000; + final int deliveryTimeout = 15000; RecordAccumulator accumulator = new RecordAccumulator(logContext, 16 * 1024, CompressionType.NONE, 0, 0L, - 15000, metrics, "", time, apiVersions, transactionManager, + deliveryTimeout, metrics, "", time, apiVersions, transactionManager, new BufferPool(1024 * 1024, 16 * 1024, metrics, time, "")); Sender sender = new Sender(logContext, this.client, this.metadata, accumulator, false, - MAX_REQUEST_SIZE, ACKS_ALL, MAX_RETRIES, new SenderMetricsRegistry(metrics), this.time, 10000, + MAX_REQUEST_SIZE, ACKS_ALL, MAX_RETRIES, new SenderMetricsRegistry(metrics), this.time, requestTimeout, 0, transactionManager, apiVersions); assertEquals(0, transactionManager.sequenceNumber(tp0).intValue()); @@ -696,7 +698,7 @@ MAX_REQUEST_SIZE, ACKS_ALL, MAX_RETRIES, new SenderMetricsRegistry(metrics), thi sender.runOnce(); assertEquals(1, transactionManager.sequenceNumber(tp0).intValue()); - time.sleep(10000); // request time out + time.sleep(requestTimeout); sender.runOnce(); assertEquals(0, client.inFlightRequestCount()); assertTrue(transactionManager.hasInflightBatches(tp0)); @@ -707,9 +709,15 @@ MAX_REQUEST_SIZE, ACKS_ALL, MAX_RETRIES, new SenderMetricsRegistry(metrics), thi assertEquals(1, transactionManager.sequenceNumber(tp0).intValue()); time.sleep(5000); // delivery time out - sender.runOnce(); // expired in accumulator + sender.runOnce(); + + // The retried request will remain inflight until the request timeout + // is reached even though the delivery timeout has expired and the + // future has completed exceptionally. + assertTrue(responseFuture1.isDone()); + TestUtils.assertFutureThrows(responseFuture1, TimeoutException.class); assertFalse(transactionManager.hasInFlightRequest()); - assertEquals(1, client.inFlightRequestCount()); // not reaching request timeout, so still in flight + assertEquals(1, client.inFlightRequestCount()); sender.runOnce(); // bump the epoch assertEquals(epoch + 1, transactionManager.producerIdAndEpoch().epoch); @@ -725,7 +733,6 @@ MAX_REQUEST_SIZE, ACKS_ALL, MAX_RETRIES, new SenderMetricsRegistry(metrics), thi time.sleep(5000); // request time out again sender.runOnce(); assertTrue(transactionManager.hasInflightBatches(tp0)); // the latter batch failed and retried - assertTrue(responseFuture1.isDone()); assertFalse(responseFuture2.isDone()); }