From ddd2d3ff5501b3d1444d319734b48c2e0a221e9f Mon Sep 17 00:00:00 2001 From: Bob Barrett Date: Fri, 12 Jul 2019 14:54:01 -0700 Subject: [PATCH 1/3] KAFKA-8635: Skip client poll in Sender loop when no request is sent This patch breaks up maybeSendTransactionalRequest() and changes it to return false if a FindCoordinatorRequest is enqueued. If this is the case, we no longer poll for because no request was actually sent. --- .../clients/producer/internals/Sender.java | 25 ++++++++--- .../producer/internals/SenderTest.java | 44 +++++++++++++------ 2 files changed, 48 insertions(+), 21 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java index 121ddb25595fa..7c1d06ba16fd0 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java @@ -302,12 +302,22 @@ void runOnce() { transactionManager.transitionToFatalError( new KafkaException("The client hasn't received acknowledgment for " + "some previously sent messages and can no longer retry them. It isn't safe to continue.")); - } else if (transactionManager.hasInFlightTransactionalRequest() || maybeSendTransactionalRequest()) { + } else if (transactionManager.hasInFlightTransactionalRequest()) { // as long as there are outstanding transactional requests, we simply wait for them to return client.poll(retryBackoffMs, time.milliseconds()); return; } + maybeBeginFlush(); + TransactionManager.TxnRequestHandler nextRequestHandler = transactionManager.nextRequestHandler(accumulator.hasIncomplete()); + if (nextRequestHandler != null) { + if (maybeSendTransactionalRequest(nextRequestHandler)) { + client.poll(retryBackoffMs, time.milliseconds()); + } + + return; + } + // do not continue sending if the transaction manager is in a failed state or if there // is no producer id (for the idempotent case). if (transactionManager.hasFatalError() || !transactionManager.hasProducerId()) { @@ -412,7 +422,7 @@ private long sendProducerData(long now) { return pollTimeout; } - private boolean maybeSendTransactionalRequest() { + private void maybeBeginFlush() { if (transactionManager.isCompleting() && accumulator.hasIncomplete()) { if (transactionManager.isAborting()) accumulator.abortUndrainedBatches(new KafkaException("Failing batch since transaction was aborted")); @@ -423,11 +433,12 @@ private boolean maybeSendTransactionalRequest() { if (!accumulator.flushInProgress()) accumulator.beginFlush(); } + } - TransactionManager.TxnRequestHandler nextRequestHandler = transactionManager.nextRequestHandler(accumulator.hasIncomplete()); - if (nextRequestHandler == null) - return false; - + /** + * Enqueue a FindCoordinatorRequest if needed, otherwise send the next request + */ + private boolean maybeSendTransactionalRequest(TransactionManager.TxnRequestHandler nextRequestHandler) { AbstractRequest.Builder requestBuilder = nextRequestHandler.requestBuilder(); while (!forceClose) { Node targetNode = null; @@ -470,7 +481,7 @@ private boolean maybeSendTransactionalRequest() { metadata.requestUpdate(); } transactionManager.retry(nextRequestHandler); - return true; + return false; } private void maybeAbortBatches(RuntimeException exception) { 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 5ae76c34a0142..194176d8e2235 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 @@ -113,6 +113,8 @@ import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; public class SenderTest { private static final int MAX_REQUEST_SIZE = 1024 * 1024; @@ -121,6 +123,7 @@ public class SenderTest { private static final double EPS = 0.0001; private static final int MAX_BLOCK_TIMEOUT = 1000; private static final int REQUEST_TIMEOUT = 1000; + private static final long RETRY_BACKOFF_MS = 50; private TopicPartition tp0 = new TopicPartition("test", 0); private TopicPartition tp1 = new TopicPartition("test", 1); @@ -315,7 +318,7 @@ public void testSenderMetricsTemplates() throws Exception { metrics = new Metrics(new MetricConfig().tags(clientTags)); SenderMetricsRegistry metricsRegistry = new SenderMetricsRegistry(metrics); Sender sender = new Sender(logContext, client, metadata, this.accumulator, false, MAX_REQUEST_SIZE, ACKS_ALL, - 1, metricsRegistry, time, REQUEST_TIMEOUT, 50, null, apiVersions); + 1, metricsRegistry, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, null, apiVersions); // Append a message so that topic metrics are created accumulator.append(tp0, 0L, "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT); @@ -343,7 +346,7 @@ public void testRetries() throws Exception { SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); try { Sender sender = new Sender(logContext, client, metadata, this.accumulator, false, MAX_REQUEST_SIZE, ACKS_ALL, - maxRetries, senderMetrics, time, REQUEST_TIMEOUT, 50, null, apiVersions); + maxRetries, senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, null, apiVersions); // do a successful retry Future future = accumulator.append(tp0, 0L, "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; sender.runOnce(); // connect @@ -401,7 +404,7 @@ public void testSendInOrder() throws Exception { try { Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, maxRetries, - senderMetrics, time, REQUEST_TIMEOUT, 50, null, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, null, apiVersions); // Create a two broker cluster, with partition 0 on broker 0 and partition 1 on broker 1 MetadataResponse metadataUpdate1 = TestUtils.metadataUpdateWith(2, Collections.singletonMap("test", 2)); client.prepareMetadataUpdate(metadataUpdate1); @@ -1172,7 +1175,7 @@ public void testResetOfProducerStateShouldAllowQueuedBatchesToDrain() throws Exc SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, maxRetries, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future failedResponse = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; @@ -1214,7 +1217,7 @@ public void testCloseWithProducerIdReset() throws Exception { SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, 10, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future failedResponse = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; @@ -1253,7 +1256,7 @@ public void testForceCloseWithProducerIdReset() throws Exception { SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, 10, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future failedResponse = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; @@ -1289,7 +1292,7 @@ public void testBatchesDrainedWithOldProducerIdShouldFailWithOutOfOrderSequenceO SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, maxRetries, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future failedResponse = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; @@ -1778,7 +1781,7 @@ public void testSequenceNumberIncrement() throws InterruptedException { SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, maxRetries, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future responseFuture = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; client.prepareResponse(new MockClient.RequestMatcher() { @@ -1820,7 +1823,7 @@ public void testAbortRetryWhenProducerIdChanges() throws InterruptedException { Metrics m = new Metrics(); SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, maxRetries, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future responseFuture = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; sender.runOnce(); // connect. @@ -1859,7 +1862,7 @@ public void testResetWhenOutOfOrderSequenceReceived() throws InterruptedExceptio SenderMetricsRegistry senderMetrics = new SenderMetricsRegistry(m); Sender sender = new Sender(logContext, client, metadata, this.accumulator, true, MAX_REQUEST_SIZE, ACKS_ALL, maxRetries, - senderMetrics, time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); Future responseFuture = accumulator.append(tp0, time.milliseconds(), "key".getBytes(), "value".getBytes(), null, null, MAX_BLOCK_TIMEOUT).future; sender.runOnce(); // connect. @@ -2205,7 +2208,7 @@ public void testTransactionalRequestsSentOnShutdown() { try { TransactionManager txnManager = new TransactionManager(logContext, "testTransactionalRequestsSentOnShutdown", 6000, 100); Sender sender = new Sender(logContext, client, metadata, this.accumulator, false, MAX_REQUEST_SIZE, ACKS_ALL, - maxRetries, senderMetrics, time, REQUEST_TIMEOUT, 50, txnManager, apiVersions); + maxRetries, senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, txnManager, apiVersions); ProducerIdAndEpoch producerIdAndEpoch = new ProducerIdAndEpoch(123456L, (short) 0); TopicPartition tp = new TopicPartition("testTransactionalRequestsSentOnShutdown", 1); @@ -2237,7 +2240,7 @@ public void testIncompleteTransactionAbortOnShutdown() { try { TransactionManager txnManager = new TransactionManager(logContext, "testIncompleteTransactionAbortOnShutdown", 6000, 100); Sender sender = new Sender(logContext, client, metadata, this.accumulator, false, MAX_REQUEST_SIZE, ACKS_ALL, - maxRetries, senderMetrics, time, REQUEST_TIMEOUT, 50, txnManager, apiVersions); + maxRetries, senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, txnManager, apiVersions); ProducerIdAndEpoch producerIdAndEpoch = new ProducerIdAndEpoch(123456L, (short) 0); TopicPartition tp = new TopicPartition("testIncompleteTransactionAbortOnShutdown", 1); @@ -2268,7 +2271,7 @@ public void testForceShutdownWithIncompleteTransaction() { try { TransactionManager txnManager = new TransactionManager(logContext, "testForceShutdownWithIncompleteTransaction", 6000, 100); Sender sender = new Sender(logContext, client, metadata, this.accumulator, false, MAX_REQUEST_SIZE, ACKS_ALL, - maxRetries, senderMetrics, time, REQUEST_TIMEOUT, 50, txnManager, apiVersions); + maxRetries, senderMetrics, time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, txnManager, apiVersions); ProducerIdAndEpoch producerIdAndEpoch = new ProducerIdAndEpoch(123456L, (short) 0); TopicPartition tp = new TopicPartition("testForceShutdownWithIncompleteTransaction", 1); @@ -2293,6 +2296,19 @@ public void testForceShutdownWithIncompleteTransaction() { } } + @Test + public void testDoNotPollWhenNoRequestSent() { + client = spy(new MockClient(time, metadata)); + + TransactionManager txnManager = new TransactionManager(logContext, "testDoNotPollWhenNoRequestSent", 6000, 100); + ProducerIdAndEpoch producerIdAndEpoch = new ProducerIdAndEpoch(123456L, (short) 0); + setupWithTransactionState(txnManager); + doInitTransactions(txnManager, producerIdAndEpoch); + + // doInitTransactions calls sender.doOnce three times, only two requests are sent, so we should only poll twice + verify(client, times(2)).poll(eq(RETRY_BACKOFF_MS), anyLong()); + } + class AssertEndTxnRequestMatcher implements MockClient.RequestMatcher { private TransactionResult requiredResult; @@ -2418,7 +2434,7 @@ private void setupWithTransactionState(TransactionManager transactionManager, bo deliveryTimeoutMs, metrics, metricGrpName, time, apiVersions, transactionManager, pool); this.senderMetricsRegistry = new SenderMetricsRegistry(this.metrics); this.sender = new Sender(logContext, this.client, this.metadata, this.accumulator, guaranteeOrder, MAX_REQUEST_SIZE, ACKS_ALL, - Integer.MAX_VALUE, this.senderMetricsRegistry, this.time, REQUEST_TIMEOUT, 50, transactionManager, apiVersions); + Integer.MAX_VALUE, this.senderMetricsRegistry, this.time, REQUEST_TIMEOUT, RETRY_BACKOFF_MS, transactionManager, apiVersions); metadata.add("test"); this.client.updateMetadata(TestUtils.metadataUpdateWith(1, Collections.singletonMap("test", 2))); From 2da11bc9565979b926885882eeaf95a92a21b98f Mon Sep 17 00:00:00 2001 From: Bob Barrett Date: Tue, 16 Jul 2019 16:57:06 -0700 Subject: [PATCH 2/3] PR feedback: consolidate awaitReady path, move in-flight check into helper method --- .../clients/producer/internals/Sender.java | 112 +++++++++--------- .../internals/TransactionManager.java | 6 +- 2 files changed, 57 insertions(+), 61 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java index 7c1d06ba16fd0..6cc41c54c309d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java @@ -46,6 +46,7 @@ import org.apache.kafka.common.record.MemoryRecords; import org.apache.kafka.common.record.RecordBatch; import org.apache.kafka.common.requests.AbstractRequest; +import org.apache.kafka.common.requests.FindCoordinatorRequest; import org.apache.kafka.common.requests.InitProducerIdRequest; import org.apache.kafka.common.requests.InitProducerIdResponse; import org.apache.kafka.common.requests.ProduceRequest; @@ -302,19 +303,7 @@ void runOnce() { transactionManager.transitionToFatalError( new KafkaException("The client hasn't received acknowledgment for " + "some previously sent messages and can no longer retry them. It isn't safe to continue.")); - } else if (transactionManager.hasInFlightTransactionalRequest()) { - // as long as there are outstanding transactional requests, we simply wait for them to return - client.poll(retryBackoffMs, time.milliseconds()); - return; - } - - maybeBeginFlush(); - TransactionManager.TxnRequestHandler nextRequestHandler = transactionManager.nextRequestHandler(accumulator.hasIncomplete()); - if (nextRequestHandler != null) { - if (maybeSendTransactionalRequest(nextRequestHandler)) { - client.poll(retryBackoffMs, time.milliseconds()); - } - + } else if (maybeSendAndPollTransactionalRequest()) { return; } @@ -422,7 +411,16 @@ private long sendProducerData(long now) { return pollTimeout; } - private void maybeBeginFlush() { + /** + * Returns true if a transactional request is sent or polled, or if a FindCoordinator request is enqueued + */ + private boolean maybeSendAndPollTransactionalRequest() { + if (transactionManager.hasInFlightTransactionalRequest()) { + // as long as there are outstanding transactional requests, we simply wait for them to return + client.poll(retryBackoffMs, time.milliseconds()); + return true; + } + if (transactionManager.isCompleting() && accumulator.hasIncomplete()) { if (transactionManager.isAborting()) accumulator.abortUndrainedBatches(new KafkaException("Failing batch since transaction was aborted")); @@ -433,55 +431,48 @@ private void maybeBeginFlush() { if (!accumulator.flushInProgress()) accumulator.beginFlush(); } - } - /** - * Enqueue a FindCoordinatorRequest if needed, otherwise send the next request - */ - private boolean maybeSendTransactionalRequest(TransactionManager.TxnRequestHandler nextRequestHandler) { - AbstractRequest.Builder requestBuilder = nextRequestHandler.requestBuilder(); - while (!forceClose) { - Node targetNode = null; - try { - if (nextRequestHandler.needsCoordinator()) { - targetNode = transactionManager.coordinator(nextRequestHandler.coordinatorType()); - if (targetNode == null) { - transactionManager.lookupCoordinator(nextRequestHandler); - break; - } - if (!NetworkClientUtils.awaitReady(client, targetNode, time, requestTimeoutMs)) { - transactionManager.lookupCoordinator(nextRequestHandler); - break; - } - } else { - targetNode = awaitLeastLoadedNodeReady(requestTimeoutMs); - } + TransactionManager.TxnRequestHandler nextRequestHandler = transactionManager.nextRequestHandler(accumulator.hasIncomplete()); + if (nextRequestHandler == null) + return false; - if (targetNode != null) { - if (nextRequestHandler.isRetry()) - time.sleep(nextRequestHandler.retryBackoffMs()); - long currentTimeMs = time.milliseconds(); - ClientRequest clientRequest = client.newClientRequest( - targetNode.idString(), requestBuilder, currentTimeMs, true, requestTimeoutMs, nextRequestHandler); - log.debug("Sending transactional request {} to node {}", requestBuilder, targetNode); - client.send(clientRequest, currentTimeMs); - transactionManager.setInFlightCorrelationId(clientRequest.correlationId()); - return true; - } - } catch (IOException e) { - log.debug("Disconnect from {} while trying to send request {}. Going " + - "to back off and retry.", targetNode, requestBuilder, e); - if (nextRequestHandler.needsCoordinator()) { - // We break here so that we pick up the FindCoordinator request immediately. - transactionManager.lookupCoordinator(nextRequestHandler); - break; - } + AbstractRequest.Builder requestBuilder = nextRequestHandler.requestBuilder(); + Node targetNode = null; + try { + targetNode = awaitLeastLoadedNodeReady(requestTimeoutMs, nextRequestHandler.coordinatorType()); + if (targetNode == null) { + lookupCoordinatorAndRetry(nextRequestHandler); + return true; } + + if (nextRequestHandler.isRetry()) + time.sleep(nextRequestHandler.retryBackoffMs()); + long currentTimeMs = time.milliseconds(); + ClientRequest clientRequest = client.newClientRequest( + targetNode.idString(), requestBuilder, currentTimeMs, true, requestTimeoutMs, nextRequestHandler); + log.debug("Sending transactional request {} to node {}", requestBuilder, targetNode); + client.send(clientRequest, currentTimeMs); + transactionManager.setInFlightCorrelationId(clientRequest.correlationId()); + client.poll(retryBackoffMs, time.milliseconds()); + return true; + } catch (IOException e) { + log.debug("Disconnect from {} while trying to send request {}. Going " + + "to back off and retry.", targetNode, requestBuilder, e); + // We break here so that we pick up the FindCoordinator request immediately. + lookupCoordinatorAndRetry(nextRequestHandler); + return true; + } + } + + private void lookupCoordinatorAndRetry(TransactionManager.TxnRequestHandler nextRequestHandler) { + if (!nextRequestHandler.needsCoordinator()) { + // For non-coordinator requests, sleep here to prevent a tight loop when no node is available time.sleep(retryBackoffMs); metadata.requestUpdate(); } + + transactionManager.lookupCoordinator(nextRequestHandler); transactionManager.retry(nextRequestHandler); - return false; } private void maybeAbortBatches(RuntimeException exception) { @@ -524,8 +515,11 @@ private ClientResponse sendAndAwaitInitProducerIdRequest(Node node) throws IOExc return NetworkClientUtils.sendAndReceive(client, request, time); } - private Node awaitLeastLoadedNodeReady(long remainingTimeMs) throws IOException { - Node node = client.leastLoadedNode(time.milliseconds()); + private Node awaitLeastLoadedNodeReady(long remainingTimeMs, FindCoordinatorRequest.CoordinatorType coordinatorType) throws IOException { + Node node = coordinatorType != null ? + transactionManager.coordinator(coordinatorType) : + client.leastLoadedNode(time.milliseconds()); + if (node != null && NetworkClientUtils.awaitReady(client, node, time, remainingTimeMs)) { return node; } @@ -536,7 +530,7 @@ private void maybeWaitForProducerId() { while (!forceClose && !transactionManager.hasProducerId() && !transactionManager.hasError()) { Node node = null; try { - node = awaitLeastLoadedNodeReady(requestTimeoutMs); + node = awaitLeastLoadedNodeReady(requestTimeoutMs, null); if (node != null) { ClientResponse response = sendAndAwaitInitProducerIdRequest(node); InitProducerIdResponse initProducerIdResponse = (InitProducerIdResponse) response.responseBody(); 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 ea4d9d0fe60be..a50531a413c0d 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 @@ -944,6 +944,9 @@ private void enqueueRequest(TxnRequestHandler requestHandler) { } private void lookupCoordinator(FindCoordinatorRequest.CoordinatorType type, String coordinatorKey) { + if (type == null) + return; + switch (type) { case GROUP: consumerGroupCoordinator = null; @@ -1057,8 +1060,7 @@ public void onComplete(ClientResponse response) { clearInFlightCorrelationId(); if (response.wasDisconnected()) { log.debug("Disconnected from {}. Will retry.", response.destination()); - if (this.needsCoordinator()) - lookupCoordinator(this.coordinatorType(), this.coordinatorKey()); + lookupCoordinator(this.coordinatorType(), this.coordinatorKey()); reenqueue(); } else if (response.versionMismatch() != null) { fatalError(response.versionMismatch()); From 0cd337605253e12007d68a90238cbfd64b69d6f1 Mon Sep 17 00:00:00 2001 From: Bob Barrett Date: Thu, 18 Jul 2019 11:10:34 -0700 Subject: [PATCH 3/3] Rename awaitLeastLoadedNode, check before calling lookupCoordinator --- .../kafka/clients/producer/internals/Sender.java | 13 +++++++------ .../producer/internals/TransactionManager.java | 6 ++---- 2 files changed, 9 insertions(+), 10 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java index 6cc41c54c309d..efa418c8174bf 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/Sender.java @@ -439,7 +439,7 @@ private boolean maybeSendAndPollTransactionalRequest() { AbstractRequest.Builder requestBuilder = nextRequestHandler.requestBuilder(); Node targetNode = null; try { - targetNode = awaitLeastLoadedNodeReady(requestTimeoutMs, nextRequestHandler.coordinatorType()); + targetNode = awaitNodeReady(nextRequestHandler.coordinatorType()); if (targetNode == null) { lookupCoordinatorAndRetry(nextRequestHandler); return true; @@ -465,13 +465,14 @@ private boolean maybeSendAndPollTransactionalRequest() { } private void lookupCoordinatorAndRetry(TransactionManager.TxnRequestHandler nextRequestHandler) { - if (!nextRequestHandler.needsCoordinator()) { + if (nextRequestHandler.needsCoordinator()) { + transactionManager.lookupCoordinator(nextRequestHandler); + } else { // For non-coordinator requests, sleep here to prevent a tight loop when no node is available time.sleep(retryBackoffMs); metadata.requestUpdate(); } - transactionManager.lookupCoordinator(nextRequestHandler); transactionManager.retry(nextRequestHandler); } @@ -515,12 +516,12 @@ private ClientResponse sendAndAwaitInitProducerIdRequest(Node node) throws IOExc return NetworkClientUtils.sendAndReceive(client, request, time); } - private Node awaitLeastLoadedNodeReady(long remainingTimeMs, FindCoordinatorRequest.CoordinatorType coordinatorType) throws IOException { + private Node awaitNodeReady(FindCoordinatorRequest.CoordinatorType coordinatorType) throws IOException { Node node = coordinatorType != null ? transactionManager.coordinator(coordinatorType) : client.leastLoadedNode(time.milliseconds()); - if (node != null && NetworkClientUtils.awaitReady(client, node, time, remainingTimeMs)) { + if (node != null && NetworkClientUtils.awaitReady(client, node, time, requestTimeoutMs)) { return node; } return null; @@ -530,7 +531,7 @@ private void maybeWaitForProducerId() { while (!forceClose && !transactionManager.hasProducerId() && !transactionManager.hasError()) { Node node = null; try { - node = awaitLeastLoadedNodeReady(requestTimeoutMs, null); + node = awaitNodeReady(null); if (node != null) { ClientResponse response = sendAndAwaitInitProducerIdRequest(node); InitProducerIdResponse initProducerIdResponse = (InitProducerIdResponse) response.responseBody(); 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 a50531a413c0d..ea4d9d0fe60be 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 @@ -944,9 +944,6 @@ private void enqueueRequest(TxnRequestHandler requestHandler) { } private void lookupCoordinator(FindCoordinatorRequest.CoordinatorType type, String coordinatorKey) { - if (type == null) - return; - switch (type) { case GROUP: consumerGroupCoordinator = null; @@ -1060,7 +1057,8 @@ public void onComplete(ClientResponse response) { clearInFlightCorrelationId(); if (response.wasDisconnected()) { log.debug("Disconnected from {}. Will retry.", response.destination()); - lookupCoordinator(this.coordinatorType(), this.coordinatorKey()); + if (this.needsCoordinator()) + lookupCoordinator(this.coordinatorType(), this.coordinatorKey()); reenqueue(); } else if (response.versionMismatch() != null) { fatalError(response.versionMismatch());