From 1e5123f4323a6c9b6b17e3bb7aa5bff8196c21bc Mon Sep 17 00:00:00 2001 From: Xiongqi Wesley Wu Date: Mon, 12 Apr 2021 18:28:41 -0700 Subject: [PATCH] [LI-HOTFIX] revert sticky metadata fetch hotfix TICKET = LI_DESCRIPTION = revert LI HOTFIX "Consumer should fetch metadata from the same broker until it is disconnected from that broker". since sticky metadata request is no longer needed. The original hotfix was introduced as a short term fix util the actual issue was found, and we no longer need this hotfix now. EXIT_CRITERIA = MANUAL ["after the original fix has been removed"] --- .../kafka/clients/CommonClientConfigs.java | 3 --- .../apache/kafka/clients/NetworkClient.java | 23 +++---------------- .../clients/consumer/ConsumerConfig.java | 5 ---- .../kafka/clients/consumer/KafkaConsumer.java | 1 - .../clients/producer/ProducerConfig.java | 5 ---- 5 files changed, 3 insertions(+), 34 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java b/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java index 03b9b31716c3a..b19704d67aca5 100644 --- a/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java +++ b/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java @@ -147,9 +147,6 @@ public class CommonClientConfigs { + "consumer's session stays active and to facilitate rebalancing when new consumers join or leave the group. " + "The value must be set lower than session.timeout.ms, but typically should be set no higher " + "than 1/3 of that value. It can be adjusted even lower to control the expected time for normal rebalances."; - public static final String ENABLE_STICKY_METADATA_FETCH_CONFIG = "enable.sticky.metadata.fetch"; - public static final String ENABLE_STICKY_METADATA_FETCH_DOC = "Fetch metadata from the least loaded broker if false. Otherwise fetch metadata " - + "from the same broker until it is disconnected."; /** * Postprocess the configuration so that exponential backoff is disabled when reconnect backoff diff --git a/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java b/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java index bc925c0625c6d..e6ca884e9903d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java @@ -115,8 +115,6 @@ private enum State { private final Time time; - private boolean enableStickyMetadataFetch = true; - /** * True if we should send an ApiVersionRequest when first connecting to a broker. */ @@ -275,10 +273,6 @@ private NetworkClient(MetadataUpdater metadataUpdater, this.state = new AtomicReference<>(State.ACTIVE); } - public void setEnableStickyMetadataFetch(boolean enableStickyMetadataFetch) { - this.enableStickyMetadataFetch = enableStickyMetadataFetch; - } - /** * Begin connecting to the given node, return true if we are already connected and ready to send to that node. * @@ -577,12 +571,6 @@ public List poll(long timeout, long now) { handleTimedOutRequests(responses, updatedNow); completeResponses(responses); - // We changed the metadataUpdater.maybeUpdate() such that it will keep sending MetadataRequest - // to the same broker instead choosing the least loaded node. If we don't try to send metadata here, it is possible that - // another request is sent to the broker before the next networkClient.poll(). This can cause starvation - // for the MetadataRequest and consumer's metadata may be stale for a long time. - metadataUpdater.maybeUpdate(updatedNow); - return responses; } @@ -994,8 +982,6 @@ class DefaultMetadataUpdater implements MetadataUpdater { /* the current cluster metadata */ private final Metadata metadata; - // Consumer needs to keep fetching metadata from the same node until that node goes down - private Node nodeToFetchMetadata; // Defined if there is a request in progress, null otherwise private Integer inProgressRequestVersion; @@ -1003,7 +989,6 @@ class DefaultMetadataUpdater implements MetadataUpdater { DefaultMetadataUpdater(Metadata metadata) { this.metadata = metadata; this.inProgressRequestVersion = null; - this.nodeToFetchMetadata = null; } @Override @@ -1033,15 +1018,13 @@ public long maybeUpdate(long now) { // Beware that the behavior of this method and the computation of timeouts for poll() are // highly dependent on the behavior of leastLoadedNode. - if (!enableStickyMetadataFetch || nodeToFetchMetadata == null || !connectionStates.isReady(nodeToFetchMetadata.idString(), now)) - nodeToFetchMetadata = leastLoadedNode(now); - - if (nodeToFetchMetadata == null) { + Node node = leastLoadedNode(now); + if (node == null) { log.debug("Give up sending metadata request since no node is available"); return reconnectBackoffMs; } - return maybeUpdate(now, nodeToFetchMetadata); + return maybeUpdate(now, node); } @Override diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java index ddfa30c9f6121..da797ecd3e8ea 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java @@ -524,11 +524,6 @@ public class ConsumerConfig extends AbstractConfig { CommonClientConfigs.DEFAULT_SECURITY_PROTOCOL, Importance.MEDIUM, CommonClientConfigs.SECURITY_PROTOCOL_DOC) - .define(CommonClientConfigs.ENABLE_STICKY_METADATA_FETCH_CONFIG, - Type.BOOLEAN, - true, - Importance.MEDIUM, - CommonClientConfigs.ENABLE_STICKY_METADATA_FETCH_DOC) .withClientSslSupport() .withClientSaslSupport(); } diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java index 7fd7d2c3bf52f..a4dfa6658867e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java @@ -762,7 +762,6 @@ else if (enableAutoCommit) apiVersions, throttleTimeSensor, logContext); - netClient.setEnableStickyMetadataFetch(config.getBoolean(CommonClientConfigs.ENABLE_STICKY_METADATA_FETCH_CONFIG)); this.client = new ConsumerNetworkClient( logContext, netClient, diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java index b8bf54012ce46..cd3ff01095c5d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java @@ -373,11 +373,6 @@ public class ProducerConfig extends AbstractConfig { null, Importance.LOW, SECURITY_PROVIDERS_DOC) - .define(CommonClientConfigs.ENABLE_STICKY_METADATA_FETCH_CONFIG, - Type.BOOLEAN, - false, - Importance.MEDIUM, - CommonClientConfigs.ENABLE_STICKY_METADATA_FETCH_DOC) .withClientSslSupport() .withClientSaslSupport() .define(ENABLE_IDEMPOTENCE_CONFIG,