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 ff2e5cd8cc57d..010fff81fa591 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 @@ -254,6 +254,17 @@ public class ConsumerConfig extends AbstractConfig { "be excluded from the subscription. It is always possible to explicitly subscribe to an internal topic."; public static final boolean DEFAULT_EXCLUDE_INTERNAL_TOPICS = true; + /** + * internal.leave.group.on.close + * Whether or not the consumer should leave the group on close. If set to false then a rebalance + * won't occur until session.timeout.ms expires. + * + *

+ * Note: this is an internal configuration and could be changed in the future in a backward incompatible way + * + */ + static final String LEAVE_GROUP_ON_CLOSE_CONFIG = "internal.leave.group.on.close"; + /** isolation.level */ public static final String ISOLATION_LEVEL_CONFIG = "isolation.level"; public static final String ISOLATION_LEVEL_DOC = "

Controls how to read messages written transactionally. If set to read_committed, consumer.poll() will only return" + @@ -476,6 +487,10 @@ public class ConsumerConfig extends AbstractConfig { DEFAULT_EXCLUDE_INTERNAL_TOPICS, Importance.MEDIUM, EXCLUDE_INTERNAL_TOPICS_DOC) + .defineInternal(LEAVE_GROUP_ON_CLOSE_CONFIG, + Type.BOOLEAN, + true, + Importance.LOW) .define(ISOLATION_LEVEL_CONFIG, Type.STRING, DEFAULT_ISOLATION_LEVEL, 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 ad7ae82127fe5..c33a52e396109 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 @@ -792,7 +792,8 @@ else if (enableAutoCommit) retryBackoffMs, enableAutoCommit, config.getInt(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG), - this.interceptors); + this.interceptors, + config.getBoolean(ConsumerConfig.LEAVE_GROUP_ON_CLOSE_CONFIG)); this.fetcher = new Fetcher<>( logContext, this.client, diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java index 3af6d057233f9..54678f7ada3eb 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java @@ -133,6 +133,7 @@ private enum MemberState { private Generation generation = Generation.NO_GENERATION; private RequestFuture findCoordinatorFuture = null; + private final boolean leaveGroupOnClose; /** * Initialize the coordination manager. @@ -147,7 +148,8 @@ public AbstractCoordinator(LogContext logContext, Metrics metrics, String metricGrpPrefix, Time time, - long retryBackoffMs) { + long retryBackoffMs, + boolean leaveGroupOnClose) { this.log = logContext.logger(AbstractCoordinator.class); this.client = client; this.time = time; @@ -159,6 +161,7 @@ public AbstractCoordinator(LogContext logContext, this.heartbeat = heartbeat; this.sensors = new GroupCoordinatorMetrics(metrics, metricGrpPrefix); this.retryBackoffMs = retryBackoffMs; + this.leaveGroupOnClose = leaveGroupOnClose; } public AbstractCoordinator(LogContext logContext, @@ -171,10 +174,11 @@ public AbstractCoordinator(LogContext logContext, Metrics metrics, String metricGrpPrefix, Time time, - long retryBackoffMs) { + long retryBackoffMs, + boolean leaveGroupOnClose) { this(logContext, client, groupId, groupInstanceId, rebalanceTimeoutMs, sessionTimeoutMs, new Heartbeat(time, sessionTimeoutMs, heartbeatIntervalMs, rebalanceTimeoutMs, retryBackoffMs), - metrics, metricGrpPrefix, time, retryBackoffMs); + metrics, metricGrpPrefix, time, retryBackoffMs, leaveGroupOnClose); } /** @@ -845,7 +849,9 @@ protected void close(Timer timer) { // Synchronize after closing the heartbeat thread since heartbeat thread // needs this lock to complete and terminate after close flag is set. synchronized (this) { - maybeLeaveGroup(); + if (leaveGroupOnClose) { + maybeLeaveGroup(); + } // At this point, there may be pending commits (async commits or sync commits that were // interrupted using wakeup) and the leave group request which have been queued, but not diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java index 3aef0c5257cf5..b03af74a5dec9 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java @@ -137,7 +137,8 @@ public ConsumerCoordinator(LogContext logContext, long retryBackoffMs, boolean autoCommitEnabled, int autoCommitIntervalMs, - ConsumerInterceptors interceptors) { + ConsumerInterceptors interceptors, + boolean leaveGroupOnClose) { super(logContext, client, groupId, @@ -148,7 +149,8 @@ public ConsumerCoordinator(LogContext logContext, metrics, metricGrpPrefix, time, - retryBackoffMs); + retryBackoffMs, + leaveGroupOnClose); this.log = logContext.logger(ConsumerCoordinator.class); this.metadata = metadata; this.metadataSnapshot = new MetadataSnapshot(subscriptions, metadata.fetch(), metadata.updateVersion()); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java index 9012ea25e9166..42cccd49a9be3 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java @@ -1902,7 +1902,8 @@ private KafkaConsumer newConsumer(Time time, retryBackoffMs, autoCommitEnabled, autoCommitIntervalMs, - interceptors); + interceptors, + true); Fetcher fetcher = new Fetcher<>( loggerFactory, diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java index 31328b390ad44..0fc5f62b69adb 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java @@ -857,7 +857,7 @@ public DummyCoordinator(ConsumerNetworkClient client, int retryBackoffMs, Optional groupInstanceId) { super(new LogContext(), client, GROUP_ID, groupInstanceId, rebalanceTimeoutMs, - SESSION_TIMEOUT_MS, HEARTBEAT_INTERVAL_MS, metrics, METRIC_GROUP_PREFIX, time, retryBackoffMs); + SESSION_TIMEOUT_MS, HEARTBEAT_INTERVAL_MS, metrics, METRIC_GROUP_PREFIX, time, retryBackoffMs, !groupInstanceId.isPresent()); } @Override diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinatorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinatorTest.java index f0214d2228764..86032c4f7c927 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinatorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinatorTest.java @@ -2219,7 +2219,8 @@ private ConsumerCoordinator buildCoordinator(final Metrics metrics, retryBackoffMs, autoCommitEnabled, autoCommitIntervalMs, - null + null, + !groupInstanceId.isPresent() ); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java index 706742aad8842..fd7c7a429f26a 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerCoordinator.java @@ -95,7 +95,8 @@ public WorkerCoordinator(LogContext logContext, metrics, metricGrpPrefix, time, - retryBackoffMs); + retryBackoffMs, + true); this.log = logContext.logger(WorkerCoordinator.class); this.restUrl = restUrl; this.configStorage = configStorage; diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index 5024c28143347..6d93b9977df87 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -720,6 +720,7 @@ public class StreamsConfig extends AbstractConfig { tempConsumerDefaultOverrides.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "1000"); tempConsumerDefaultOverrides.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); tempConsumerDefaultOverrides.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + tempConsumerDefaultOverrides.put("internal.leave.group.on.close", false); CONSUMER_DEFAULT_OVERRIDES = Collections.unmodifiableMap(tempConsumerDefaultOverrides); } diff --git a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java index c202c93cec776..5f053bca1ca62 100644 --- a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java @@ -31,6 +31,7 @@ import org.apache.kafka.streams.processor.FailOnInvalidTimestamp; import org.apache.kafka.streams.processor.TimestampExtractor; import org.apache.kafka.streams.processor.internals.StreamsPartitionAssignor; +import org.hamcrest.CoreMatchers; import org.junit.Before; import org.junit.Test; @@ -424,6 +425,13 @@ public void testGetGlobalConsumerConfigsWithGlobalConsumerOverridenPrefix() { assertEquals("50", returnedProps.get(ConsumerConfig.MAX_POLL_RECORDS_CONFIG)); } + @Test + public void shouldSetInternalLeaveGroupOnCloseConfigToFalseInConsumer() { + final StreamsConfig streamsConfig = new StreamsConfig(props); + final Map consumerConfigs = streamsConfig.getMainConsumerConfigs(groupId, clientId, threadIdx); + assertThat(consumerConfigs.get("internal.leave.group.on.close"), CoreMatchers.equalTo(false)); + } + @Test public void shouldAcceptAtLeastOnce() { // don't use `StreamsConfig.AT_LEAST_ONCE` to actually do a useful test