From 8037122c79984b57e3c74ccb53b52d0e277d95eb Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Wed, 27 Oct 2021 16:20:49 -0700 Subject: [PATCH 01/10] KAFKA-13412; Ensure initTransactions() safe for retry after timeout --- .../internals/TransactionManager.java | 16 +++++--- .../internals/TransactionalRequestResult.java | 10 +++++ .../clients/producer/KafkaProducerTest.java | 41 +++++++++++++++++++ .../internals/TransactionManagerTest.java | 10 ++++- 4 files changed, 71 insertions(+), 6 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 2de31a03c5694..bee0a72078800 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 @@ -241,6 +241,7 @@ void resetSequenceNumbers(Consumer resetSequence) { private Node consumerGroupCoordinator; private boolean coordinatorSupportsBumpingEpoch; + private volatile State priorState = null; private volatile State currentState = State.UNINITIALIZED; private volatile RuntimeException lastError = null; private volatile ProducerIdAndEpoch producerIdAndEpoch; @@ -1090,6 +1091,7 @@ private void transitionTo(State target, RuntimeException error) { else log.debug("Transition from state {} to {}", currentState, target); + priorState = currentState; currentState = target; } @@ -1184,15 +1186,19 @@ private TxnOffsetCommitHandler txnOffsetCommitHandler(TransactionalRequestResult } private TransactionalRequestResult handleCachedTransactionRequestResult( - Supplier transactionalRequestResultSupplier, - State targetState) { + Supplier transactionalRequestResultSupplier, + State transientState + ) { ensureTransactional(); - if (pendingResult != null && currentState == targetState) { + if (pendingResult != null && !pendingResult.isAcked()) { TransactionalRequestResult result = pendingResult; - if (result.isCompleted()) + if (result.isCompleted() && priorState == transientState) { pendingResult = null; - return result; + return result; + } else if (currentState == transientState) { + return result; + } } pendingResult = transactionalRequestResultSupplier.get(); diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java index d442b188d4520..e94de39ec9058 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java @@ -29,6 +29,7 @@ public final class TransactionalRequestResult { private final CountDownLatch latch; private volatile RuntimeException error = null; private final String operation; + private volatile boolean isAcked = false; public TransactionalRequestResult(String operation) { this(new CountDownLatch(1), operation); @@ -60,6 +61,8 @@ public void await() { } } + isAcked = true; + if (!isSuccessful()) throw error(); } @@ -68,10 +71,13 @@ public void await(long timeout, TimeUnit unit) { try { boolean success = latch.await(timeout, unit); if (!isSuccessful()) { + isAcked = true; throw error(); } if (!success) { throw new TimeoutException("Timeout expired after " + timeout + " " + unit.name().toLowerCase(Locale.ROOT) + " while awaiting " + operation); + } else { + isAcked = true; } } catch (InterruptedException e) { throw new InterruptException("Received interrupt while awaiting " + operation, e); @@ -90,4 +96,8 @@ public boolean isCompleted() { return latch.getCount() == 0L; } + public boolean isAcked() { + return isAcked; + } + } diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java index 1e45a58b98efe..2abd376910cb3 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java @@ -903,6 +903,47 @@ public void testPartitionsForWithNullTopic() { } } + @Test + public void testInitTransactionsResponseAfterTimeout() throws Exception { + int maxBlockMs = 500; + + Map configs = new HashMap<>(); + configs.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "bad-transaction"); + configs.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, maxBlockMs); + configs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9000"); + + Time time = new MockTime(); + MetadataResponse initialUpdateResponse = RequestTestUtils.metadataUpdateWith(1, singletonMap("topic", 1)); + ProducerMetadata metadata = newMetadata(0, Long.MAX_VALUE); + metadata.updateWithCurrentRequestVersion(initialUpdateResponse, false, time.milliseconds()); + + MockClient client = new MockClient(time, metadata); + + ExecutorService executor = Executors.newFixedThreadPool(1); + + Producer producer = kafkaProducer(configs, new StringSerializer(), + new StringSerializer(), metadata, client, null, time); + try { + client.prepareResponse( + request -> request instanceof FindCoordinatorRequest && + ((FindCoordinatorRequest) request).data().keyType() == FindCoordinatorRequest.CoordinatorType.TRANSACTION.id(), + FindCoordinatorResponse.prepareResponse(Errors.NONE, "bad-transaction", host1)); + + Future future = executor.submit(producer::initTransactions); + TestUtils.waitForCondition(client::hasInFlightRequests, "blah blah"); + + time.sleep(maxBlockMs); + TestUtils.assertFutureThrows(future, TimeoutException.class); + + client.respond(initProducerIdResponse(1L, (short) 5, Errors.NONE)); + + Thread.sleep(1000); + producer.initTransactions(); + } finally { + producer.close(Duration.ZERO); + } + } + @Test public void testInitTransactionTimeout() { Map configs = new HashMap<>(); 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 6c1e2fdf1e785..b647e3fe77dc9 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 @@ -1661,6 +1661,10 @@ public void testInvalidProducerEpochConvertToProducerFencedInEndTxn() throws Int runUntil(commitResult::isCompleted); runUntil(responseFuture::isDone); + assertThrows(KafkaException.class, commitResult::await); + assertFalse(commitResult.isSuccessful()); + assertTrue(commitResult.isAcked()); + // make sure the exception was thrown directly from the follow-up calls. assertThrows(KafkaException.class, () -> transactionManager.beginTransaction()); assertThrows(KafkaException.class, () -> transactionManager.beginCommit()); @@ -3522,13 +3526,17 @@ private void doInitTransactions() { } private void doInitTransactions(long producerId, short epoch) { - transactionManager.initializeTransactions(); + TransactionalRequestResult result = transactionManager.initializeTransactions(); prepareFindCoordinatorResponse(Errors.NONE, false, CoordinatorType.TRANSACTION, transactionalId); runUntil(() -> transactionManager.coordinator(CoordinatorType.TRANSACTION) != null); assertEquals(brokerNode, transactionManager.coordinator(CoordinatorType.TRANSACTION)); prepareInitPidResponse(Errors.NONE, false, producerId, epoch); runUntil(transactionManager::hasProducerId); + + result.await(); + assertTrue(result.isSuccessful()); + assertTrue(result.isAcked()); } private void assertAbortableError(Class cause) { From a70563a9a020ff917d511944212e78ed499b3917 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Tue, 16 Nov 2021 18:06:30 -0800 Subject: [PATCH 02/10] Attempt to simplify pending transition bookkeeping --- .../internals/TransactionManager.java | 49 ++++++++++++------- .../internals/TransactionalRequestResult.java | 34 ++++--------- .../producer/internals/SenderTest.java | 3 +- .../internals/TransactionManagerTest.java | 27 +++++----- 4 files changed, 57 insertions(+), 56 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 bee0a72078800..085e1d67b5e40 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 @@ -225,7 +225,7 @@ void resetSequenceNumbers(Consumer resetSequence) { private final Set newPartitionsInTransaction; private final Set pendingPartitionsInTransaction; private final Set partitionsInTransaction; - private TransactionalRequestResult pendingResult; + private PendingStateTransition pendingTransition; // This is used by the TxnRequestHandlers to control how long to back off before a given request is retried. // For instance, this value is lowered by the AddPartitionsToTxnHandler when it receives a CONCURRENT_TRANSACTIONS @@ -241,7 +241,6 @@ void resetSequenceNumbers(Consumer resetSequence) { private Node consumerGroupCoordinator; private boolean coordinatorSupportsBumpingEpoch; - private volatile State priorState = null; private volatile State currentState = State.UNINITIALIZED; private volatile RuntimeException lastError = null; private volatile ProducerIdAndEpoch producerIdAndEpoch; @@ -501,8 +500,8 @@ synchronized void transitionToFatalError(RuntimeException exception) { log.info("Transiting to fatal error state due to {}", exception.toString()); transitionTo(State.FATAL_ERROR, exception); - if (pendingResult != null) { - pendingResult.fail(exception); + if (pendingTransition != null) { + pendingTransition.result.fail(exception); } } @@ -920,8 +919,8 @@ synchronized void close() { KafkaException shutdownException = new KafkaException("The producer closed forcefully"); pendingRequests.forEach(handler -> handler.fatalError(shutdownException)); - if (pendingResult != null) { - pendingResult.fail(shutdownException); + if (pendingTransition != null) { + pendingTransition.result.fail(shutdownException); } } @@ -1091,7 +1090,6 @@ private void transitionTo(State target, RuntimeException error) { else log.debug("Transition from state {} to {}", currentState, target); - priorState = currentState; currentState = target; } @@ -1187,22 +1185,24 @@ private TxnOffsetCommitHandler txnOffsetCommitHandler(TransactionalRequestResult private TransactionalRequestResult handleCachedTransactionRequestResult( Supplier transactionalRequestResultSupplier, - State transientState + State nextState ) { ensureTransactional(); - if (pendingResult != null && !pendingResult.isAcked()) { - TransactionalRequestResult result = pendingResult; - if (result.isCompleted() && priorState == transientState) { - pendingResult = null; - return result; - } else if (currentState == transientState) { - return result; + if (pendingTransition != null) { + if (pendingTransition.result.isAcked()) { + pendingTransition = null; + } else if (nextState != pendingTransition.state) { + throw new KafkaException("Unexpected transition to " + nextState + + " while awaiting transition from " + pendingTransition.state); + } else { + return pendingTransition.result; } } - pendingResult = transactionalRequestResultSupplier.get(); - return pendingResult; + TransactionalRequestResult result = transactionalRequestResultSupplier.get(); + pendingTransition = new PendingStateTransition(result, nextState); + return result; } // package-private for testing @@ -1768,4 +1768,19 @@ private boolean isFatalException(Errors error) { || error == Errors.PRODUCER_FENCED || error == Errors.UNSUPPORTED_FOR_MESSAGE_FORMAT; } + + private static final class PendingStateTransition { + private final TransactionalRequestResult result; + private final State state; + + private PendingStateTransition( + TransactionalRequestResult result, + State state + ) { + this.result = result; + this.state = state; + } + } + + } diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java index e94de39ec9058..6739da815273e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionalRequestResult.java @@ -20,12 +20,10 @@ import org.apache.kafka.common.errors.InterruptException; import org.apache.kafka.common.errors.TimeoutException; -import java.util.Locale; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; public final class TransactionalRequestResult { - private final CountDownLatch latch; private volatile RuntimeException error = null; private final String operation; @@ -50,34 +48,20 @@ public void done() { } public void await() { - boolean completed = false; - - while (!completed) { - try { - latch.await(); - completed = true; - } catch (InterruptedException e) { - // Keep waiting until done, we have no other option for these transactional requests. - } - } - - isAcked = true; - - if (!isSuccessful()) - throw error(); + this.await(Long.MAX_VALUE, TimeUnit.MILLISECONDS); } public void await(long timeout, TimeUnit unit) { try { boolean success = latch.await(timeout, unit); - if (!isSuccessful()) { - isAcked = true; - throw error(); - } if (!success) { - throw new TimeoutException("Timeout expired after " + timeout + " " + unit.name().toLowerCase(Locale.ROOT) + " while awaiting " + operation); - } else { - isAcked = true; + throw new TimeoutException("Timeout expired after " + unit.toMillis(timeout) + + "ms while awaiting " + operation); + } + + isAcked = true; + if (error != null) { + throw error; } } catch (InterruptedException e) { throw new InterruptException("Received interrupt while awaiting " + operation, e); @@ -89,7 +73,7 @@ public RuntimeException error() { } public boolean isSuccessful() { - return error == null; + return isCompleted() && error == null; } public boolean isCompleted() { diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java index 34e4f180e7866..10baa47997680 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java @@ -3191,7 +3191,7 @@ private InitProducerIdResponse initProducerIdResponse(long producerId, short pro } private void doInitTransactions(TransactionManager transactionManager, ProducerIdAndEpoch producerIdAndEpoch) { - transactionManager.initializeTransactions(); + TransactionalRequestResult result = transactionManager.initializeTransactions(); prepareFindCoordinatorResponse(Errors.NONE, transactionManager.transactionalId()); sender.runOnce(); sender.runOnce(); @@ -3199,6 +3199,7 @@ private void doInitTransactions(TransactionManager transactionManager, ProducerI prepareInitProducerResponse(Errors.NONE, producerIdAndEpoch.producerId, producerIdAndEpoch.epoch); sender.runOnce(); assertTrue(transactionManager.hasProducerId()); + result.await(); } private void prepareFindCoordinatorResponse(Errors error, String txnid) { 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 b647e3fe77dc9..08932db3f1688 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 @@ -1075,7 +1075,7 @@ public void testTransactionalIdAuthorizationFailureInFindCoordinator() { assertTrue(transactionManager.hasFatalError()); assertTrue(transactionManager.lastError() instanceof TransactionalIdAuthorizationException); assertFalse(initPidResult.isSuccessful()); - assertTrue(initPidResult.error() instanceof TransactionalIdAuthorizationException); + assertThrows(TransactionalIdAuthorizationException.class, initPidResult::await); assertFatalError(TransactionalIdAuthorizationException.class); } @@ -1090,8 +1090,7 @@ public void testTransactionalIdAuthorizationFailureInInitProducerId() { runUntil(transactionManager::hasError); assertTrue(initPidResult.isCompleted()); assertFalse(initPidResult.isSuccessful()); - assertTrue(initPidResult.error() instanceof TransactionalIdAuthorizationException); - + assertThrows(TransactionalIdAuthorizationException.class, initPidResult::await); assertFatalError(TransactionalIdAuthorizationException.class); } @@ -1296,7 +1295,7 @@ public void testRecoveryFromAbortableErrorTransactionNotStarted() throws Excepti runUntil(() -> !client.hasPendingResponses()); assertTrue(transactionManager.hasAbortableError()); - transactionManager.beginAbort(); + TransactionalRequestResult abortResult = transactionManager.beginAbort(); runUntil(responseFuture::isDone); assertProduceFutureFailed(responseFuture); @@ -1304,6 +1303,8 @@ public void testRecoveryFromAbortableErrorTransactionNotStarted() throws Excepti runUntil(transactionManager::isReady); assertFalse(transactionManager.hasPartitionsToAdd()); assertFalse(accumulator.hasIncomplete()); + assertTrue(abortResult.isSuccessful()); + abortResult.await(); // ensure we can now start a new transaction @@ -1351,13 +1352,15 @@ public void testRecoveryFromAbortableErrorTransactionStarted() throws Exception assertFalse(unauthorizedTopicProduceFuture.isDone()); prepareEndTxnResponse(Errors.NONE, TransactionResult.ABORT, producerId, epoch); - transactionManager.beginAbort(); + TransactionalRequestResult result = transactionManager.beginAbort(); runUntil(transactionManager::isReady); // neither produce request has been sent, so they should both be failed immediately assertProduceFutureFailed(authorizedTopicProduceFuture); assertProduceFutureFailed(unauthorizedTopicProduceFuture); assertFalse(transactionManager.hasPartitionsToAdd()); assertFalse(accumulator.hasIncomplete()); + assertTrue(result.isSuccessful()); + result.await(); // ensure we can now start a new transaction @@ -1417,12 +1420,14 @@ public void testRecoveryFromAbortableErrorProduceRequestInRetry() throws Excepti assertTrue(authorizedTopicProduceFuture.isDone()); prepareEndTxnResponse(Errors.NONE, TransactionResult.ABORT, producerId, epoch); - transactionManager.beginAbort(); + TransactionalRequestResult abortResult = transactionManager.beginAbort(); runUntil(transactionManager::isReady); // neither produce request has been sent, so they should both be failed immediately assertTrue(transactionManager.isReady()); assertFalse(transactionManager.hasPartitionsToAdd()); assertFalse(accumulator.hasIncomplete()); + assertTrue(abortResult.isSuccessful()); + abortResult.await(); // ensure we can now start a new transaction @@ -1563,7 +1568,7 @@ private void verifyProducerFencedForInitProducerId(Errors error) { runUntil(transactionManager::hasError); - assertEquals(ProducerFencedException.class, result.error().getClass()); + assertThrows(ProducerFencedException.class, result::await); assertThrows(ProducerFencedException.class, () -> transactionManager.beginTransaction()); assertThrows(ProducerFencedException.class, () -> transactionManager.beginCommit()); @@ -2508,6 +2513,7 @@ public void testDropCommitOnBatchExpiry() throws InterruptedException { } runUntil(commitResult::isCompleted); // the commit shouldn't be completed without being sent since the produce request failed. assertFalse(commitResult.isSuccessful()); // the commit shouldn't succeed since the produce request failed. + assertThrows(TimeoutException.class, commitResult::await); assertTrue(transactionManager.hasAbortableError()); assertTrue(transactionManager.hasOngoingTransaction()); @@ -3288,12 +3294,7 @@ private void verifyCommitOrAbortTransactionRetriable(TransactionResult firstTran prepareEndTxnResponse(Errors.NONE, firstTransactionResult, producerId, epoch, true); runUntil(() -> !client.hasPendingResponses()); assertFalse(result.isCompleted()); - - try { - result.await(MAX_BLOCK_TIMEOUT, TimeUnit.MILLISECONDS); - fail("Should have raised TimeoutException"); - } catch (TimeoutException ignored) { - } + assertThrows(TimeoutException.class, () -> result.await(MAX_BLOCK_TIMEOUT, TimeUnit.MILLISECONDS)); prepareFindCoordinatorResponse(Errors.NONE, false, CoordinatorType.TRANSACTION, transactionalId); runUntil(() -> !client.hasPendingResponses()); From f67f377c95a05a27c3ac0701678f1adab20815e6 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 22 Nov 2021 15:22:22 -0800 Subject: [PATCH 03/10] Additional test cases and cleanup --- .../kafka/clients/producer/KafkaProducer.java | 8 +- .../internals/TransactionManager.java | 83 +++-- .../producer/internals/SenderTest.java | 29 +- .../internals/TransactionManagerTest.java | 345 ++++++++++-------- 4 files changed, 256 insertions(+), 209 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java index dbb908d4cbf80..a7f9389219add 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java @@ -964,9 +964,10 @@ private Future doSend(ProducerRecord record, Callback call // producer callback will make sure to call both 'callback' and interceptor callback Callback interceptCallback = new InterceptorCallback<>(callback, this.interceptors, tp); - if (transactionManager != null && transactionManager.isTransactional()) { - transactionManager.failIfNotReadyForSend(); + if (transactionManager != null) { + transactionManager.maybeAddPartition(tp); } + RecordAccumulator.RecordAppendResult result = accumulator.append(tp, timestamp, serializedKey, serializedValue, headers, interceptCallback, remainingWaitMs, true, nowMs); @@ -985,9 +986,6 @@ private Future doSend(ProducerRecord record, Callback call serializedValue, headers, interceptCallback, remainingWaitMs, false, nowMs); } - if (transactionManager != null && transactionManager.isTransactional()) - transactionManager.maybeAddPartitionToTransaction(tp); - if (result.batchIsFull || result.newBatchCreated) { log.trace("Waking up the sender since topic {} partition {} is either full or getting a new batch", record.topic(), partition); this.sender.wakeup(); 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 085e1d67b5e40..5f05fc92e4c44 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 @@ -347,11 +347,12 @@ synchronized TransactionalRequestResult initializeTransactions(ProducerIdAndEpoc isEpochBump); enqueueRequest(handler); return handler.result; - }, State.INITIALIZING); + }, State.INITIALIZING, "initTransactions"); } public synchronized void beginTransaction() { ensureTransactional(); + throwIfPendingState("beginTransaction"); maybeFailWithError(); transitionTo(State.IN_TRANSACTION); } @@ -361,7 +362,7 @@ public synchronized TransactionalRequestResult beginCommit() { maybeFailWithError(); transitionTo(State.COMMITTING_TRANSACTION); return beginCompletingTransaction(TransactionResult.COMMIT); - }, State.COMMITTING_TRANSACTION); + }, State.COMMITTING_TRANSACTION, "commitTransaction"); } public synchronized TransactionalRequestResult beginAbort() { @@ -373,7 +374,7 @@ public synchronized TransactionalRequestResult beginAbort() { // We're aborting the transaction, so there should be no need to add new partitions newPartitionsInTransaction.clear(); return beginCompletingTransaction(TransactionResult.ABORT); - }, State.ABORTING_TRANSACTION); + }, State.ABORTING_TRANSACTION, "abortTransaction"); } private TransactionalRequestResult beginCompletingTransaction(TransactionResult transactionResult) { @@ -404,10 +405,12 @@ private TransactionalRequestResult beginCompletingTransaction(TransactionResult public synchronized TransactionalRequestResult sendOffsetsToTransaction(final Map offsets, final ConsumerGroupMetadata groupMetadata) { ensureTransactional(); + throwIfPendingState("sendOffsetsToTransaction"); maybeFailWithError(); - if (currentState != State.IN_TRANSACTION) - throw new KafkaException("Cannot send offsets to transaction either because the producer is not in an " + - "active transaction"); + + if (currentState != State.IN_TRANSACTION) { + throw new KafkaException("Cannot send offsets to transaction either because in state " + currentState); + } log.debug("Begin adding offsets {} for consumer group {} to transaction", offsets, groupMetadata); AddOffsetsToTxnRequest.Builder builder = new AddOffsetsToTxnRequest.Builder( @@ -423,34 +426,31 @@ public synchronized TransactionalRequestResult sendOffsetsToTransaction(final Ma return handler.result; } - public synchronized void maybeAddPartitionToTransaction(TopicPartition topicPartition) { - if (isPartitionAdded(topicPartition) || isPartitionPendingAdd(topicPartition)) - return; + public synchronized void maybeAddPartition(TopicPartition topicPartition) { + maybeFailWithError(); + throwIfPendingState("send"); - log.debug("Begin adding new partition {} to transaction", topicPartition); - topicPartitionBookkeeper.addPartition(topicPartition); - newPartitionsInTransaction.add(topicPartition); + if (isTransactional()) { + if (!hasProducerId()) { + throw new KafkaException("Cannot add partition " + topicPartition + + "to transaction before completing a call to initTransactions"); + } else if (currentState != State.IN_TRANSACTION) { + throw new KafkaException("Cannot add partition " + topicPartition + + " to transaction while in state " + currentState); + } else if (isPartitionAdded(topicPartition) || isPartitionPendingAdd(topicPartition)) { + return; + } else { + log.debug("Begin adding new partition {} to transaction", topicPartition); + topicPartitionBookkeeper.addPartition(topicPartition); + newPartitionsInTransaction.add(topicPartition); + } + } } RuntimeException lastError() { return lastError; } - public synchronized void failIfNotReadyForSend() { - if (hasError()) - throw new KafkaException("Cannot perform send because at least one previous transactional or " + - "idempotent request has failed with errors.", lastError); - - if (isTransactional()) { - if (!hasProducerId()) - throw new IllegalStateException("Cannot perform a 'send' before completing a call to initTransactions " + - "when transactions are enabled."); - - if (currentState != State.IN_TRANSACTION) - throw new IllegalStateException("Cannot call send in state " + currentState); - } - } - synchronized boolean isSendToPartitionAllowed(TopicPartition tp) { if (hasFatalError()) return false; @@ -1183,9 +1183,26 @@ private TxnOffsetCommitHandler txnOffsetCommitHandler(TransactionalRequestResult return new TxnOffsetCommitHandler(result, builder); } + private void throwPendingTransitionError(String operation) { + throw new KafkaException("Cannot attempt operation `" + operation + "` " + + "because the previous call to `" + pendingTransition.operation + "` " + + "timed out and must be retried"); + } + + private void throwIfPendingState(String operation) { + if (pendingTransition != null) { + if (pendingTransition.result.isAcked()) { + pendingTransition = null; + } else { + throwPendingTransitionError(operation); + } + } + } + private TransactionalRequestResult handleCachedTransactionRequestResult( Supplier transactionalRequestResultSupplier, - State nextState + State nextState, + String operation ) { ensureTransactional(); @@ -1193,15 +1210,14 @@ private TransactionalRequestResult handleCachedTransactionRequestResult( if (pendingTransition.result.isAcked()) { pendingTransition = null; } else if (nextState != pendingTransition.state) { - throw new KafkaException("Unexpected transition to " + nextState + - " while awaiting transition from " + pendingTransition.state); + throwPendingTransitionError(operation); } else { return pendingTransition.result; } } TransactionalRequestResult result = transactionalRequestResultSupplier.get(); - pendingTransition = new PendingStateTransition(result, nextState); + pendingTransition = new PendingStateTransition(result, nextState, operation); return result; } @@ -1772,13 +1788,16 @@ private boolean isFatalException(Errors error) { private static final class PendingStateTransition { private final TransactionalRequestResult result; private final State state; + private final String operation; private PendingStateTransition( TransactionalRequestResult result, - State state + State state, + String operation ) { this.result = result; this.state = state; + this.operation = operation; } } diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java index 10baa47997680..60e9f06186255 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java @@ -1465,8 +1465,7 @@ public void testUnresolvedSequencesAreNotFatal() throws Exception { doInitTransactions(txnManager, producerIdAndEpoch); txnManager.beginTransaction(); - txnManager.failIfNotReadyForSend(); - txnManager.maybeAddPartitionToTransaction(tp0); + txnManager.maybeAddPartition(tp0); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp0, Errors.NONE))); sender.runOnce(); @@ -1751,7 +1750,8 @@ public void testTransactionalUnknownProducerHandlingWhenRetentionLimitReached() doInitTransactions(transactionManager, new ProducerIdAndEpoch(producerId, (short) 0)); assertTrue(transactionManager.hasProducerId()); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.beginTransaction(); + transactionManager.maybeAddPartition(tp0); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp0, Errors.NONE))); sender.runOnce(); // Receive AddPartitions response @@ -2307,8 +2307,7 @@ public void testTransactionalSplitBatchAndSend() throws Exception { doInitTransactions(txnManager, producerIdAndEpoch); txnManager.beginTransaction(); - txnManager.failIfNotReadyForSend(); - txnManager.maybeAddPartitionToTransaction(tp); + txnManager.maybeAddPartition(tp); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp, Errors.NONE))); sender.runOnce(); @@ -2654,8 +2653,7 @@ public void testTransactionalRequestsSentOnShutdown() { doInitTransactions(txnManager, producerIdAndEpoch); txnManager.beginTransaction(); - txnManager.failIfNotReadyForSend(); - txnManager.maybeAddPartitionToTransaction(tp); + txnManager.maybeAddPartition(tp); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp, Errors.NONE))); sender.runOnce(); sender.initiateClose(); @@ -2697,7 +2695,7 @@ public void testRecordsFlushedImmediatelyOnTransactionCompletion() throws Except // Now begin the commit and assert that the Produce request is sent immediately // without waiting for the linger. - txnManager.beginCommit(); + TransactionalRequestResult commitResult = txnManager.beginCommit(); runUntil(sender, client::hasInFlightRequests); // Respond to the produce request and wait for the EndTxn request to be sent. @@ -2708,6 +2706,9 @@ public void testRecordsFlushedImmediatelyOnTransactionCompletion() throws Except respondToEndTxn(Errors.NONE); runUntil(sender, txnManager::isReady); + assertTrue(commitResult.isSuccessful()); + commitResult.await(); + // Finally, we want to assert that the linger time is still effective // when the new transaction begins. txnManager.beginTransaction(); @@ -2772,7 +2773,7 @@ public void testAwaitPendingRecordsBeforeCommittingTransaction() throws Exceptio } private void addPartitionToTxn(Sender sender, TransactionManager txnManager, TopicPartition tp) { - txnManager.maybeAddPartitionToTransaction(tp); + txnManager.maybeAddPartition(tp); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp, Errors.NONE))); runUntil(sender, () -> txnManager.isPartitionAdded(tp)); assertFalse(txnManager.hasInFlightRequest()); @@ -2813,8 +2814,7 @@ public void testIncompleteTransactionAbortOnShutdown() { doInitTransactions(txnManager, producerIdAndEpoch); txnManager.beginTransaction(); - txnManager.failIfNotReadyForSend(); - txnManager.maybeAddPartitionToTransaction(tp); + txnManager.maybeAddPartition(tp); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp, Errors.NONE))); sender.runOnce(); sender.initiateClose(); @@ -2848,8 +2848,7 @@ public void testForceShutdownWithIncompleteTransaction() { doInitTransactions(txnManager, producerIdAndEpoch); txnManager.beginTransaction(); - txnManager.failIfNotReadyForSend(); - txnManager.maybeAddPartitionToTransaction(tp); + txnManager.maybeAddPartition(tp); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp, Errors.NONE))); sender.runOnce(); @@ -2874,7 +2873,7 @@ public void testTransactionAbortedExceptionOnAbortWithoutError() throws Interrup doInitTransactions(txnManager, producerIdAndEpoch); // Begin the transaction txnManager.beginTransaction(); - txnManager.maybeAddPartitionToTransaction(tp0); + txnManager.maybeAddPartition(tp0); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp0, Errors.NONE))); // Run it once so that the partition is added to the transaction. sender.runOnce(); @@ -2912,7 +2911,7 @@ public void testTooLargeBatchesAreSafelyRemoved() throws InterruptedException { doInitTransactions(txnManager, producerIdAndEpoch); txnManager.beginTransaction(); - txnManager.maybeAddPartitionToTransaction(tp0); + txnManager.maybeAddPartition(tp0); client.prepareResponse(new AddPartitionsToTxnResponse(0, Collections.singletonMap(tp0, Errors.NONE))); sender.runOnce(); 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 08932db3f1688..aef007dd14f70 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 @@ -103,6 +103,7 @@ import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -184,8 +185,7 @@ public void testSenderShutdownWithPendingTransactions() throws Exception { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); FutureRecordMetadata sendFuture = appendToAccumulator(tp0); prepareAddPartitionsToTxn(tp0, Errors.NONE); @@ -206,8 +206,7 @@ public void testEndTxnNotSentIfIncompleteBatches() { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxn(tp0, Errors.NONE); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -218,26 +217,26 @@ public void testEndTxnNotSentIfIncompleteBatches() { @Test public void testFailIfNotReadyForSendNoProducerId() { - assertThrows(IllegalStateException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test public void testFailIfNotReadyForSendIdempotentProducer() { initializeTransactionManager(Optional.empty()); - transactionManager.failIfNotReadyForSend(); + transactionManager.maybeAddPartition(tp0); } @Test public void testFailIfNotReadyForSendIdempotentProducerFatalError() { initializeTransactionManager(Optional.empty()); transactionManager.transitionToFatalError(new KafkaException()); - assertThrows(KafkaException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test public void testFailIfNotReadyForSendNoOngoingTransaction() { doInitTransactions(); - assertThrows(IllegalStateException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -245,14 +244,14 @@ public void testFailIfNotReadyForSendAfterAbortableError() { doInitTransactions(); transactionManager.beginTransaction(); transactionManager.transitionToAbortableError(new KafkaException()); - assertThrows(KafkaException.class, transactionManager::failIfNotReadyForSend); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test public void testFailIfNotReadyForSendAfterFatalError() { doInitTransactions(); transactionManager.transitionToFatalError(new KafkaException()); - assertThrows(KafkaException.class, transactionManager::failIfNotReadyForSend); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -266,8 +265,7 @@ public void testHasOngoingTransactionSuccessfulAbort() { transactionManager.beginTransaction(); assertTrue(transactionManager.hasOngoingTransaction()); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); runUntil(transactionManager::hasOngoingTransaction); prepareAddPartitionsToTxn(partition, Errors.NONE); @@ -291,8 +289,7 @@ public void testHasOngoingTransactionSuccessfulCommit() { transactionManager.beginTransaction(); assertTrue(transactionManager.hasOngoingTransaction()); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasOngoingTransaction()); prepareAddPartitionsToTxn(partition, Errors.NONE); @@ -316,8 +313,7 @@ public void testHasOngoingTransactionAbortableError() { transactionManager.beginTransaction(); assertTrue(transactionManager.hasOngoingTransaction()); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasOngoingTransaction()); prepareAddPartitionsToTxn(partition, Errors.NONE); @@ -344,8 +340,7 @@ public void testHasOngoingTransactionFatalError() { transactionManager.beginTransaction(); assertTrue(transactionManager.hasOngoingTransaction()); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasOngoingTransaction()); prepareAddPartitionsToTxn(partition, Errors.NONE); @@ -361,8 +356,7 @@ public void testMaybeAddPartitionToTransaction() { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasPartitionsToAdd()); assertFalse(transactionManager.isPartitionAdded(partition)); assertTrue(transactionManager.isPartitionPendingAdd(partition)); @@ -375,8 +369,7 @@ public void testMaybeAddPartitionToTransaction() { assertFalse(transactionManager.isPartitionPendingAdd(partition)); // adding the partition again should not have any effect - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertFalse(transactionManager.hasPartitionsToAdd()); assertTrue(transactionManager.isPartitionAdded(partition)); assertFalse(transactionManager.isPartitionPendingAdd(partition)); @@ -388,8 +381,7 @@ public void testAddPartitionToTransactionOverridesRetryBackoffForConcurrentTrans doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasPartitionsToAdd()); assertFalse(transactionManager.isPartitionAdded(partition)); assertTrue(transactionManager.isPartitionPendingAdd(partition)); @@ -408,8 +400,7 @@ public void testAddPartitionToTransactionRetainsRetryBackoffForRegularRetriableE doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasPartitionsToAdd()); assertFalse(transactionManager.isPartitionAdded(partition)); assertTrue(transactionManager.isPartitionPendingAdd(partition)); @@ -428,8 +419,7 @@ public void testAddPartitionToTransactionRetainsRetryBackoffWhenPartitionsAlread doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(partition); + transactionManager.maybeAddPartition(partition); assertTrue(transactionManager.hasPartitionsToAdd()); assertFalse(transactionManager.isPartitionAdded(partition)); assertTrue(transactionManager.isPartitionPendingAdd(partition)); @@ -438,8 +428,7 @@ public void testAddPartitionToTransactionRetainsRetryBackoffWhenPartitionsAlread runUntil(() -> transactionManager.isPartitionAdded(partition)); TopicPartition otherPartition = new TopicPartition("foo", 1); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(otherPartition); + transactionManager.maybeAddPartition(otherPartition); prepareAddPartitionsToTxn(otherPartition, Errors.CONCURRENT_TRANSACTIONS); TransactionManager.TxnRequestHandler handler = transactionManager.nextRequest(false); @@ -449,13 +438,13 @@ public void testAddPartitionToTransactionRetainsRetryBackoffWhenPartitionsAlread @Test public void testNotReadyForSendBeforeInitTransactions() { - assertThrows(IllegalStateException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test public void testNotReadyForSendBeforeBeginTransaction() { doInitTransactions(); - assertThrows(IllegalStateException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -463,14 +452,14 @@ public void testNotReadyForSendAfterAbortableError() { doInitTransactions(); transactionManager.beginTransaction(); transactionManager.transitionToAbortableError(new KafkaException()); - assertThrows(KafkaException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test public void testNotReadyForSendAfterFatalError() { doInitTransactions(); transactionManager.transitionToFatalError(new KafkaException()); - assertThrows(KafkaException.class, () -> transactionManager.failIfNotReadyForSend()); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -478,8 +467,7 @@ public void testIsSendToPartitionAllowedWithPendingPartitionAfterAbortableError( doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); transactionManager.transitionToAbortableError(new KafkaException()); assertFalse(transactionManager.isSendToPartitionAllowed(tp0)); @@ -491,8 +479,7 @@ public void testIsSendToPartitionAllowedWithInFlightPartitionAddAfterAbortableEr doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); // Send the AddPartitionsToTxn request and leave it in-flight runUntil(transactionManager::hasInFlightRequest); @@ -507,8 +494,7 @@ public void testIsSendToPartitionAllowedWithPendingPartitionAfterFatalError() { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); transactionManager.transitionToFatalError(new KafkaException()); assertFalse(transactionManager.isSendToPartitionAllowed(tp0)); @@ -520,8 +506,7 @@ public void testIsSendToPartitionAllowedWithInFlightPartitionAddAfterFatalError( doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); // Send the AddPartitionsToTxn request and leave it in-flight runUntil(transactionManager::hasInFlightRequest); @@ -537,8 +522,7 @@ public void testIsSendToPartitionAllowedWithAddedPartitionAfterAbortableError() transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); runUntil(() -> !transactionManager.hasPartitionsToAdd()); @@ -553,8 +537,7 @@ public void testIsSendToPartitionAllowedWithAddedPartitionAfterFatalError() { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); runUntil(() -> !transactionManager.hasPartitionsToAdd()); @@ -747,8 +730,7 @@ public void testBasicTransaction() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1201,10 +1183,8 @@ public void testTopicAuthorizationFailureInAddPartitions() throws InterruptedExc doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp0); + transactionManager.maybeAddPartition(tp1); FutureRecordMetadata firstPartitionAppend = appendToAccumulator(tp0); FutureRecordMetadata secondPartitionAppend = appendToAccumulator(tp1); @@ -1241,8 +1221,8 @@ public void testCommitWithTopicAuthorizationFailureInAddPartitionsInFlight() thr // Begin a transaction, send two records, and begin commit transactionManager.beginTransaction(); - transactionManager.maybeAddPartitionToTransaction(tp0); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp0); + transactionManager.maybeAddPartition(tp1); FutureRecordMetadata firstPartitionAppend = appendToAccumulator(tp0); FutureRecordMetadata secondPartitionAppend = appendToAccumulator(tp1); TransactionalRequestResult commitResult = transactionManager.beginCommit(); @@ -1286,8 +1266,7 @@ public void testRecoveryFromAbortableErrorTransactionNotStarted() throws Excepti doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(unauthorizedPartition); + transactionManager.maybeAddPartition(unauthorizedPartition); Future responseFuture = appendToAccumulator(unauthorizedPartition); @@ -1309,8 +1288,7 @@ public void testRecoveryFromAbortableErrorTransactionNotStarted() throws Excepti // ensure we can now start a new transaction transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); responseFuture = appendToAccumulator(tp0); @@ -1327,6 +1305,104 @@ public void testRecoveryFromAbortableErrorTransactionNotStarted() throws Excepti runUntil(transactionManager::isReady); } + @Test + public void testRetryAbortTransactionAfterTimeout() throws Exception { + doInitTransactions(); + + transactionManager.beginTransaction(); + transactionManager.maybeAddPartition(tp0); + + prepareAddPartitionsToTxn(tp0, Errors.NONE); + appendToAccumulator(tp0); + runUntil(() -> transactionManager.isPartitionAdded(tp0)); + + TransactionalRequestResult result = transactionManager.beginAbort(); + assertThrows(TimeoutException.class, () -> result.await(0, TimeUnit.MILLISECONDS)); + + prepareEndTxnResponse(Errors.NONE, TransactionResult.ABORT, producerId, epoch); + runUntil(transactionManager::isReady); + assertTrue(result.isSuccessful()); + assertFalse(result.isAcked()); + assertFalse(transactionManager.hasOngoingTransaction()); + + assertThrows(KafkaException.class, transactionManager::initializeTransactions); + assertThrows(KafkaException.class, transactionManager::beginTransaction); + assertThrows(KafkaException.class, transactionManager::beginCommit); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + + assertSame(result, transactionManager.beginAbort()); + result.await(); + + transactionManager.beginTransaction(); + assertTrue(transactionManager.hasOngoingTransaction()); + } + + @Test + public void testRetryCommitTransactionAfterTimeout() throws Exception { + doInitTransactions(); + + transactionManager.beginTransaction(); + transactionManager.maybeAddPartition(tp0); + + prepareAddPartitionsToTxn(tp0, Errors.NONE); + prepareProduceResponse(Errors.NONE, producerId, epoch); + + appendToAccumulator(tp0); + runUntil(() -> transactionManager.isPartitionAdded(tp0)); + + TransactionalRequestResult result = transactionManager.beginCommit(); + assertThrows(TimeoutException.class, () -> result.await(0, TimeUnit.MILLISECONDS)); + + prepareEndTxnResponse(Errors.NONE, TransactionResult.COMMIT, producerId, epoch); + runUntil(transactionManager::isReady); + assertTrue(result.isSuccessful()); + assertFalse(result.isAcked()); + assertFalse(transactionManager.hasOngoingTransaction()); + + assertThrows(KafkaException.class, transactionManager::initializeTransactions); + assertThrows(KafkaException.class, transactionManager::beginTransaction); + assertThrows(KafkaException.class, transactionManager::beginAbort); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + + assertSame(result, transactionManager.beginCommit()); + result.await(); + + transactionManager.beginTransaction(); + assertTrue(transactionManager.hasOngoingTransaction()); + } + + @Test + public void testRetryInitTransactionsAfterTimeout() { + TransactionalRequestResult result = transactionManager.initializeTransactions(); + prepareFindCoordinatorResponse(Errors.NONE, false, CoordinatorType.TRANSACTION, transactionalId); + runUntil(() -> transactionManager.coordinator(CoordinatorType.TRANSACTION) != null); + assertEquals(brokerNode, transactionManager.coordinator(CoordinatorType.TRANSACTION)); + + assertThrows(TimeoutException.class, () -> result.await(0, TimeUnit.MILLISECONDS)); + + prepareInitPidResponse(Errors.NONE, false, producerId, epoch); + runUntil(transactionManager::hasProducerId); + assertTrue(result.isSuccessful()); + assertFalse(result.isAcked()); + + // At this point, the InitProducerId call has returned, but the user has yet + // to complete the call to `initTransactions`. Other transitions should be + // rejected until they do. + + assertThrows(KafkaException.class, transactionManager::beginTransaction); + assertThrows(KafkaException.class, transactionManager::beginAbort); + assertThrows(KafkaException.class, transactionManager::beginCommit); + assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + + assertSame(result, transactionManager.initializeTransactions()); + result.await(); + assertTrue(result.isAcked()); + assertThrows(KafkaException.class, transactionManager::initializeTransactions); + + transactionManager.beginTransaction(); + assertTrue(transactionManager.hasOngoingTransaction()); + } + @Test public void testRecoveryFromAbortableErrorTransactionStarted() throws Exception { final TopicPartition unauthorizedPartition = new TopicPartition("foo", 0); @@ -1334,15 +1410,13 @@ public void testRecoveryFromAbortableErrorTransactionStarted() throws Exception doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxn(tp0, Errors.NONE); Future authorizedTopicProduceFuture = appendToAccumulator(unauthorizedPartition); runUntil(() -> transactionManager.isPartitionAdded(tp0)); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(unauthorizedPartition); + transactionManager.maybeAddPartition(unauthorizedPartition); Future unauthorizedTopicProduceFuture = appendToAccumulator(unauthorizedPartition); prepareAddPartitionsToTxn(singletonMap(unauthorizedPartition, Errors.TOPIC_AUTHORIZATION_FAILED)); runUntil(transactionManager::hasAbortableError); @@ -1365,8 +1439,7 @@ public void testRecoveryFromAbortableErrorTransactionStarted() throws Exception // ensure we can now start a new transaction transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); FutureRecordMetadata nextTransactionFuture = appendToAccumulator(tp0); @@ -1390,8 +1463,7 @@ public void testRecoveryFromAbortableErrorProduceRequestInRetry() throws Excepti doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxn(tp0, Errors.NONE); Future authorizedTopicProduceFuture = appendToAccumulator(tp0); @@ -1403,8 +1475,7 @@ public void testRecoveryFromAbortableErrorProduceRequestInRetry() throws Excepti assertFalse(authorizedTopicProduceFuture.isDone()); assertTrue(accumulator.hasIncomplete()); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(unauthorizedPartition); + transactionManager.maybeAddPartition(unauthorizedPartition); Future unauthorizedTopicProduceFuture = appendToAccumulator(unauthorizedPartition); prepareAddPartitionsToTxn(singletonMap(unauthorizedPartition, Errors.TOPIC_AUTHORIZATION_FAILED)); runUntil(transactionManager::hasAbortableError); @@ -1432,8 +1503,7 @@ public void testRecoveryFromAbortableErrorProduceRequestInRetry() throws Excepti // ensure we can now start a new transaction transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); FutureRecordMetadata nextTransactionFuture = appendToAccumulator(tp0); @@ -1457,8 +1527,7 @@ public void testTransactionalIdAuthorizationFailureInAddPartitions() { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp); + transactionManager.maybeAddPartition(tp); prepareAddPartitionsToTxn(tp, Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED); runUntil(transactionManager::hasError); @@ -1472,8 +1541,7 @@ public void testFlushPendingPartitionsOnCommit() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1511,8 +1579,7 @@ public void testMultipleAddPartitionsPerForOneProduce() throws InterruptedExcept transactionManager.beginTransaction(); // User does one producer.send - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1525,8 +1592,7 @@ public void testMultipleAddPartitionsPerForOneProduce() throws InterruptedExcept runUntil(() -> transactionManager.transactionContainsPartition(tp0)); // In the mean time, the user does a second produce to a different partition - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp1); Future secondResponseFuture = appendToAccumulator(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp1, epoch, producerId); @@ -1591,8 +1657,7 @@ private void verifyProducerFencedForAddPartitionsToTxn(Errors error) throws Inte doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1616,7 +1681,6 @@ private void verifyProducerFencedForAddOffsetsToTxn(Errors error) throws Interru doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); transactionManager.sendOffsetsToTransaction(Collections.emptyMap(), new ConsumerGroupMetadata(consumerGroupId)); Future responseFuture = appendToAccumulator(tp0); @@ -1652,8 +1716,7 @@ public void testInvalidProducerEpochConvertToProducerFencedInEndTxn() throws Int doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); TransactionalRequestResult commitResult = transactionManager.beginCommit(); Future responseFuture = appendToAccumulator(tp0); @@ -1683,8 +1746,7 @@ public void testInvalidProducerEpochFromProduce() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1718,8 +1780,7 @@ public void testDisallowCommitOnProduceFailure() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1758,8 +1819,7 @@ public void testAllowAbortOnProduceFailure() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1782,8 +1842,7 @@ public void testAbortableErrorWhileAbortInProgress() throws InterruptedException doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1813,8 +1872,7 @@ public void testCommitTransactionWithUnsentProduceRequest() throws Exception { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1857,8 +1915,7 @@ public void testCommitTransactionWithInFlightProduceRequest() throws Exception { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1900,8 +1957,7 @@ public void testFindCoordinatorAllowedInAbortableErrorState() throws Interrupted doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1923,8 +1979,7 @@ public void testCancelUnsentAddPartitionsAndProduceOnAbort() throws InterruptedE doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -1950,8 +2005,7 @@ public void testAbortResendsAddPartitionErrorIfRetried() throws InterruptedExcep doInitTransactions(producerId, epoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.UNKNOWN_TOPIC_OR_PARTITION, tp0, epoch, producerId); Future responseFuture = appendToAccumulator(tp0); @@ -1982,8 +2036,7 @@ public void testAbortResendsProduceRequestIfRetried() throws Exception { doInitTransactions(producerId, epoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); prepareProduceResponse(Errors.REQUEST_TIMED_OUT, producerId, epoch); @@ -2011,8 +2064,7 @@ public void testHandlingOfUnknownTopicPartitionErrorOnAddPartitions() throws Int doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -2125,8 +2177,7 @@ public void shouldNotAddPartitionsToTransactionWhenTopicAuthorizationFailed() th doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); assertFalse(responseFuture.isDone()); @@ -2140,8 +2191,7 @@ public void shouldNotSendAbortTxnRequestWhenOnlyAddPartitionsRequestFailed() { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.TOPIC_AUTHORIZATION_FAILED, tp0, epoch, producerId); runUntil(() -> !client.hasPendingResponses()); @@ -2250,7 +2300,6 @@ private TransactionalRequestResult prepareGroupMetadataCommit(Runnable prepareTx assertFalse(addOffsetsResult.isCompleted()); // The request should complete only after the TxnOffsetCommit completes prepareFindCoordinatorResponse(Errors.NONE, false, CoordinatorType.GROUP, consumerGroupId); -// prepareTxnOffsetCommitResponse(consumerGroupId, producerId, epoch, groupInstanceId, memberId, generationId, txnOffsetCommitResponse); prepareTxnCommitResponse.run(); assertNull(transactionManager.coordinator(CoordinatorType.GROUP)); @@ -2265,11 +2314,9 @@ private TransactionalRequestResult prepareGroupMetadataCommit(Runnable prepareTx public void testNoDrainWhenPartitionsPending() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); appendToAccumulator(tp0); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp1); appendToAccumulator(tp1); assertFalse(transactionManager.isSendToPartitionAllowed(tp0)); @@ -2300,13 +2347,11 @@ public void testNoDrainWhenPartitionsPending() throws InterruptedException { public void testAllowDrainInAbortableErrorState() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp1); prepareAddPartitionsToTxn(tp1, Errors.NONE); runUntil(() -> transactionManager.transactionContainsPartition(tp1)); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxn(tp0, Errors.TOPIC_AUTHORIZATION_FAILED); runUntil(transactionManager::hasAbortableError); assertTrue(transactionManager.isSendToPartitionAllowed(tp1)); @@ -2354,8 +2399,7 @@ public void resendFailedProduceRequestAfterAbortableError() throws Exception { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -2376,8 +2420,7 @@ public void testTransitionToAbortableErrorOnBatchExpiry() throws InterruptedExce doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -2417,10 +2460,8 @@ public void testTransitionToAbortableErrorOnMultipleBatchExpiry() throws Interru doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp0); + transactionManager.maybeAddPartition(tp1); Future firstBatchResponse = appendToAccumulator(tp0); Future secondBatchResponse = appendToAccumulator(tp1); @@ -2477,8 +2518,7 @@ public void testDropCommitOnBatchExpiry() throws InterruptedException { doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -2545,8 +2585,7 @@ public void testTransitionToFatalErrorWhenRetriedBatchIsExpired() throws Interru doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -2745,8 +2784,7 @@ public void testAbortTransactionAndReuseSequenceNumberOnError() throws Interrupt doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture0 = appendToAccumulator(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); @@ -2768,11 +2806,11 @@ public void testAbortTransactionAndReuseSequenceNumberOnError() throws Interrupt prepareEndTxnResponse(Errors.NONE, TransactionResult.ABORT, producerId, epoch); runUntil(abortResult::isCompleted); assertTrue(abortResult.isSuccessful()); + abortResult.await(); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); // Send AddPartitionsRequest @@ -2798,16 +2836,15 @@ public void testAbortTransactionAndResetSequenceNumberOnUnknownProducerId() thro doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp1); + transactionManager.maybeAddPartition(tp1); Future successPartitionResponseFuture = appendToAccumulator(tp1); prepareAddPartitionsToTxnResponse(Errors.NONE, tp1, epoch, producerId); prepareProduceResponse(Errors.NONE, producerId, epoch, tp1); runUntil(successPartitionResponseFuture::isDone); assertTrue(transactionManager.isPartitionAdded(tp1)); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture0 = appendToAccumulator(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); prepareProduceResponse(Errors.NONE, producerId, epoch); @@ -2829,11 +2866,11 @@ public void testAbortTransactionAndResetSequenceNumberOnUnknownProducerId() thro prepareEndTxnResponse(Errors.NONE, TransactionResult.ABORT, producerId, epoch); runUntil(abortResult::isCompleted); assertTrue(abortResult.isSuccessful()); + abortResult.await(); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, epoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -2850,8 +2887,7 @@ public void testBumpTransactionalEpochOnAbortableError() throws InterruptedExcep doInitTransactions(producerId, initialEpoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, initialEpoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -2877,11 +2913,11 @@ public void testBumpTransactionalEpochOnAbortableError() throws InterruptedExcep assertTrue(abortResult.isCompleted()); assertTrue(abortResult.isSuccessful()); + abortResult.await(); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, bumpedEpoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -2897,8 +2933,7 @@ public void testBumpTransactionalEpochOnUnknownProducerIdError() throws Interrup doInitTransactions(producerId, initialEpoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, initialEpoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -2925,11 +2960,11 @@ public void testBumpTransactionalEpochOnUnknownProducerIdError() throws Interrup assertTrue(abortResult.isCompleted()); assertTrue(abortResult.isSuccessful()); + abortResult.await(); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, bumpedEpoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -2945,8 +2980,7 @@ public void testBumpTransactionalEpochOnTimeout() throws InterruptedException { doInitTransactions(producerId, initialEpoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, initialEpoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -2985,11 +3019,11 @@ public void testBumpTransactionalEpochOnTimeout() throws InterruptedException { assertTrue(abortResult.isCompleted()); assertTrue(abortResult.isSuccessful()); + abortResult.await(); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.NONE, tp0, bumpedEpoch, producerId); runUntil(() -> transactionManager.isPartitionAdded(tp0)); @@ -3005,8 +3039,7 @@ public void testBumpTransactionalEpochOnRecoverableAddPartitionRequestError() { doInitTransactions(producerId, initialEpoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); prepareAddPartitionsToTxnResponse(Errors.INVALID_PRODUCER_ID_MAPPING, tp0, initialEpoch, producerId); runUntil(transactionManager::hasAbortableError); TransactionalRequestResult abortResult = transactionManager.beginAbort(); @@ -3026,8 +3059,7 @@ public void testBumpTransactionalEpochOnRecoverableAddOffsetsRequestError() thro doInitTransactions(producerId, initialEpoch); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); Future responseFuture = appendToAccumulator(tp0); @@ -3280,8 +3312,7 @@ private void verifyCommitOrAbortTransactionRetriable(TransactionResult firstTran doInitTransactions(); transactionManager.beginTransaction(); - transactionManager.failIfNotReadyForSend(); - transactionManager.maybeAddPartitionToTransaction(tp0); + transactionManager.maybeAddPartition(tp0); appendToAccumulator(tp0); From 4ba2fda2f4b964ad499ba8d07193cba299825d1e Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 22 Nov 2021 15:25:20 -0800 Subject: [PATCH 04/10] Fill in test failure message --- .../org/apache/kafka/clients/producer/KafkaProducerTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java index 2abd376910cb3..4a550ee0117ba 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java @@ -930,7 +930,8 @@ public void testInitTransactionsResponseAfterTimeout() throws Exception { FindCoordinatorResponse.prepareResponse(Errors.NONE, "bad-transaction", host1)); Future future = executor.submit(producer::initTransactions); - TestUtils.waitForCondition(client::hasInFlightRequests, "blah blah"); + TestUtils.waitForCondition(client::hasInFlightRequests, + "Timed out while waiting for expected `InitProducerId` request to be sent"); time.sleep(maxBlockMs); TestUtils.assertFutureThrows(future, TimeoutException.class); From 15d0954cbc2a75748e700f3c58d367ace4f22762 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 29 Nov 2021 18:01:54 -0800 Subject: [PATCH 05/10] Fix error messages --- .../kafka/clients/producer/internals/TransactionManager.java | 5 +++-- 1 file changed, 3 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 5f05fc92e4c44..08625080e76a8 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 @@ -409,7 +409,8 @@ public synchronized TransactionalRequestResult sendOffsetsToTransaction(final Ma maybeFailWithError(); if (currentState != State.IN_TRANSACTION) { - throw new KafkaException("Cannot send offsets to transaction either because in state " + currentState); + throw new KafkaException("Cannot send offsets if a transaction is not in progress " + + "(currentState= " + currentState + ")"); } log.debug("Begin adding offsets {} for consumer group {} to transaction", offsets, groupMetadata); @@ -433,7 +434,7 @@ public synchronized void maybeAddPartition(TopicPartition topicPartition) { if (isTransactional()) { if (!hasProducerId()) { throw new KafkaException("Cannot add partition " + topicPartition + - "to transaction before completing a call to initTransactions"); + " to transaction before completing a call to initTransactions"); } else if (currentState != State.IN_TRANSACTION) { throw new KafkaException("Cannot add partition " + topicPartition + " to transaction while in state " + currentState); From 636e9fbdbd79973033bb0d2cd8b2adae4d0dffe5 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Tue, 30 Nov 2021 10:12:37 -0800 Subject: [PATCH 06/10] Use IllegalStateException for illegal state transitions --- .../internals/TransactionManager.java | 22 +++--- .../clients/producer/KafkaProducerTest.java | 2 +- .../internals/TransactionManagerTest.java | 69 +++++++------------ 3 files changed, 35 insertions(+), 58 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 08625080e76a8..add97df99018d 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 @@ -409,7 +409,7 @@ public synchronized TransactionalRequestResult sendOffsetsToTransaction(final Ma maybeFailWithError(); if (currentState != State.IN_TRANSACTION) { - throw new KafkaException("Cannot send offsets if a transaction is not in progress " + + throw new IllegalStateException("Cannot send offsets if a transaction is not in progress " + "(currentState= " + currentState + ")"); } @@ -433,10 +433,10 @@ public synchronized void maybeAddPartition(TopicPartition topicPartition) { if (isTransactional()) { if (!hasProducerId()) { - throw new KafkaException("Cannot add partition " + topicPartition + + throw new IllegalStateException("Cannot add partition " + topicPartition + " to transaction before completing a call to initTransactions"); } else if (currentState != State.IN_TRANSACTION) { - throw new KafkaException("Cannot add partition " + topicPartition + + throw new IllegalStateException("Cannot add partition " + topicPartition + " to transaction while in state " + currentState); } else if (isPartitionAdded(topicPartition) || isPartitionPendingAdd(topicPartition)) { return; @@ -1074,7 +1074,7 @@ private void transitionTo(State target) { private void transitionTo(State target, RuntimeException error) { if (!currentState.isTransitionValid(currentState, target)) { String idString = transactionalId == null ? "" : "TransactionalId " + transactionalId + ": "; - throw new KafkaException(idString + "Invalid transition attempted from state " + throw new IllegalStateException(idString + "Invalid transition attempted from state " + currentState.name() + " to state " + target.name()); } @@ -1184,18 +1184,14 @@ private TxnOffsetCommitHandler txnOffsetCommitHandler(TransactionalRequestResult return new TxnOffsetCommitHandler(result, builder); } - private void throwPendingTransitionError(String operation) { - throw new KafkaException("Cannot attempt operation `" + operation + "` " - + "because the previous call to `" + pendingTransition.operation + "` " - + "timed out and must be retried"); - } - private void throwIfPendingState(String operation) { if (pendingTransition != null) { if (pendingTransition.result.isAcked()) { pendingTransition = null; } else { - throwPendingTransitionError(operation); + throw new IllegalStateException("Cannot attempt operation `" + operation + "` " + + "because the previous call to `" + pendingTransition.operation + "` " + + "timed out and must be retried"); } } } @@ -1211,7 +1207,9 @@ private TransactionalRequestResult handleCachedTransactionRequestResult( if (pendingTransition.result.isAcked()) { pendingTransition = null; } else if (nextState != pendingTransition.state) { - throwPendingTransitionError(operation); + throw new IllegalStateException("Cannot attempt operation `" + operation + "` " + + "because the previous call to `" + pendingTransition.operation + "` " + + "timed out and must be retried"); } else { return pendingTransition.result; } diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java index 4a550ee0117ba..96a30341b842f 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java @@ -1291,7 +1291,7 @@ public void testOnlyCanExecuteCloseAfterInitTransactionsTimeout() { assertThrows(TimeoutException.class, producer::initTransactions); // other transactional operations should not be allowed if we catch the error after initTransactions failed try { - assertThrows(KafkaException.class, producer::beginTransaction); + assertThrows(IllegalStateException.class, producer::beginTransaction); } finally { producer.close(Duration.ofMillis(0)); } 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 aef007dd14f70..4227db5e61e62 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 @@ -217,7 +217,7 @@ public void testEndTxnNotSentIfIncompleteBatches() { @Test public void testFailIfNotReadyForSendNoProducerId() { - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -236,7 +236,7 @@ public void testFailIfNotReadyForSendIdempotentProducerFatalError() { @Test public void testFailIfNotReadyForSendNoOngoingTransaction() { doInitTransactions(); - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -438,13 +438,13 @@ public void testAddPartitionToTransactionRetainsRetryBackoffWhenPartitionsAlread @Test public void testNotReadyForSendBeforeInitTransactions() { - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test public void testNotReadyForSendBeforeBeginTransaction() { doInitTransactions(); - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); } @Test @@ -795,7 +795,7 @@ public void testDisconnectAndRetry() { public void testInitializeTransactionsTwiceRaisesError() { doInitTransactions(producerId, epoch); assertTrue(transactionManager.hasProducerId()); - assertThrows(KafkaException.class, () -> transactionManager.initializeTransactions()); + assertThrows(IllegalStateException.class, () -> transactionManager.initializeTransactions()); } @Test @@ -1325,10 +1325,10 @@ public void testRetryAbortTransactionAfterTimeout() throws Exception { assertFalse(result.isAcked()); assertFalse(transactionManager.hasOngoingTransaction()); - assertThrows(KafkaException.class, transactionManager::initializeTransactions); - assertThrows(KafkaException.class, transactionManager::beginTransaction); - assertThrows(KafkaException.class, transactionManager::beginCommit); - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, transactionManager::initializeTransactions); + assertThrows(IllegalStateException.class, transactionManager::beginTransaction); + assertThrows(IllegalStateException.class, transactionManager::beginCommit); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); assertSame(result, transactionManager.beginAbort()); result.await(); @@ -1359,10 +1359,10 @@ public void testRetryCommitTransactionAfterTimeout() throws Exception { assertFalse(result.isAcked()); assertFalse(transactionManager.hasOngoingTransaction()); - assertThrows(KafkaException.class, transactionManager::initializeTransactions); - assertThrows(KafkaException.class, transactionManager::beginTransaction); - assertThrows(KafkaException.class, transactionManager::beginAbort); - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, transactionManager::initializeTransactions); + assertThrows(IllegalStateException.class, transactionManager::beginTransaction); + assertThrows(IllegalStateException.class, transactionManager::beginAbort); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); assertSame(result, transactionManager.beginCommit()); result.await(); @@ -1389,15 +1389,15 @@ public void testRetryInitTransactionsAfterTimeout() { // to complete the call to `initTransactions`. Other transitions should be // rejected until they do. - assertThrows(KafkaException.class, transactionManager::beginTransaction); - assertThrows(KafkaException.class, transactionManager::beginAbort); - assertThrows(KafkaException.class, transactionManager::beginCommit); - assertThrows(KafkaException.class, () -> transactionManager.maybeAddPartition(tp0)); + assertThrows(IllegalStateException.class, transactionManager::beginTransaction); + assertThrows(IllegalStateException.class, transactionManager::beginAbort); + assertThrows(IllegalStateException.class, transactionManager::beginCommit); + assertThrows(IllegalStateException.class, () -> transactionManager.maybeAddPartition(tp0)); assertSame(result, transactionManager.initializeTransactions()); result.await(); assertTrue(result.isAcked()); - assertThrows(KafkaException.class, transactionManager::initializeTransactions); + assertThrows(IllegalStateException.class, transactionManager::initializeTransactions); transactionManager.beginTransaction(); assertTrue(transactionManager.hasOngoingTransaction()); @@ -1791,19 +1791,8 @@ public void testDisallowCommitOnProduceFailure() throws InterruptedException { runUntil(commitResult::isCompleted); // commit should be cancelled with exception without being sent. - try { - commitResult.await(); - fail(); // the get() must throw an exception. - } catch (KafkaException e) { - // Expected - } - - try { - responseFuture.get(); - fail("Expected produce future to raise an exception"); - } catch (ExecutionException e) { - assertTrue(e.getCause() instanceof OutOfOrderSequenceException); - } + assertThrows(KafkaException.class, commitResult::await); + TestUtils.assertFutureThrows(responseFuture, OutOfOrderSequenceException.class); // Commit is not allowed, so let's abort and try again. TransactionalRequestResult abortResult = transactionManager.beginAbort(); @@ -1992,12 +1981,7 @@ public void testCancelUnsentAddPartitionsAndProduceOnAbort() throws InterruptedE assertTrue(abortResult.isSuccessful()); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. - try { - responseFuture.get(); - fail("Expected produce future to raise an exception"); - } catch (ExecutionException e) { - assertTrue(e.getCause() instanceof KafkaException); - } + TestUtils.assertFutureThrows(responseFuture, KafkaException.class); } @Test @@ -2023,12 +2007,7 @@ public void testAbortResendsAddPartitionErrorIfRetried() throws InterruptedExcep assertTrue(abortResult.isSuccessful()); assertTrue(transactionManager.isReady()); // make sure we are ready for a transaction now. - try { - responseFuture.get(); - fail("Expected produce future to raise an exception"); - } catch (ExecutionException e) { - assertTrue(e.getCause() instanceof KafkaException); - } + TestUtils.assertFutureThrows(responseFuture, KafkaException.class); } @Test @@ -3190,12 +3169,12 @@ public void testRetryCommitTransaction() throws InterruptedException { @Test public void testRetryAbortTransactionAfterCommitTimeout() { - assertThrows(KafkaException.class, () -> verifyCommitOrAbortTransactionRetriable(TransactionResult.COMMIT, TransactionResult.ABORT)); + assertThrows(IllegalStateException.class, () -> verifyCommitOrAbortTransactionRetriable(TransactionResult.COMMIT, TransactionResult.ABORT)); } @Test public void testRetryCommitTransactionAfterAbortTimeout() { - assertThrows(KafkaException.class, () -> verifyCommitOrAbortTransactionRetriable(TransactionResult.ABORT, TransactionResult.COMMIT)); + assertThrows(IllegalStateException.class, () -> verifyCommitOrAbortTransactionRetriable(TransactionResult.ABORT, TransactionResult.COMMIT)); } @Test From b50738904fb0f8e35244cbfdb62929e6215e1d1b Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Wed, 5 Jan 2022 13:12:25 -0800 Subject: [PATCH 07/10] Fix failing authorization tests --- .../kafka/api/AuthorizerIntegrationTest.scala | 65 ++++++++++--------- 1 file changed, 36 insertions(+), 29 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala index 6efb860675d96..1da4fcd587edd 100644 --- a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala @@ -58,7 +58,7 @@ import org.apache.kafka.common.resource.{PatternType, Resource, ResourcePattern, import org.apache.kafka.common.security.auth.{AuthenticationContext, KafkaPrincipal, SecurityProtocol} import org.apache.kafka.common.security.authenticator.DefaultKafkaPrincipalBuilder import org.apache.kafka.common.utils.Utils -import org.apache.kafka.common.{ElectionType, IsolationLevel, Node, TopicPartition, Uuid, requests} +import org.apache.kafka.common.{ElectionType, IsolationLevel, KafkaException, Node, TopicPartition, Uuid, requests} import org.apache.kafka.test.{TestUtils => JTestUtils} import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test, TestInfo} @@ -1883,31 +1883,38 @@ class AuthorizerIntegrationTest extends BaseRequestTest { def testIdempotentProducerNoIdempotentWriteAclInInitProducerId(): Unit = { createTopic(topic) addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WildcardHost, READ, ALLOW)), topicResource) - shouldIdempotentProducerFailInInitProducerId(true) + assertIdempotentSendAuthorizationFailure() } - def shouldIdempotentProducerFailInInitProducerId(expectAuthException: Boolean): Unit = { + private def assertIdempotentSendSuccess(): Unit = { val producer = buildIdempotentProducer() - try { + producer.send(new ProducerRecord[Array[Byte], Array[Byte]](topic, "hi".getBytes)).get() + } + + private def assertIdempotentSendAuthorizationFailure(): Unit = { + val producer = buildIdempotentProducer() + + def assertClusterAuthFailure(): Unit = { // the InitProducerId is sent asynchronously, so we expect the error either in the callback // or raised from send itself - producer.send(new ProducerRecord[Array[Byte], Array[Byte]](topic, "hi".getBytes)).get() - if (expectAuthException) - fail("Should have raised ClusterAuthorizationException") - } catch { - case e: ExecutionException => - assertTrue(e.getCause.isInstanceOf[ClusterAuthorizationException]) - } - try { - // the second time, the call to send itself should fail (the producer becomes unusable - // if no producerId can be obtained) - producer.send(new ProducerRecord[Array[Byte], Array[Byte]](topic, "hi".getBytes)).get() - if (expectAuthException) - fail("Should have raised ClusterAuthorizationException") - } catch { - case e: ExecutionException => - assertTrue(e.getCause.isInstanceOf[ClusterAuthorizationException]) + val exception = assertThrows(classOf[Exception], () => { + val future = producer.send(new ProducerRecord[Array[Byte], Array[Byte]](topic, "hi".getBytes)) + future.get() + }) + + exception match { + case e@ (_: KafkaException | _: ExecutionException) => + assertTrue(exception.getCause.isInstanceOf[ClusterAuthorizationException]) + case _ => + fail(s"Unexpected exception type raised from send: ${exception.getClass}") + } } + + assertClusterAuthFailure() + + // the second time, the call to send itself should fail (the producer becomes unusable + // if no producerId can be obtained) + assertClusterAuthFailure() } @Test @@ -2119,14 +2126,14 @@ class AuthorizerIntegrationTest extends BaseRequestTest { for (_ <- 1 to 3) { addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WildcardHost, DESCRIBE, ALLOW)), topicResource) - shouldIdempotentProducerFailInInitProducerId(true) + assertIdempotentSendAuthorizationFailure() addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WildcardHost, WRITE, ALLOW)), topicResource) - shouldIdempotentProducerFailInInitProducerId(false) + assertIdempotentSendSuccess() removeAllClientAcls() addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WildcardHost, DESCRIBE, ALLOW)), topicResource) - shouldIdempotentProducerFailInInitProducerId(true) + assertIdempotentSendAuthorizationFailure() } } @@ -2149,7 +2156,7 @@ class AuthorizerIntegrationTest extends BaseRequestTest { addAndVerifyAcls(Set(acl1, acl4, acl5), topicResource) addAndVerifyAcls(Set(acl2, acl3), unrelatedTopicResource) addAndVerifyAcls(Set(acl2, acl3), unrelatedGroupResource) - shouldIdempotentProducerFailInInitProducerId(false) + assertIdempotentSendSuccess() } @Test @@ -2157,11 +2164,11 @@ class AuthorizerIntegrationTest extends BaseRequestTest { createTopic(topic) val allowWriteAce = new AccessControlEntry(clientPrincipalString, WildcardHost, WRITE, ALLOW) addAndVerifyAcls(Set(allowWriteAce), topicResource) - shouldIdempotentProducerFailInInitProducerId(false) + assertIdempotentSendSuccess() val denyWriteAce = new AccessControlEntry(clientPrincipalString, WildcardHost, WRITE, DENY) addAndVerifyAcls(Set(denyWriteAce), topicResource) - shouldIdempotentProducerFailInInitProducerId(true) + assertIdempotentSendAuthorizationFailure() } @Test @@ -2175,10 +2182,10 @@ class AuthorizerIntegrationTest extends BaseRequestTest { addAndVerifyAcls(Set(allowWriteAce), prefixed) addAndVerifyAcls(Set(allowWriteAce), literal) - shouldIdempotentProducerFailInInitProducerId(false) + assertIdempotentSendSuccess() addAndVerifyAcls(Set(denyWriteAce), wildcard) - shouldIdempotentProducerFailInInitProducerId(true) + assertIdempotentSendAuthorizationFailure() } @Test @@ -2191,7 +2198,7 @@ class AuthorizerIntegrationTest extends BaseRequestTest { addAndVerifyAcls(Set(denyWriteAce), prefixed) addAndVerifyAcls(Set(allowWriteAce), literal) - shouldIdempotentProducerFailInInitProducerId(true) + assertIdempotentSendAuthorizationFailure() } @Test From d121f615132faca87370444b1a663f1544396884 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Wed, 5 Jan 2022 13:12:41 -0800 Subject: [PATCH 08/10] Cleanup some error messages --- .../clients/producer/internals/TransactionManager.java | 8 +++++--- .../org/apache/kafka/common/utils/ProducerIdAndEpoch.java | 2 +- 2 files changed, 6 insertions(+), 4 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 add97df99018d..77c5b75e20451 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 @@ -1104,10 +1104,12 @@ private void maybeFailWithError() { // for ProducerFencedException, do not wrap it as a KafkaException // but create a new instance without the call trace since it was not thrown because of the current call if (lastError instanceof ProducerFencedException) { - throw new ProducerFencedException("The producer has been rejected from the broker because " + - "it tried to use an old epoch with the transactionalId"); + throw new ProducerFencedException("Producer with transactionalId '" + transactionalId + + "' and " + producerIdAndEpoch + " has been fenced by another producer " + + "with the same transactionalId"); } else if (lastError instanceof InvalidProducerEpochException) { - throw new InvalidProducerEpochException("Producer attempted to produce with an old epoch " + producerIdAndEpoch); + throw new InvalidProducerEpochException("Producer with transactionalId '" + transactionalId + + "' and " + producerIdAndEpoch + " attempted to produce with an old epoch"); } else { throw new KafkaException("Cannot execute transactional method because we are in an error state", lastError); } diff --git a/clients/src/main/java/org/apache/kafka/common/utils/ProducerIdAndEpoch.java b/clients/src/main/java/org/apache/kafka/common/utils/ProducerIdAndEpoch.java index 674b42352b787..5061da1343cdc 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/ProducerIdAndEpoch.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/ProducerIdAndEpoch.java @@ -35,7 +35,7 @@ public boolean isValid() { @Override public String toString() { - return "(producerId=" + producerId + ", epoch=" + epoch + ")"; + return "ProducerIdAndEpoch(producerId=" + producerId + ", epoch=" + epoch + ")"; } @Override From 3138edebda371d44f44d1b95b966de58b2ec4f78 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Wed, 5 Jan 2022 15:43:06 -0800 Subject: [PATCH 09/10] Fail immediately in initTransactions if the client is in an error state --- .../kafka/clients/producer/internals/TransactionManager.java | 2 ++ 1 file changed, 2 insertions(+) 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 77c5b75e20451..521d5da4f8ed6 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 @@ -329,6 +329,8 @@ public synchronized TransactionalRequestResult initializeTransactions() { } synchronized TransactionalRequestResult initializeTransactions(ProducerIdAndEpoch producerIdAndEpoch) { + maybeFailWithError(); + boolean isEpochBump = producerIdAndEpoch != ProducerIdAndEpoch.NONE; return handleCachedTransactionRequestResult(() -> { // If this is an epoch bump, we will transition the state as part of handling the EndTxnRequest From 23f17289473bf4876bf884a15c193ec5c431402b Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Wed, 5 Jan 2022 18:48:02 -0800 Subject: [PATCH 10/10] Fix one more test assertion due to exception change --- .../test/scala/integration/kafka/api/TransactionsTest.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/TransactionsTest.scala b/core/src/test/scala/integration/kafka/api/TransactionsTest.scala index 1fbba9e7ee394..2d8689fa64d9f 100644 --- a/core/src/test/scala/integration/kafka/api/TransactionsTest.scala +++ b/core/src/test/scala/integration/kafka/api/TransactionsTest.scala @@ -30,7 +30,7 @@ import kafka.utils.TestUtils.consumeRecords import org.apache.kafka.clients.consumer.{ConsumerConfig, ConsumerGroupMetadata, KafkaConsumer, OffsetAndMetadata} import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import org.apache.kafka.common.errors.{InvalidProducerEpochException, ProducerFencedException, TimeoutException} -import org.apache.kafka.common.{KafkaException, TopicPartition} +import org.apache.kafka.common.TopicPartition import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test, TestInfo} @@ -604,7 +604,7 @@ class TransactionsTest extends KafkaServerTestHarness { val producer = createTransactionalProducer(transactionalId = "normalProducer") producer.initTransactions() - assertThrows(classOf[KafkaException], () => producer.initTransactions()) + assertThrows(classOf[IllegalStateException], () => producer.initTransactions()) } @Test