This is an automated email from the ASF dual-hosted git repository.
mjsax pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new b91561ac8ab KAFKA-20744: Add back `rack.aware.assignment.tags` config
(#22213)
b91561ac8ab is described below
commit b91561ac8ab325986a80e34ee20864e5d5392327
Author: gabriellafu <[email protected]>
AuthorDate: Sat Jul 25 21:07:58 2026 -0400
KAFKA-20744: Add back `rack.aware.assignment.tags` config (#22213)
This PR adds broker- and client-side plumbing for
`rack.aware.assignment.tags`
to support rack-aware standby task assignment in Kafka Streams (new
stream" group protocol via KIP-1071). This PR does not include the assignor
logic change -- only the config propagation and client-side validation.
Adds `group.streams.rack.aware.assignment.tags` as a broker-level config
with a per-group override `streams.rack.aware.assignment.tags`. Adds
RackAwareAssignmentTags field to the StreamsGroupHeartbeatResponse
RPC schema and populate it in the heartbeat response. On the client
side, compares the broker's required tags against the client's client.tag.*
keys and logs a WARN for any missing tags.
Reviewers: ChickenchickenLove (@chickenchickenlove ), Lucas Brutschy
<[email protected]>, Matthias J. Sax <[email protected]>
---
checkstyle/suppressions.xml | 2 +-
.../StreamsGroupHeartbeatRequestManager.java | 43 ++-
.../requests/StreamsGroupHeartbeatResponse.java | 3 +-
.../message/StreamsGroupHeartbeatRequest.json | 2 +
.../message/StreamsGroupHeartbeatResponse.json | 4 +-
.../StreamsGroupHeartbeatRequestManagerTest.java | 180 ++++++++++
.../scala/unit/kafka/server/KafkaApisTest.scala | 3 +-
.../scala/unit/kafka/server/KafkaConfigTest.scala | 1 +
.../server/StreamsGroupHeartbeatRequestTest.scala | 57 ++-
.../kafka/coordinator/group/GroupConfig.java | 52 ++-
.../coordinator/group/GroupCoordinatorConfig.java | 47 ++-
.../coordinator/group/GroupMetadataManager.java | 40 ++-
.../kafka/coordinator/group/GroupConfigTest.java | 37 ++
.../group/GroupCoordinatorConfigTest.java | 38 ++
.../group/GroupMetadataManagerTest.java | 382 +++++++++++++++++----
15 files changed, 804 insertions(+), 87 deletions(-)
diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index b1c71df9861..3823ff81ffc 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -89,7 +89,7 @@
files="(Sender|Fetcher|FetchRequestManager|OffsetFetcher|KafkaConsumer|Metrics|RequestResponse|TransactionManager|KafkaAdminClient|Message|KafkaProducer)Test.java"/>
<suppress checks="ClassFanOutComplexity"
-
files="(ConsumerCoordinator|KafkaConsumer|RequestResponse|Fetcher|FetchRequestManager|KafkaAdminClient|Message|KafkaProducer|NetworkClient)Test.java"/>
+
files="(ConsumerCoordinator|KafkaConsumer|RequestResponse|Fetcher|FetchRequestManager|KafkaAdminClient|Message|KafkaProducer|NetworkClient|StreamsGroupHeartbeatRequestManager)Test.java"/>
<suppress checks="ClassFanOutComplexity"
files="MockAdminClient.java"/>
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
index 5ba7d1deb20..3b6aff6b30d 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
@@ -376,6 +376,8 @@ public class StreamsGroupHeartbeatRequestManager implements
RequestManager {
private final StreamsRebalanceData streamsRebalanceData;
+ private String lastMissingClientTagsDetail = null;
+
/**
* Timer for tracking the time since the last consumer poll. If the timer
expires, the consumer will stop
* sending heartbeat until the next poll.
@@ -678,20 +680,41 @@ public class StreamsGroupHeartbeatRequestManager
implements RequestManager {
streamsRebalanceData.setPartitionsByHost(convertHostInfoMap(data));
}
- List<StreamsGroupHeartbeatResponseData.Status> statuses =
data.status();
- if (statuses != null) {
- streamsRebalanceData.setStatuses(statuses);
- if (!statuses.isEmpty()) {
- String statusDetails = statuses.stream()
- .map(status -> "(" + status.statusCode() + ") " +
status.statusDetail())
- .collect(Collectors.joining(", "));
- logger.warn("Membership is in the following statuses: {}",
statusDetails);
- }
- }
+ maybeLogStatuses(data.status());
membershipManager.onHeartbeatSuccess(response);
}
+ private void maybeLogStatuses(final
List<StreamsGroupHeartbeatResponseData.Status> statuses) {
+ if (statuses == null) {
+ return;
+ }
+ streamsRebalanceData.setStatuses(statuses);
+ // The broker recomputes and returns the full set of statuses on every
heartbeat, so a response without a
+ // MISSING_CLIENT_TAGS status means the condition no longer holds.
+ boolean hasMissingClientTagsStatus = false;
+ List<String> statusesToLog = new ArrayList<>();
+ for (StreamsGroupHeartbeatResponseData.Status status : statuses) {
+ if (status.statusCode() ==
StreamsGroupHeartbeatResponse.Status.MISSING_CLIENT_TAGS.code()) {
+ hasMissingClientTagsStatus = true;
+ if
(!status.statusDetail().equals(lastMissingClientTagsDetail)) {
+ lastMissingClientTagsDetail = status.statusDetail();
+ statusesToLog.add("(" + status.statusCode() + ") " +
status.statusDetail());
+ }
+ } else {
+ statusesToLog.add("(" + status.statusCode() + ") " +
status.statusDetail());
+ }
+ }
+ // Reset the de-duplication marker once the MISSING_CLIENT_TAGS status
clears, so that a later recurrence
+ // (even with the same detail) is logged again rather than silently
suppressed.
+ if (!hasMissingClientTagsStatus) {
+ lastMissingClientTagsDetail = null;
+ }
+ if (!statusesToLog.isEmpty()) {
+ logger.warn("Membership is in the following statuses: {}",
String.join(", ", statusesToLog));
+ }
+ }
+
// Renders a coordinator-provided config value for logging, or a note when
the broker did not provide it. An older
// broker leaves these at their protocol defaults (intervals 0,
acceptableRecoveryLag -1); a value below minValid
// means "not provided".
diff --git
a/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatResponse.java
b/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatResponse.java
index 68f2a460c4f..3edceabb8b0 100644
---
a/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatResponse.java
+++
b/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatResponse.java
@@ -83,7 +83,8 @@ public class StreamsGroupHeartbeatResponse extends
AbstractResponse {
INCORRECTLY_PARTITIONED_TOPICS((byte) 2, "One or more topics expected
to be copartitioned are not copartitioned."),
MISSING_INTERNAL_TOPICS((byte) 3, "One or more internal topics are
missing."),
SHUTDOWN_APPLICATION((byte) 4, "A client requested the shutdown of the
whole application."),
- ASSIGNMENT_DELAYED((byte) 5, "The assignment was delayed by the
coordinator.");
+ ASSIGNMENT_DELAYED((byte) 5, "The assignment was delayed by the
coordinator."),
+ MISSING_CLIENT_TAGS((byte) 6, "One or more required client tags for
rack-aware standby task assignment are missing.");
private static final Map<Byte, Status> CODE_TO_STATUS;
diff --git
a/clients/src/main/resources/common/message/StreamsGroupHeartbeatRequest.json
b/clients/src/main/resources/common/message/StreamsGroupHeartbeatRequest.json
index e901ad5fa56..86d2b722147 100644
---
a/clients/src/main/resources/common/message/StreamsGroupHeartbeatRequest.json
+++
b/clients/src/main/resources/common/message/StreamsGroupHeartbeatRequest.json
@@ -20,6 +20,8 @@
"name": "StreamsGroupHeartbeatRequest",
// Version 1 is the same as version 0; bumped together with
StreamsGroupHeartbeatResponse v1,
// which adds TopologyDescriptionRequired (KIP-1331). Required so the
response v1 is negotiated.
+ // Version 1 is also required for the broker to return the
MISSING_CLIENT_TAGS status (code 6),
+ // since version 0 clients do not understand it (KAFKA-20744).
"validVersions": "0-1",
"flexibleVersions": "0+",
"fields": [
diff --git
a/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
b/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
index 34bdaf2e2d3..a63ddbb4af2 100644
---
a/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
+++
b/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
@@ -17,7 +17,8 @@
"apiKey": 88,
"type": "response",
"name": "StreamsGroupHeartbeatResponse",
- // Version 1 adds TopologyDescriptionRequired (KIP-1331).
+ // Version 1 adds TopologyDescriptionRequired (KIP-1331), and allows the
broker to return the
+ // MISSING_CLIENT_TAGS status (code 6), which version 0 clients do not
understand (KAFKA-20744).
"validVersions": "0-1",
"flexibleVersions": "0+",
// Supported errors:
@@ -102,6 +103,7 @@
// topic creation, this will be
indicated in StatusDetail.
// 4 - SHUTDOWN_APPLICATION - A client requested the shutdown
of the whole application.
// 5 - ASSIGNMENT_DELAYED - No assignment was provided
because assignment computation was delayed.
+ // 6 - MISSING_CLIENT_TAGS - One or more required client
tags for rack-aware standby task assignment are missing.
{ "name": "StatusCode", "type": "int8", "versions": "0+",
"about": "A code to indicate that a particular status is active for
the group membership" },
{ "name": "StatusDetail", "type": "string", "versions": "0+",
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
index 860d3439f8f..c87087c67a2 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
@@ -2707,6 +2707,186 @@ class StreamsGroupHeartbeatRequestManagerTest {
);
}
+ @Test
+ public void testMissingClientTagsStatusLogsWarningOnlyOnce() {
+ try (
+ final MockedConstruction<HeartbeatRequestState> ignored =
mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+ final LogCaptureAppender logAppender =
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+ ) {
+
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class,
Level.WARN);
+ final StreamsGroupHeartbeatRequestManager heartbeatRequestManager
= createStreamsGroupHeartbeatRequestManager();
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+ when(membershipManager.groupId()).thenReturn(GROUP_ID);
+ when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+ when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+ final String statusDetail = "Missing required client tags for
rack-aware standby assignment: [zone, cluster]";
+
+ // First heartbeat with MISSING_CLIENT_TAGS status
+ final NetworkClientDelegate.PollResult result1 =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result1.unsentRequests.size());
+
+ final ClientResponse response1 = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null, "-1", time.milliseconds(), time.milliseconds(), false,
null, null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setStatus(List.of(new
StreamsGroupHeartbeatResponseData.Status()
+
.setStatusCode(StreamsGroupHeartbeatResponse.Status.MISSING_CLIENT_TAGS.code())
+ .setStatusDetail(statusDetail)))
+ )
+ );
+ result1.unsentRequests.get(0).handler().onComplete(response1);
+
+ long firstWarnCount = logAppender.getMessages("WARN").stream()
+ .filter(m -> m.contains("Missing required client tags"))
+ .count();
+ assertEquals(1, firstWarnCount);
+ assertTrue(logAppender.getMessages("WARN").stream().anyMatch(m ->
m.contains("[zone, cluster]")),
+ "The logged warning should contain the missing client tags
detail [zone, cluster]");
+
+ // Second heartbeat with the same status — should NOT log again
+ final NetworkClientDelegate.PollResult result2 =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result2.unsentRequests.size());
+
+ final ClientResponse response2 = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null, "-1", time.milliseconds(), time.milliseconds(), false,
null, null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setStatus(List.of(new
StreamsGroupHeartbeatResponseData.Status()
+
.setStatusCode(StreamsGroupHeartbeatResponse.Status.MISSING_CLIENT_TAGS.code())
+ .setStatusDetail(statusDetail)))
+ )
+ );
+ result2.unsentRequests.get(0).handler().onComplete(response2);
+
+ long secondWarnCount = logAppender.getMessages("WARN").stream()
+ .filter(m -> m.contains("Missing required client tags"))
+ .count();
+ assertEquals(1, secondWarnCount, "MISSING_CLIENT_TAGS warning
should not be logged again for the same detail");
+
+ // Third heartbeat with a DIFFERENT status detail — should log
again
+ final String changedStatusDetail = "Missing required client tags
for rack-aware standby assignment: [zone]";
+
+ final NetworkClientDelegate.PollResult result3 =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result3.unsentRequests.size());
+
+ final ClientResponse response3 = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null, "-1", time.milliseconds(), time.milliseconds(), false,
null, null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setStatus(List.of(new
StreamsGroupHeartbeatResponseData.Status()
+
.setStatusCode(StreamsGroupHeartbeatResponse.Status.MISSING_CLIENT_TAGS.code())
+ .setStatusDetail(changedStatusDetail)))
+ )
+ );
+ result3.unsentRequests.get(0).handler().onComplete(response3);
+
+ List<String> missingTagWarnings =
logAppender.getMessages("WARN").stream()
+ .filter(m -> m.contains("Missing required client tags"))
+ .collect(Collectors.toList());
+ assertEquals(2, missingTagWarnings.size(),
+ "MISSING_CLIENT_TAGS warning should be logged again when the
detail changes");
+ // The second log line must reflect only the new detail: it
contains [zone] and must no
+ // longer report the previous [zone, cluster] detail.
+ String secondWarning = missingTagWarnings.get(1);
+ assertTrue(secondWarning.contains("[zone]"),
+ "The second logged warning should contain the changed missing
client tags detail [zone]");
+ assertFalse(secondWarning.contains("[zone, cluster]"),
+ "The second logged warning should not contain the stale detail
[zone, cluster]");
+
+ // Fourth heartbeat with the status cleared (e.g. broker reverted
its required tags) — nothing to log,
+ // but the de-duplication marker should be reset.
+ final NetworkClientDelegate.PollResult result4 =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result4.unsentRequests.size());
+
+ final ClientResponse response4 = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null, "-1", time.milliseconds(), time.milliseconds(), false,
null, null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setStatus(List.of())
+ )
+ );
+ result4.unsentRequests.get(0).handler().onComplete(response4);
+
+ long fourthWarnCount = logAppender.getMessages("WARN").stream()
+ .filter(m -> m.contains("Missing required client tags"))
+ .count();
+ assertEquals(2, fourthWarnCount, "Clearing the status should not
log a new warning");
+
+ // Fifth heartbeat with the status recurring with the
previously-seen detail — should log again because
+ // the marker was reset when the status cleared.
+ final NetworkClientDelegate.PollResult result5 =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result5.unsentRequests.size());
+
+ final ClientResponse response5 = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null, "-1", time.milliseconds(), time.milliseconds(), false,
null, null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setStatus(List.of(new
StreamsGroupHeartbeatResponseData.Status()
+
.setStatusCode(StreamsGroupHeartbeatResponse.Status.MISSING_CLIENT_TAGS.code())
+ .setStatusDetail(changedStatusDetail)))
+ )
+ );
+ result5.unsentRequests.get(0).handler().onComplete(response5);
+
+ long fifthWarnCount = logAppender.getMessages("WARN").stream()
+ .filter(m -> m.contains("Missing required client tags"))
+ .count();
+ assertEquals(3, fifthWarnCount, "MISSING_CLIENT_TAGS warning
should be logged again after the status cleared and recurred");
+ }
+ }
+
+ @Test
+ public void testNoWarningWhenClientTagsPresent() {
+ try (
+ final MockedConstruction<HeartbeatRequestState> ignored =
mockConstruction(
+ HeartbeatRequestState.class,
+ (mock, context) ->
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+ final LogCaptureAppender logAppender =
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+ ) {
+
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class,
Level.WARN);
+ final StreamsGroupHeartbeatRequestManager heartbeatRequestManager
= createStreamsGroupHeartbeatRequestManager();
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+ when(membershipManager.groupId()).thenReturn(GROUP_ID);
+ when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+ when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+ // The client supplies all required rack-aware tags (e.g. zone and
cluster), so the broker returns a
+ // heartbeat with no MISSING_CLIENT_TAGS status and nothing should
be logged.
+ final NetworkClientDelegate.PollResult result =
heartbeatRequestManager.poll(time.milliseconds());
+ assertEquals(1, result.unsentRequests.size());
+
+ final ClientResponse response = new ClientResponse(
+ new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1,
"", 1),
+ null, "-1", time.milliseconds(), time.milliseconds(), false,
null, null,
+ new StreamsGroupHeartbeatResponse(
+ new StreamsGroupHeartbeatResponseData()
+ .setHeartbeatIntervalMs((int)
RECEIVED_HEARTBEAT_INTERVAL_MS)
+ .setStatus(List.of())
+ )
+ );
+ result.unsentRequests.get(0).handler().onComplete(response);
+
+ assertTrue(logAppender.getMessages("WARN").stream()
+ .noneMatch(m -> m.contains("Missing required client
tags")),
+ "No MISSING_CLIENT_TAGS warning should be logged when the
client provides the required tags");
+ }
+ }
+
private static void assertTaskIdsEquals(final
List<StreamsGroupHeartbeatRequestData.TaskIds> expected,
final
List<StreamsGroupHeartbeatRequestData.TaskIds> actual) {
List<StreamsGroupHeartbeatRequestData.TaskIds> sortedExpected =
expected.stream()
diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
index 8d4e6f2a671..c6154958dd0 100644
--- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
@@ -80,7 +80,7 @@ import org.apache.kafka.common.utils.Utils
import org.apache.kafka.common.utils.internals.ImplicitLinkedHashCollection
import org.apache.kafka.common.utils.internals.ProducerIdAndEpoch
import org.apache.kafka.common.utils.internals.SecurityUtils
-import
org.apache.kafka.coordinator.group.GroupConfig.{CONSUMER_ASSIGNMENT_INTERVAL_MS_CONFIG,
CONSUMER_ASSIGNOR_OFFLOAD_ENABLE_CONFIG,
CONSUMER_HEARTBEAT_INTERVAL_MS_CONFIG, CONSUMER_SESSION_TIMEOUT_MS_CONFIG,
SHARE_ASSIGNMENT_INTERVAL_MS_CONFIG, SHARE_ASSIGNOR_OFFLOAD_ENABLE_CONFIG,
SHARE_AUTO_OFFSET_RESET_CONFIG, SHARE_DELIVERY_COUNT_LIMIT_CONFIG,
SHARE_HEARTBEAT_INTERVAL_MS_CONFIG, SHARE_ISOLATION_LEVEL_CONFIG,
SHARE_PARTITION_MAX_RECORD_LOCKS_CONFIG, SHARE_RECORD_LOCK_DURATION_MS_CO [...]
+import
org.apache.kafka.coordinator.group.GroupConfig.{CONSUMER_ASSIGNMENT_INTERVAL_MS_CONFIG,
CONSUMER_ASSIGNOR_OFFLOAD_ENABLE_CONFIG,
CONSUMER_HEARTBEAT_INTERVAL_MS_CONFIG, CONSUMER_SESSION_TIMEOUT_MS_CONFIG,
SHARE_ASSIGNMENT_INTERVAL_MS_CONFIG, SHARE_ASSIGNOR_OFFLOAD_ENABLE_CONFIG,
SHARE_AUTO_OFFSET_RESET_CONFIG, SHARE_DELIVERY_COUNT_LIMIT_CONFIG,
SHARE_HEARTBEAT_INTERVAL_MS_CONFIG, SHARE_ISOLATION_LEVEL_CONFIG,
SHARE_PARTITION_MAX_RECORD_LOCKS_CONFIG, SHARE_RECORD_LOCK_DURATION_MS_CO [...]
import org.apache.kafka.coordinator.group.modern.share.ShareGroupConfig
import org.apache.kafka.coordinator.group.{GroupConfig, GroupConfigManager,
GroupCoordinator, GroupCoordinatorConfig}
import org.apache.kafka.coordinator.group.streams.StreamsGroupHeartbeatResult
@@ -383,6 +383,7 @@ class KafkaApisTest extends Logging {
cgConfigs.put(STREAMS_ASSIGNOR_OFFLOAD_ENABLE_CONFIG, "false")
cgConfigs.put(STREAMS_TASK_OFFSET_INTERVAL_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_TASK_OFFSET_INTERVAL_MS_DEFAULT.toString)
cgConfigs.put(STREAMS_NUM_WARMUP_REPLICAS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_NUM_WARMUP_REPLICAS_DEFAULT.toString)
+ cgConfigs.put(STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DEFAULT)
cgConfigs.put(STREAMS_ACCEPTABLE_RECOVERY_LAG_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_DEFAULT.toString)
cgConfigs.put(ERRORS_DEADLETTERQUEUE_TOPIC_NAME_CONFIG, "")
cgConfigs.put(ERRORS_DEADLETTERQUEUE_COPY_RECORD_ENABLE_CONFIG, "false")
diff --git a/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
b/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
index 523675c5293..5cdb38c79af 100755
--- a/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
@@ -1110,6 +1110,7 @@ class KafkaConfigTest {
case
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG =>
assertPropertyInvalid(baseProperties, name, "not_a_number", -1)
case
GroupCoordinatorConfig.STREAMS_GROUP_TASK_OFFSET_INTERVAL_MS_CONFIG =>
assertPropertyInvalid(baseProperties, name, "not_a_number", -1)
case
GroupCoordinatorConfig.STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_CONFIG =>
assertPropertyInvalid(baseProperties, name, "not_a_number", -1)
+ case
GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG => //
ignore list
case
GroupCoordinatorConfig.STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_CONFIG =>
assertPropertyInvalid(baseProperties, name, "not_a_number", -1)
/** Share coordinator configs */
diff --git
a/core/src/test/scala/unit/kafka/server/StreamsGroupHeartbeatRequestTest.scala
b/core/src/test/scala/unit/kafka/server/StreamsGroupHeartbeatRequestTest.scala
index 16811eee758..281abf59e88 100644
---
a/core/src/test/scala/unit/kafka/server/StreamsGroupHeartbeatRequestTest.scala
+++
b/core/src/test/scala/unit/kafka/server/StreamsGroupHeartbeatRequestTest.scala
@@ -23,11 +23,12 @@ import
org.apache.kafka.common.message.{StreamsGroupHeartbeatRequestData, Stream
import org.apache.kafka.common.protocol.Errors
import org.apache.kafka.common.test.ClusterInstance
import org.apache.kafka.common.test.api.{ClusterConfigProperty,
ClusterFeature, ClusterTest, ClusterTestDefaults, Type}
-import org.apache.kafka.coordinator.group.GroupCoordinatorConfig
-import org.apache.kafka.common.errors.UnsupportedVersionException
+import org.apache.kafka.coordinator.group.{GroupConfig, GroupCoordinatorConfig}
+import org.apache.kafka.common.errors.{InvalidConfigurationException,
UnsupportedVersionException}
import org.apache.kafka.server.common.Feature
import org.junit.jupiter.api.Assertions.{assertEquals, assertNotNull,
assertNull, assertThrows, assertTrue}
+import java.util.concurrent.ExecutionException
import scala.jdk.CollectionConverters._
object StreamsGroupHeartbeatRequestTest {
@@ -873,6 +874,58 @@ class StreamsGroupHeartbeatRequestTest(cluster:
ClusterInstance) extends GroupCo
}
}
+ @ClusterTest
+ def testAlterRackAwareAssignmentTagsGroupConfig(): Unit = {
+ val admin = cluster.admin()
+ val groupId = "test-group"
+
+ try {
+ TestUtils.createOffsetsTopicWithAdmin(
+ admin = admin,
+ brokers = cluster.brokers.values().asScala.toSeq,
+ controllers = cluster.controllers().values().asScala.toSeq
+ )
+
+ val groupConfigResource = new ConfigResource(ConfigResource.Type.GROUP,
groupId)
+
+ // Valid case: a list of distinct client tag keys is accepted.
+ val validAlterOp = new AlterConfigOp(
+ new ConfigEntry(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster"),
+ AlterConfigOp.OpType.SET
+ )
+ admin.incrementalAlterConfigs(
+ Map(groupConfigResource -> List(validAlterOp).asJavaCollection).asJava
+ ).all().get()
+
+ // Verify the config was persisted. Group config propagation is
asynchronous, so wait for it.
+ TestUtils.waitUntilTrue(() => {
+ val describedConfigs =
admin.describeConfigs(List(groupConfigResource).asJava).all().get()
+ val rackAwareTags =
describedConfigs.get(groupConfigResource).get(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG)
+ rackAwareTags != null && rackAwareTags.value() == "zone,cluster"
+ }, s"${GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG} was not
updated to the expected value within the timeout period.")
+
+ // Invalid case: an empty tag key (here an empty element between commas)
is rejected with INVALID_CONFIG.
+ val invalidAlterOp = new AlterConfigOp(
+ new ConfigEntry(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,,cluster"),
+ AlterConfigOp.OpType.SET
+ )
+ val executionException = assertThrows(classOf[ExecutionException], () =>
+ admin.incrementalAlterConfigs(
+ Map(groupConfigResource ->
List(invalidAlterOp).asJavaCollection).asJava
+ ).all().get()
+ )
+
assertTrue(executionException.getCause.isInstanceOf[InvalidConfigurationException],
+ s"Expected InvalidConfigurationException but got
${executionException.getCause}")
+
+ // The invalid alter is rejected at validation time, so the previously
persisted valid value is unchanged.
+ val describedConfigsAfter =
admin.describeConfigs(List(groupConfigResource).asJava).all().get()
+ assertEquals("zone,cluster",
+
describedConfigsAfter.get(groupConfigResource).get(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG).value())
+ } finally {
+ admin.close()
+ }
+ }
+
@ClusterTest(
types = Array(Type.KRAFT),
serverProperties = Array(
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupConfig.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupConfig.java
index 9d4a8ad9f9b..24b547a633b 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupConfig.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupConfig.java
@@ -27,16 +27,19 @@ import
org.apache.kafka.coordinator.group.modern.share.ShareGroupConfig;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Optional;
import java.util.Properties;
import java.util.Set;
+import static org.apache.kafka.common.config.ConfigDef.Importance.LOW;
import static org.apache.kafka.common.config.ConfigDef.Importance.MEDIUM;
import static org.apache.kafka.common.config.ConfigDef.Range.atLeast;
import static org.apache.kafka.common.config.ConfigDef.Type.BOOLEAN;
import static org.apache.kafka.common.config.ConfigDef.Type.INT;
+import static org.apache.kafka.common.config.ConfigDef.Type.LIST;
import static org.apache.kafka.common.config.ConfigDef.Type.LONG;
import static org.apache.kafka.common.config.ConfigDef.Type.STRING;
import static org.apache.kafka.common.config.ConfigDef.ValidString.in;
@@ -102,6 +105,8 @@ public final class GroupConfig extends AbstractConfig {
public static final String STREAMS_ASSIGNMENT_INTERVAL_MS_CONFIG =
"streams.assignment.interval.ms";
+ public static final String STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG =
"streams.rack.aware.assignment.tags";
+
public static final String STREAMS_ASSIGNOR_OFFLOAD_ENABLE_CONFIG =
"streams.assignor.offload.enable";
public static final String STREAMS_TASK_OFFSET_INTERVAL_MS_CONFIG =
"streams.task.offset.interval.ms";
@@ -152,6 +157,8 @@ public final class GroupConfig extends AbstractConfig {
private final Optional<Integer> streamsAssignmentIntervalMs;
+ private final Optional<List<String>> streamsRackAwareAssignmentTags;
+
private final Optional<Boolean> streamsAssignorOffloadEnable;
private final Optional<Integer> streamsTaskOffsetIntervalMs;
@@ -297,6 +304,12 @@ public final class GroupConfig extends AbstractConfig {
atLeast(0),
MEDIUM,
GroupCoordinatorConfig.STREAMS_GROUP_NUM_WARMUP_REPLICAS_DOC)
+ .define(STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
+ LIST,
+
GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DEFAULT,
+ ConfigDef.ValidList.anyNonDuplicateValues(true, false),
+ LOW,
+
GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DOC)
.define(STREAMS_ACCEPTABLE_RECOVERY_LAG_CONFIG,
LONG,
GroupCoordinatorConfig.STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_DEFAULT,
@@ -348,6 +361,7 @@ public final class GroupConfig extends AbstractConfig {
Map.entry(STREAMS_ASSIGNOR_OFFLOAD_ENABLE_CONFIG,
Optional.of(GroupCoordinatorConfig.STREAMS_GROUP_ASSIGNOR_OFFLOAD_ENABLE_CONFIG)),
Map.entry(STREAMS_TASK_OFFSET_INTERVAL_MS_CONFIG,
Optional.of(GroupCoordinatorConfig.STREAMS_GROUP_TASK_OFFSET_INTERVAL_MS_CONFIG)),
Map.entry(STREAMS_NUM_WARMUP_REPLICAS_CONFIG,
Optional.of(GroupCoordinatorConfig.STREAMS_GROUP_NUM_WARMUP_REPLICAS_CONFIG)),
+ Map.entry(STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
Optional.of(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG)),
Map.entry(STREAMS_ACCEPTABLE_RECOVERY_LAG_CONFIG,
Optional.of(GroupCoordinatorConfig.STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_CONFIG)),
// DLQ configs
@@ -395,6 +409,7 @@ public final class GroupConfig extends AbstractConfig {
this.streamsNumStandbyReplicas =
optionalInt(STREAMS_NUM_STANDBY_REPLICAS_CONFIG);
this.streamsInitialRebalanceDelayMs =
optionalInt(STREAMS_INITIAL_REBALANCE_DELAY_MS_CONFIG);
this.streamsAssignmentIntervalMs =
optionalInt(STREAMS_ASSIGNMENT_INTERVAL_MS_CONFIG);
+ this.streamsRackAwareAssignmentTags =
optionalList(STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG);
this.streamsAssignorOffloadEnable =
optionalBoolean(STREAMS_ASSIGNOR_OFFLOAD_ENABLE_CONFIG);
this.streamsTaskOffsetIntervalMs =
optionalInt(STREAMS_TASK_OFFSET_INTERVAL_MS_CONFIG);
this.streamsNumWarmupReplicas =
optionalInt(STREAMS_NUM_WARMUP_REPLICAS_CONFIG);
@@ -422,6 +437,10 @@ public final class GroupConfig extends AbstractConfig {
return originals().containsKey(key) ? Optional.of(getString(key)) :
Optional.empty();
}
+ private Optional<List<String>> optionalList(String key) {
+ return originals().containsKey(key) ? Optional.of(getList(key)) :
Optional.empty();
+ }
+
public static Optional<Type> configType(String configName) {
return
Optional.ofNullable(CONFIG_DEF.configKeys().get(configName)).map(c -> c.type);
}
@@ -435,7 +454,7 @@ public final class GroupConfig extends AbstractConfig {
*
* @param newGroupConfig The new group config overrides.
*/
- public static void validateNames(Map<String, ?> newGroupConfig) {
+ public static void validateNames(Map<String, String> newGroupConfig) {
Set<String> names = configNames();
for (String name : newGroupConfig.keySet()) {
if (!names.contains(name)) {
@@ -453,11 +472,14 @@ public final class GroupConfig extends AbstractConfig {
* @param shareGroupConfig The share group config.
*/
public static void validate(
- Map<String, ?> newGroupConfig,
+ Map<String, String> newGroupConfig,
GroupCoordinatorConfig groupCoordinatorConfig,
ShareGroupConfig shareGroupConfig
) {
validateNames(newGroupConfig);
+ // ConfigDef silently de-duplicates LIST values during parsing, so
inspect the raw value to make
+ // sure an alterConfigs request with duplicate rack-aware assignment
tags is rejected.
+ validateNoDuplicateRackAwareAssignmentTags(newGroupConfig);
var parsed = CONFIG_DEF.parse(newGroupConfig);
parsed.keySet().retainAll(newGroupConfig.keySet());
validateValues(
@@ -467,6 +489,25 @@ public final class GroupConfig extends AbstractConfig {
);
}
+ /**
+ * Rejects an alterConfigs request whose rack-aware assignment tags
contain duplicate keys, inspecting the
+ * raw value because {@link ConfigDef} removes duplicates from LIST values
before they can be validated.
+ *
+ * @param newGroupConfig The new unparsed group config overrides.
+ */
+ private static void validateNoDuplicateRackAwareAssignmentTags(Map<String,
String> newGroupConfig) {
+ String rawValue =
newGroupConfig.get(STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG);
+ if (rawValue == null) {
+ return;
+ }
+ String trimmed = rawValue.trim();
+ List<String> rawTags = trimmed.isEmpty() ? List.of() :
List.of(trimmed.split("\\s*,\\s*", -1));
+ if (Set.copyOf(rawTags).size() != rawTags.size()) {
+ throw new InvalidConfigurationException(
+ STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG + " must not contain
duplicate tag keys.");
+ }
+ }
+
/**
* Validates the parsed values against broker-level bounds.
* Only configs explicitly present in the parsed map are validated.
@@ -1159,6 +1200,13 @@ public final class GroupConfig extends AbstractConfig {
return streamsNumWarmupReplicas;
}
+ /**
+ * The list of client tag keys used for rack-aware standby task assignment.
+ */
+ public Optional<List<String>> streamsRackAwareAssignmentTags() {
+ return streamsRackAwareAssignmentTags;
+ }
+
/**
* The acceptable recovery lag for streams groups.
*/
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
index 31454863565..14646c197a5 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
@@ -386,6 +386,10 @@ public class GroupCoordinatorConfig {
public static final String STREAMS_GROUP_MAX_ASSIGNMENT_INTERVAL_MS_DOC =
"The maximum interval between assignment updates for a streams group.";
public static final int STREAMS_GROUP_MAX_ASSIGNMENT_INTERVAL_MS_DEFAULT =
15000;
+ public static final String STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG
= "group.streams.rack.aware.assignment.tags";
+ public static final String
STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DEFAULT = "";
+ public static final String STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DOC =
"List of client tag keys used to distribute standby replicas across Kafka
Streams instances. When configured, and the used broker-side assignor supports
it, it will make a best-effort to distribute standby tasks over each client tag
dimension.";
+
public static final String STREAMS_GROUP_ASSIGNOR_OFFLOAD_ENABLE_CONFIG =
"group.streams.assignor.offload.enable";
public static final String STREAMS_GROUP_ASSIGNOR_OFFLOAD_ENABLE_DOC =
"Whether to offload streams group assignment to a group coordinator background
thread.";
public static final boolean STREAMS_GROUP_ASSIGNOR_OFFLOAD_ENABLE_DEFAULT
= true;
@@ -507,6 +511,7 @@ public class GroupCoordinatorConfig {
.define(STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_CONFIG, INT,
STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_DEFAULT, atLeast(1), MEDIUM,
STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_DOC)
.define(STREAMS_GROUP_NUM_WARMUP_REPLICAS_CONFIG, INT,
STREAMS_GROUP_NUM_WARMUP_REPLICAS_DEFAULT, atLeast(0), MEDIUM,
STREAMS_GROUP_NUM_WARMUP_REPLICAS_DOC)
.define(STREAMS_GROUP_MAX_WARMUP_REPLICAS_CONFIG, INT,
STREAMS_GROUP_MAX_WARMUP_REPLICAS_DEFAULT, atLeast(0), MEDIUM,
STREAMS_GROUP_MAX_WARMUP_REPLICAS_DOC)
+ .define(STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG, LIST,
STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DEFAULT,
ConfigDef.ValidList.anyNonDuplicateValues(true, false), LOW,
STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_DOC)
.define(STREAMS_GROUP_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS_CONFIG, CLASS,
STREAMS_GROUP_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS_DEFAULT, MEDIUM,
STREAMS_GROUP_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS_DOC)
.define(STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_CONFIG, LONG,
STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_DEFAULT, atLeast(0L), MEDIUM,
STREAMS_GROUP_ACCEPTABLE_RECOVERY_LAG_DOC);
@@ -576,6 +581,7 @@ public class GroupCoordinatorConfig {
private final int streamsGroupMinTaskOffsetIntervalMs;
private final int streamsGroupNumWarmupReplicas;
private final int streamsGroupMaxWarmupReplicas;
+ private final List<String> streamsGroupRackAwareAssignmentTags;
private final long streamsGroupAcceptableRecoveryLag;
private final AbstractConfig config;
@@ -639,6 +645,7 @@ public class GroupCoordinatorConfig {
this.streamsGroupMaxSize =
config.getInt(GroupCoordinatorConfig.STREAMS_GROUP_MAX_SIZE_CONFIG);
this.streamsGroupNumStandbyReplicas =
config.getInt(GroupCoordinatorConfig.STREAMS_GROUP_NUM_STANDBY_REPLICAS_CONFIG);
this.streamsGroupMaxStandbyReplicas =
config.getInt(GroupCoordinatorConfig.STREAMS_GROUP_MAX_STANDBY_REPLICAS_CONFIG);
+ this.streamsGroupRackAwareAssignmentTags =
config.getList(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG);
this.streamsGroupInitialRebalanceDelayMs =
config.getInt(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG);
this.streamsGroupMinAssignmentIntervalMs =
config.getInt(GroupCoordinatorConfig.STREAMS_GROUP_MIN_ASSIGNMENT_INTERVAL_MS_CONFIG);
this.streamsGroupMaxAssignmentIntervalMs =
config.getInt(GroupCoordinatorConfig.STREAMS_GROUP_MAX_ASSIGNMENT_INTERVAL_MS_CONFIG);
@@ -653,7 +660,12 @@ public class GroupCoordinatorConfig {
}
private void checkConstraints() {
- // New group coordinator configs validation.
+ verifyConsumerGroupConfigs();
+ verifyShareGroupConfigs();
+ verifyStreamsGroupConfigs();
+ }
+
+ private void verifyConsumerGroupConfigs() {
require(consumerGroupMaxHeartbeatIntervalMs >=
consumerGroupMinHeartbeatIntervalMs,
String.format("%s must be greater than or equal to %s",
CONSUMER_GROUP_MAX_HEARTBEAT_INTERVAL_MS_CONFIG,
CONSUMER_GROUP_MIN_HEARTBEAT_INTERVAL_MS_CONFIG));
require(consumerGroupHeartbeatIntervalMs >=
consumerGroupMinHeartbeatIntervalMs,
@@ -676,9 +688,9 @@ public class GroupCoordinatorConfig {
String.format("%s must be greater than or equal to %s",
CONSUMER_GROUP_ASSIGNMENT_INTERVAL_MS_CONFIG,
CONSUMER_GROUP_MIN_ASSIGNMENT_INTERVAL_MS_CONFIG));
require(consumerGroupAssignmentIntervalMs() <=
consumerGroupMaxAssignmentIntervalMs,
String.format("%s must be less than or equal to %s",
CONSUMER_GROUP_ASSIGNMENT_INTERVAL_MS_CONFIG,
CONSUMER_GROUP_MAX_ASSIGNMENT_INTERVAL_MS_CONFIG));
+ }
-
- // Share group configs validation.
+ private void verifyShareGroupConfigs() {
require(shareGroupMaxHeartbeatIntervalMs >=
shareGroupMinHeartbeatIntervalMs,
String.format("%s must be greater than or equal to %s",
SHARE_GROUP_MAX_HEARTBEAT_INTERVAL_MS_CONFIG,
SHARE_GROUP_MIN_HEARTBEAT_INTERVAL_MS_CONFIG));
@@ -714,9 +726,9 @@ public class GroupCoordinatorConfig {
require(shareGroupAssignmentIntervalMs() <=
shareGroupMaxAssignmentIntervalMs,
String.format("%s must be less than or equal to %s",
SHARE_GROUP_ASSIGNMENT_INTERVAL_MS_CONFIG,
SHARE_GROUP_MAX_ASSIGNMENT_INTERVAL_MS_CONFIG));
+ }
-
- // Streams group configs validation.
+ private void verifyStreamsGroupConfigs() {
require(streamsGroupMaxHeartbeatIntervalMs >=
streamsGroupMinHeartbeatIntervalMs,
String.format("%s must be greater than or equal to %s",
STREAMS_GROUP_MAX_HEARTBEAT_INTERVAL_MS_CONFIG,
STREAMS_GROUP_MIN_HEARTBEAT_INTERVAL_MS_CONFIG));
@@ -754,6 +766,24 @@ public class GroupCoordinatorConfig {
String.format("%s must be greater than or equal to %s",
STREAMS_GROUP_TASK_OFFSET_INTERVAL_MS_CONFIG,
STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_CONFIG));
require(streamsGroupNumWarmupReplicas <= streamsGroupMaxWarmupReplicas,
String.format("%s must be less than or equal to %s",
STREAMS_GROUP_NUM_WARMUP_REPLICAS_CONFIG,
STREAMS_GROUP_MAX_WARMUP_REPLICAS_CONFIG));
+
+ // ConfigDef silently de-duplicates LIST values during parsing, so
inspect the raw value to make
+ // sure a broker configured with duplicate rack-aware assignment tags
refuses to start.
+ List<String> rawRackAwareAssignmentTags = rawRackAwareAssignmentTags();
+ require(Set.copyOf(rawRackAwareAssignmentTags).size() ==
rawRackAwareAssignmentTags.size(),
+ String.format("%s must not contain duplicate tag keys.",
STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG));
+ }
+
+ /**
+ * Returns the rack-aware assignment tags exactly as configured, before
{@link ConfigDef} removes duplicates.
+ */
+ private List<String> rawRackAwareAssignmentTags() {
+ String rawValue = (String)
config.originals().get(STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG);
+ if (rawValue == null) {
+ return List.of();
+ }
+ String trimmed = rawValue.trim();
+ return trimmed.isEmpty() ? List.of() :
List.of(trimmed.split("\\s*,\\s*", -1));
}
/**
@@ -1341,6 +1371,13 @@ public class GroupCoordinatorConfig {
return streamsGroupMaxStandbyReplicas;
}
+ /**
+ * The list of client tag keys used for rack-aware standby task assignment.
+ */
+ public List<String> streamsGroupRackAwareAssignmentTags() {
+ return streamsGroupRackAwareAssignmentTags;
+ }
+
/**
* The initial rebalance delay for streams groups.
*/
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
index 43ef8e90331..eaba792fdac 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
@@ -185,6 +185,7 @@ import org.slf4j.Logger;
import java.nio.ByteBuffer;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
@@ -2049,6 +2050,7 @@ public class GroupMetadataManager {
* @param clientTags Used for rack-aware assignment algorithm, or
null.
* @param shutdownApplication Whether all Streams clients in the group
should shut down.
* @param memberEndpointEpoch The last endpoint information epoch seen be
the group member.
+ * @param requestApiVersion The api version of the StreamsGroupHeartbeat
request.
* @return A result containing the StreamsGroupHeartbeat response and a
list of records to update the state machine.
*/
private CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
streamsGroupHeartbeat(
@@ -2070,7 +2072,8 @@ public class GroupMetadataManager {
Endpoint userEndpoint,
List<KeyValue> clientTags,
boolean shutdownApplication,
- int memberEndpointEpoch
+ int memberEndpointEpoch,
+ int requestApiVersion
) throws ApiException {
final long currentTimeMs = time.milliseconds();
final List<CoordinatorRecord> records = new ArrayList<>();
@@ -2355,6 +2358,27 @@ public class GroupMetadataManager {
)
));
+ String rackAwareTagsValue =
currentAssignmentConfigs.getOrDefault("rack.aware.assignment.tags", "").trim();
+ // The MISSING_CLIENT_TAGS status (code 6) requires version 1 of the
RPC: version 0 clients
+ // throw on unknown status codes, so it must not be sent to them.
+ if (requestApiVersion >= 1 && !rackAwareTagsValue.isEmpty()) {
+ List<String> requiredTags =
Arrays.asList(rackAwareTagsValue.split("\\s*,\\s*", -1));
+ Set<String> memberTagKeys = updatedMember.clientTags().keySet();
+ List<String> missingTags = requiredTags.stream()
+ .filter(tag -> !memberTagKeys.contains(tag))
+ .collect(Collectors.toList());
+ if (!missingTags.isEmpty()) {
+ returnedStatus.add(
+ new Status()
+
.setStatusCode(StreamsGroupHeartbeatResponse.Status.MISSING_CLIENT_TAGS.code())
+ .setStatusDetail(
+ String.format("Missing required client tags for
rack-aware standby assignment: %s. " +
+ "Configure them via 'client.tag.<tagKey>' in
your Streams config.", missingTags)
+ )
+ );
+ }
+ }
+
response.setStatus(returnedStatus);
return new CoordinatorResult<>(records, new
StreamsGroupHeartbeatResult(
@@ -5483,7 +5507,8 @@ public class GroupMetadataManager {
request.userEndpoint(),
request.clientTags(),
request.shutdownApplication(),
- request.endpointInformationEpoch()
+ request.endpointInformationEpoch(),
+ context.requestVersion()
);
}
}
@@ -9738,9 +9763,14 @@ public class GroupMetadataManager {
Optional<GroupConfig> groupConfig =
groupConfigManager.groupConfig(groupId);
final Integer numStandbyReplicas =
groupConfig.flatMap(GroupConfig::streamsNumStandbyReplicas)
.orElse(config.streamsGroupNumStandbyReplicas());
- return new TreeMap<>(Map.of(
- "num.standby.replicas", numStandbyReplicas.toString()
- ));
+ final List<String> rackAwareAssignmentTags =
groupConfig.flatMap(GroupConfig::streamsRackAwareAssignmentTags)
+ .orElse(config.streamsGroupRackAwareAssignmentTags());
+ Map<String, String> configs = new TreeMap<>();
+ configs.put("num.standby.replicas", numStandbyReplicas.toString());
+ if (!rackAwareAssignmentTags.isEmpty()) {
+ configs.put("rack.aware.assignment.tags", String.join(",",
rackAwareAssignmentTags));
+ }
+ return configs;
}
private static boolean hasUserEndpointChanged(StreamsGroupMember
maybeOldMember, StreamsGroupMember updatedMember) {
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupConfigTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupConfigTest.java
index 7c3549cc1ab..5436c7c4466 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupConfigTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupConfigTest.java
@@ -30,6 +30,7 @@ import org.junit.jupiter.params.provider.MethodSource;
import java.util.HashMap;
import java.util.HashSet;
+import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Properties;
@@ -132,6 +133,10 @@ public class GroupConfigTest {
assertPropertyInvalid(name, "not_a_number", "1.0");
} else if
(GroupConfig.ERRORS_DEADLETTERQUEUE_COPY_RECORD_ENABLE_CONFIG.equals(name)) {
assertPropertyInvalid(name, "not_a_boolean");
+ } else if
(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG.equals(name)) {
+ // This is a free-form list of tag keys, so values like
"not_a_number" are valid. Only an
+ // empty tag key (an empty element between commas) is rejected.
+ assertPropertyInvalid(name, "tag1,,tag2");
} else if
(!GroupConfig.ERRORS_DEADLETTERQUEUE_TOPIC_NAME_CONFIG.equals(name)) {
assertPropertyInvalid(name, "not_a_number", "-0.1");
}
@@ -345,6 +350,35 @@ public class GroupConfigTest {
doTestInvalidProps(props, ConfigException.class);
}
+ @Test
+ public void testStreamsRackAwareAssignmentTagsValidation() {
+ // Duplicate rack-aware assignment tags are rejected rather than
silently de-duplicated.
+ Map<String, String> duplicateProps = createValidGroupConfig();
+
duplicateProps.put(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,zone");
+ assertEquals("streams.rack.aware.assignment.tags must not contain
duplicate tag keys.",
+ assertThrows(InvalidConfigurationException.class,
+ () -> GroupConfig.validate(duplicateProps,
createGroupCoordinatorConfig(), createShareGroupConfig())).getMessage());
+
+ // Duplicates are detected regardless of surrounding whitespace.
+ Map<String, String> whitespaceDuplicateProps =
createValidGroupConfig();
+
whitespaceDuplicateProps.put(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
" zone , zone ");
+ assertEquals("streams.rack.aware.assignment.tags must not contain
duplicate tag keys.",
+ assertThrows(InvalidConfigurationException.class,
+ () -> GroupConfig.validate(whitespaceDuplicateProps,
createGroupCoordinatorConfig(), createShareGroupConfig())).getMessage());
+
+ // Distinct rack-aware assignment tags are accepted.
+ Map<String, String> distinctProps = createValidGroupConfig();
+
distinctProps.put(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster");
+ doTestValidProps(distinctProps);
+
+ // Surrounding whitespace is trimmed: " zone , cluster " parses and
returns two clean tags.
+ Map<String, String> whitespaceProps = createValidGroupConfig();
+
whitespaceProps.put(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG, "
zone , cluster ");
+ doTestValidProps(whitespaceProps);
+ assertEquals(Optional.of(List.of("zone", "cluster")),
+ new GroupConfig(whitespaceProps).streamsRackAwareAssignmentTags());
+ }
+
private void doTestInvalidProps(Map<String, String> props, Class<? extends
Exception> exceptionClassName) {
assertThrows(exceptionClassName, () -> GroupConfig.validate(props,
createGroupCoordinatorConfig(), createShareGroupConfig()));
}
@@ -509,6 +543,7 @@ public class GroupConfigTest {
assertEquals(Optional.empty(), config.streamsAssignmentIntervalMs());
assertEquals(Optional.empty(), config.streamsAssignorOffloadEnable());
assertEquals(Optional.empty(), config.streamsTaskOffsetIntervalMs());
+ assertEquals(Optional.empty(),
config.streamsRackAwareAssignmentTags());
// DLQ configs - have defaults from CONFIG_DEF
assertEquals("", config.errorsDLQTopicName());
@@ -540,6 +575,7 @@ public class GroupConfigTest {
props.put(GroupConfig.STREAMS_ASSIGNMENT_INTERVAL_MS_CONFIG, "1250");
props.put(GroupConfig.STREAMS_ASSIGNOR_OFFLOAD_ENABLE_CONFIG, "false");
props.put(GroupConfig.STREAMS_TASK_OFFSET_INTERVAL_MS_CONFIG, "30000");
+ props.put(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster");
props.put(GroupConfig.ERRORS_DEADLETTERQUEUE_TOPIC_NAME_CONFIG,
"my-dlq-topic");
props.put(GroupConfig.ERRORS_DEADLETTERQUEUE_COPY_RECORD_ENABLE_CONFIG, "true");
@@ -571,6 +607,7 @@ public class GroupConfigTest {
assertEquals(Optional.of(1250), config.streamsAssignmentIntervalMs());
assertEquals(Optional.of(false),
config.streamsAssignorOffloadEnable());
assertEquals(Optional.of(30000), config.streamsTaskOffsetIntervalMs());
+ assertEquals(Optional.of(List.of("zone", "cluster")),
config.streamsRackAwareAssignmentTags());
// DLQ configs
assertEquals("my-dlq-topic", config.errorsDLQTopicName());
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfigTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfigTest.java
index 999e9b78684..297019d70f1 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfigTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfigTest.java
@@ -603,6 +603,44 @@ public class GroupCoordinatorConfigTest {
configs.put(GroupCoordinatorConfig.STREAMS_GROUP_MAX_WARMUP_REPLICAS_CONFIG,
-1);
assertEquals("Invalid value -1 for configuration
group.streams.max.warmup.replicas: Value must be at least 0",
assertThrows(ConfigException.class, () ->
createConfig(configs)).getMessage());
+
+
+ // group.streams.rack.aware.assignment.tags
+
+ // default is empty list
+ configs.clear();
+ GroupCoordinatorConfig defaultTagsConfig = createConfig(configs);
+ assertEquals(List.of(),
defaultTagsConfig.streamsGroupRackAwareAssignmentTags());
+
+ // can parse a non-empty value
+ configs.clear();
+
configs.put(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster");
+ GroupCoordinatorConfig nonEmptyTagsConfig = createConfig(configs);
+ assertEquals(List.of("zone", "cluster"),
nonEmptyTagsConfig.streamsGroupRackAwareAssignmentTags());
+
+ // surrounding whitespace is trimmed: " zone , cluster " parses and
returns two clean tags
+ configs.clear();
+
configs.put(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
" zone , cluster ");
+ GroupCoordinatorConfig whitespaceTagsConfig = createConfig(configs);
+ assertEquals(List.of("zone", "cluster"),
whitespaceTagsConfig.streamsGroupRackAwareAssignmentTags());
+
+ // rejects empty entries in the list
+ configs.clear();
+
configs.put(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone, ");
+ assertEquals("Configuration 'group.streams.rack.aware.assignment.tags'
values must not be empty.",
+ assertThrows(ConfigException.class, () ->
createConfig(configs)).getMessage());
+
+ // duplicate tag keys make the broker refuse to start
+ configs.clear();
+
configs.put(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,zone");
+ assertEquals("group.streams.rack.aware.assignment.tags must not
contain duplicate tag keys.",
+ assertThrows(IllegalArgumentException.class, () ->
createConfig(configs)).getMessage());
+
+ // duplicate tag keys are detected regardless of surrounding whitespace
+ configs.clear();
+
configs.put(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
" zone , zone ");
+ assertEquals("group.streams.rack.aware.assignment.tags must not
contain duplicate tag keys.",
+ assertThrows(IllegalArgumentException.class, () ->
createConfig(configs)).getMessage());
}
@Test
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
index 24fcb7a308f..32c60e5384c 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
@@ -21,6 +21,7 @@ import
org.apache.kafka.clients.consumer.internals.ConsumerProtocol;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.config.AbstractConfig;
+import org.apache.kafka.common.config.ConfigException;
import org.apache.kafka.common.errors.CoordinatorNotAvailableException;
import org.apache.kafka.common.errors.FencedInstanceIdException;
import org.apache.kafka.common.errors.FencedMemberEpochException;
@@ -1181,7 +1182,7 @@ public class GroupMetadataManagerTest {
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupTopologyRecord(groupId,
topology));
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord(groupId,
100, computeGroupHash(Map.of(
fooTopicName, fooTopicHash
- )), 1, new TreeMap<>(Map.of("num.standby.replicas", "0")), -1, -1));
+ )), 1, getDefaultAssignmentConfigs(), -1, -1));
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentRecord(groupId,
memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)
@@ -19241,9 +19242,7 @@ public class GroupMetadataManagerTest {
2,
groupMetadataHash,
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -19259,6 +19258,305 @@ public class GroupMetadataManagerTest {
assertRecordsEquals(expectedRecords, result.records());
}
+ @Test
+ public void testMissingClientTagsStatusWhenRackAwareTagsConfigured() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(metadataImage)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
0)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster")
+ .build();
+
+ assignor.prepareGroupAssignment(Map.of(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2, 3, 4, 5)
+ )));
+
+ // Member joins without any client tags — should get
MISSING_CLIENT_TAGS status
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setProcessId("process-id")
+ .setRebalanceTimeoutMs(1500)
+ .setTopology(topology)
+ .setActiveTasks(List.of())
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of()));
+
+ assertTrue(result.response().data().status().stream()
+ .anyMatch(s -> s.statusCode() == Status.MISSING_CLIENT_TAGS.code()
+ && s.statusDetail().contains("zone")
+ && s.statusDetail().contains("cluster")));
+ }
+
+ @Test
+ public void testMissingClientTagsStatusNotReturnedToVersion0Clients() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(metadataImage)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
0)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster")
+ .build();
+
+ assignor.prepareGroupAssignment(Map.of(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2, 3, 4, 5)
+ )));
+
+ // Member joins without any client tags using version 0 of the RPC.
Version 0 clients throw on
+ // the unknown MISSING_CLIENT_TAGS status code, so the broker must not
return it to them.
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setProcessId("process-id")
+ .setRebalanceTimeoutMs(1500)
+ .setTopology(topology)
+ .setActiveTasks(List.of())
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of()),
+ (short) 0);
+
+ assertTrue(result.response().data().status().stream()
+ .noneMatch(s -> s.statusCode() ==
Status.MISSING_CLIENT_TAGS.code()));
+ }
+
+ @Test
+ public void testMissingClientTagsStatusReportsOnlyMissingTags() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(metadataImage)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
0)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,cluster")
+ .build();
+
+ assignor.prepareGroupAssignment(Map.of(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2, 3, 4, 5)
+ )));
+
+ // Member joins with only the "cluster" tag — only the missing "zone"
tag should be reported.
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setProcessId("process-id")
+ .setRebalanceTimeoutMs(1500)
+ .setTopology(topology)
+ .setClientTags(List.of(
+ new
StreamsGroupHeartbeatRequestData.KeyValue().setKey("cluster").setValue("c1")))
+ .setActiveTasks(List.of())
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of()));
+
+ StreamsGroupHeartbeatResponseData.Status status =
result.response().data().status().stream()
+ .filter(s -> s.statusCode() == Status.MISSING_CLIENT_TAGS.code())
+ .findFirst()
+ .orElseThrow(() -> new AssertionError("Expected a
MISSING_CLIENT_TAGS status"));
+ assertTrue(status.statusDetail().contains("[zone]"));
+ assertFalse(status.statusDetail().contains("cluster"));
+ }
+
+ @Test
+ public void testNoMissingClientTagsStatusWhenAllTagsPresent() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(metadataImage)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
0)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone")
+ .build();
+
+ assignor.prepareGroupAssignment(Map.of(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2, 3, 4, 5)
+ )));
+
+ // Member joins with the required client tag — should NOT get
MISSING_CLIENT_TAGS status
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setProcessId("process-id")
+ .setRebalanceTimeoutMs(1500)
+ .setTopology(topology)
+ .setClientTags(List.of(
+ new
StreamsGroupHeartbeatRequestData.KeyValue().setKey("zone").setValue("us-east-1a")))
+ .setActiveTasks(List.of())
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of()));
+
+ assertTrue(result.response().data().status().stream()
+ .noneMatch(s -> s.statusCode() ==
Status.MISSING_CLIENT_TAGS.code()));
+ }
+
+ @Test
+ public void testRackAwareAssignmentTagsRejectsEmptyTags() {
+ assertThrows(ConfigException.class, () -> new
GroupMetadataManagerTestContext.Builder()
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone,,cluster")
+ .build());
+ }
+
+ @Test
+ public void testNoMissingClientTagsStatusWhenRackAwareTagsNotConfigured() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+ // No rack-aware assignment tags configured (broker default is empty).
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(metadataImage)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
0)
+ .build();
+
+ assignor.prepareGroupAssignment(Map.of(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2, 3, 4, 5)
+ )));
+
+ // Member joins without any client tags. Since the feature is not
configured, no MISSING_CLIENT_TAGS
+ // status should be returned.
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setProcessId("process-id")
+ .setRebalanceTimeoutMs(1500)
+ .setTopology(topology)
+ .setActiveTasks(List.of())
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of()));
+
+ assertTrue(result.response().data().status().stream()
+ .noneMatch(s -> s.statusCode() ==
Status.MISSING_CLIENT_TAGS.code()));
+ }
+
+ @Test
+ public void testGroupConfigOverridesBrokerRackAwareAssignmentTags() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+ // Broker default requires only the "zone" tag.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(metadataImage)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
0)
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"zone")
+ .build();
+
+ // The per-group override requires "rack" instead, taking precedence
over the broker default.
+ Properties newConfig = new Properties();
+ newConfig.put(GroupConfig.STREAMS_RACK_AWARE_ASSIGNMENT_TAGS_CONFIG,
"rack");
+ context.groupConfigManager.updateGroupConfig(groupId, newConfig);
+
+ assignor.prepareGroupAssignment(Map.of(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2, 3, 4, 5)
+ )));
+
+ // Member reports the broker-default tag "zone" but not the
group-override tag "rack". Because the
+ // group override takes precedence, "rack" must be reported as missing
(and "zone" must not be).
+ CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord>
result = context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(0)
+ .setProcessId("process-id")
+ .setRebalanceTimeoutMs(1500)
+ .setTopology(topology)
+ .setClientTags(List.of(
+ new
StreamsGroupHeartbeatRequestData.KeyValue().setKey("zone").setValue("us-east-1a")))
+ .setActiveTasks(List.of())
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of()));
+
+ StreamsGroupHeartbeatResponseData.Status status =
result.response().data().status().stream()
+ .filter(s -> s.statusCode() == Status.MISSING_CLIENT_TAGS.code())
+ .findFirst()
+ .orElseThrow(() -> new AssertionError("Expected a
MISSING_CLIENT_TAGS status"));
+ assertTrue(status.statusDetail().contains("[rack]"),
+ "The group-level override tag 'rack' should be reported as
missing");
+ assertFalse(status.statusDetail().contains("zone"),
+ "The broker-default tag 'zone' should not be used once the group
override is set");
+ }
+
@Test
public void testJoinEmptyStreamsGroupAndDescribe() {
String groupId = "fooup";
@@ -19430,9 +19728,7 @@ public class GroupMetadataManagerTest {
2,
computeGroupHash(Map.of(fooTopicName,
computeTopicHash(fooTopicName, metadataImage))),
-1,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -19526,9 +19822,7 @@ public class GroupMetadataManagerTest {
2,
computeGroupHash(Map.of(fooTopicName,
computeTopicHash(fooTopicName, metadataImage))),
-1,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -19622,9 +19916,7 @@ public class GroupMetadataManagerTest {
barTopicName, computeTopicHash(barTopicName, metadataImage)
)),
-1,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -19731,9 +20023,7 @@ public class GroupMetadataManagerTest {
barTopicName, computeTopicHash(barTopicName, metadataImage)
)),
1,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -20048,9 +20338,7 @@ public class GroupMetadataManagerTest {
11,
groupMetadataHash,
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -20183,9 +20471,7 @@ public class GroupMetadataManagerTest {
barTopicName, computeTopicHash(barTopicName,
newMetadataImage)
)),
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -21485,9 +21771,7 @@ public class GroupMetadataManagerTest {
11,
computeGroupHash(Map.of(fooTopicName,
computeTopicHash(fooTopicName, metadataImage))),
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -21617,9 +21901,7 @@ public class GroupMetadataManagerTest {
11,
computeGroupHash(Map.of(fooTopicName,
computeTopicHash(fooTopicName, metadataImage))),
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -21778,9 +22060,7 @@ public class GroupMetadataManagerTest {
3,
groupMetadataHash,
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -22113,9 +22393,7 @@ public class GroupMetadataManagerTest {
4,
groupMetadataHash,
0,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
)
@@ -22693,9 +22971,7 @@ public class GroupMetadataManagerTest {
2,
0,
-1,
- new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )),
+ getDefaultAssignmentConfigs(),
-1,
-1
),
@@ -23131,9 +23407,7 @@ public class GroupMetadataManagerTest {
.setWarmupTasks(List.of()));
assertEquals(2, result.response().data().memberEpoch());
assertEquals(
- Map.of(
- "num.standby.replicas", "0"
- ),
+ getDefaultAssignmentConfigs(),
assignor.lastPassedAssignmentConfigs()
);
@@ -23172,9 +23446,7 @@ public class GroupMetadataManagerTest {
// Verify that the new number of standby replicas is used
assertEquals(
- Map.of(
- "num.standby.replicas", "2"
- ),
+ Map.of("num.standby.replicas", "2"),
assignor.lastPassedAssignmentConfigs()
);
@@ -23522,9 +23794,7 @@ public class GroupMetadataManagerTest {
context.assertSessionTimeout(groupId, memberId,
GroupCoordinatorConfig.STREAMS_GROUP_SESSION_TIMEOUT_MS_DEFAULT);
assertEquals(
- Map.of(
- "num.standby.replicas",
String.valueOf(GroupCoordinatorConfig.STREAMS_GROUP_NUM_STANDBY_REPLICAS_DEFAULT)
- ),
+ getDefaultAssignmentConfigs(),
assignor.lastPassedAssignmentConfigs());
assertEquals(GroupCoordinatorConfig.STREAMS_GROUP_TASK_OFFSET_INTERVAL_MS_DEFAULT,
result.response().data().taskOffsetIntervalMs());
@@ -23564,9 +23834,7 @@ public class GroupMetadataManagerTest {
// Verify that the number of standby replicas is evaluated to max,
// and task offset interval is evaluated to min
assertEquals(
- Map.of(
- "num.standby.replicas",
String.valueOf(GroupCoordinatorConfig.STREAMS_GROUP_MAX_STANDBY_REPLICAS_DEFAULT)
- ),
+ Map.of("num.standby.replicas",
String.valueOf(GroupCoordinatorConfig.STREAMS_GROUP_MAX_STANDBY_REPLICAS_DEFAULT)),
assignor.lastPassedAssignmentConfigs());
assertEquals(GroupCoordinatorConfig.STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_DEFAULT,
result.response().data().taskOffsetIntervalMs());
@@ -23863,7 +24131,7 @@ public class GroupMetadataManagerTest {
// The group still exists but the member is already gone. Replaying the
// StreamsGroupMemberMetadata tombstone should be a no-op.
-
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, Map.of("num.standby.replicas", "0"), -1, -1));
+
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, getDefaultAssignmentConfigs(), -1, -1));
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMemberTombstoneRecord("foo",
"m1"));
assertThrows(UnknownMemberIdException.class, () ->
context.groupMetadataManager.streamsGroup("foo").getMemberOrThrow("m1"));
@@ -23919,7 +24187,7 @@ public class GroupMetadataManagerTest {
.build();
// The group is created if it does not exist.
-
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, Map.of("num.standby.replicas", "0"), -1, -1));
+
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, getDefaultAssignmentConfigs(), -1, -1));
assertEquals(10,
context.groupMetadataManager.streamsGroup("foo").groupEpoch());
}
@@ -24116,7 +24384,7 @@ public class GroupMetadataManagerTest {
// The group still exists, but the member is already gone. Replaying
the
// StreamsGroupCurrentMemberAssignment tombstone should be a no-op.
-
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, Map.of("num.standby.replicas", "0"), -1, -1));
+
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, getDefaultAssignmentConfigs(), -1, -1));
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentTombstoneRecord("foo",
"m1"));
assertThrows(UnknownMemberIdException.class, () ->
context.groupMetadataManager.streamsGroup("foo").getMemberOrThrow("m1"));
@@ -24192,7 +24460,7 @@ public class GroupMetadataManagerTest {
// The group still exists, but the member is already gone. Replaying
the
// StreamsGroupTopology tombstone should be a no-op.
-
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, Map.of("num.standby.replicas", "0"), -1, -1));
+
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord("foo",
10, 0, 0, getDefaultAssignmentConfigs(), -1, -1));
context.replay(StreamsCoordinatorRecordHelpers.newStreamsGroupTopologyRecordTombstone("foo"));
assertTrue(context.groupMetadataManager.streamsGroup("foo").topology().isEmpty());
@@ -29594,9 +29862,7 @@ public class GroupMetadataManagerTest {
StreamsCoordinatorRecordHelpers.newStreamsGroupTopologyRecord(groupId,
topology),
StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord(groupId, 2,
computeGroupHash(Map.of(
fooTopicName, computeTopicHash(fooTopicName, metadataImage)
- )), 0, new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )), -1, -1),
+ )), 0, getDefaultAssignmentConfigs(), -1, -1),
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentRecord(groupId,
memberId1,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
TaskAssignmentTestUtil.mkTasks(subtopology, 0, 1, 2,
3, 4, 5)
@@ -29659,9 +29925,7 @@ public class GroupMetadataManagerTest {
StreamsCoordinatorRecordHelpers.newStreamsGroupMemberRecord(groupId,
expectedMember2),
StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord(groupId, 3,
computeGroupHash(Map.of(
fooTopicName, computeTopicHash(fooTopicName, metadataImage)
- )), 0, new TreeMap<>(Map.of(
- "num.standby.replicas", "0"
- )), -1, -1),
+ )), 0, getDefaultAssignmentConfigs(), -1, -1),
StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentRecord(groupId,
expectedMember2)
),
result2.records()