This is an automated email from the ASF dual-hosted git repository.
squah-confluent 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 0d0ec5af1e0 KAFKA-20843: Custom consumer group assignors may silently
use duplicate or built-in assignor names (#22982)
0d0ec5af1e0 is described below
commit 0d0ec5af1e08d4a4245e6feaaa48bd014598cde3
Author: gabriellafu <[email protected]>
AuthorDate: Wed Jul 29 17:55:31 2026 -0400
KAFKA-20843: Custom consumer group assignors may silently use duplicate or
built-in assignor names (#22982)
Fixed an issue where group.consumer.assignors did not check for assignor
name collisions. Previously, configuring duplicate assignors (such as
using both a short name and a fully qualified class name) or using a
custom assignor with a built-in name passed startup validation, but
caused runtime crashes later during group coordinator loading.
Key Changes:
Added broker-side duplicate name validation that throws ConfigException
at startup if two entries resolve to the same name.
Prevented custom assignors from reusing reserved built-in names (e.g.,
uniform, range).
Preserved support for configuring built-in assignors via their fully
qualified class names.
Added corresponding unit tests in GroupCoordinatorConfigTest and
verified existing integration tests.
---
.../coordinator/group/GroupCoordinatorConfig.java | 24 +++++-
.../group/GroupCoordinatorConfigTest.java | 85 ++++++++++++++++++++++
2 files changed, 107 insertions(+), 2 deletions(-)
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 0501a7c9e9c..012ca319878 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
@@ -36,6 +36,7 @@ import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -853,17 +854,24 @@ public class GroupCoordinatorConfig {
protected List<ConsumerGroupPartitionAssignor> consumerGroupAssignors(
AbstractConfig config
) {
- Map<String, ConsumerGroupPartitionAssignor> defaultAssignors =
CONSUMER_GROUP_BUILTIN_ASSIGNORS
+ Map<String, ConsumerGroupPartitionAssignor> builtInAssignors =
CONSUMER_GROUP_BUILTIN_ASSIGNORS
.stream()
.collect(Collectors.toMap(ConsumerGroupPartitionAssignor::name,
Function.identity()));
+ // A built-in may be configured either by its name or by its class
name, so it is recognised
+ // by class rather than by how it was resolved below.
+ Set<Class<? extends ConsumerGroupPartitionAssignor>>
builtInAssignorClasses = CONSUMER_GROUP_BUILTIN_ASSIGNORS
+ .stream()
+ .map(ConsumerGroupPartitionAssignor::getClass)
+ .collect(Collectors.toSet());
List<ConsumerGroupPartitionAssignor> assignors = new ArrayList<>();
+ Set<String> assignorNames = new HashSet<>();
try {
// `configuredAssignor` is either the name of a built-in assignor,
// or a fully qualified class name of a custom assignor
for (String configuredAssignor :
config.getList(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG)) {
- ConsumerGroupPartitionAssignor assignor =
defaultAssignors.get(configuredAssignor);
+ ConsumerGroupPartitionAssignor assignor =
builtInAssignors.get(configuredAssignor);
if (assignor == null) {
try {
assignor = Utils.newInstance(configuredAssignor,
ConsumerGroupPartitionAssignor.class);
@@ -882,6 +890,18 @@ public class GroupCoordinatorConfig {
assignors.add(assignor);
+ if (!builtInAssignorClasses.contains(assignor.getClass()) &&
builtInAssignors.containsKey(assignor.name())) {
+ throw new ConfigException(CONSUMER_GROUP_ASSIGNORS_CONFIG,
configuredAssignor,
+ "Assignor name '" + assignor.name() + "' is reserved
by a built-in assignor. " +
+ "A custom assignor must not reuse the name of a
built-in assignor");
+ }
+
+ if (!assignorNames.add(assignor.name())) {
+ throw new ConfigException(CONSUMER_GROUP_ASSIGNORS_CONFIG,
configuredAssignor,
+ "Assignor name '" + assignor.name() + "' is already
registered by another configured assignor. " +
+ "Assignor names, whether built-in or custom, must
be unique");
+ }
+
if (assignor instanceof Configurable configurable) {
configurable.configure(config.originals());
}
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 c53c297d22e..8443624f83c 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
@@ -72,6 +72,38 @@ public class GroupCoordinatorConfigTest {
}
}
+ public static class DuplicateNameAssignor implements
ConsumerGroupPartitionAssignor {
+ @Override
+ public String name() {
+ // Collides with CustomAssignor.
+ return "CustomAssignor";
+ }
+
+ @Override
+ public GroupAssignment assign(
+ GroupSpec groupSpec,
+ SubscribedTopicDescriber subscribedTopicDescriber
+ ) throws PartitionAssignorException {
+ return null;
+ }
+ }
+
+ public static class UniformNamedAssignor implements
ConsumerGroupPartitionAssignor {
+ @Override
+ public String name() {
+ // Collides with the built-in "uniform" assignor.
+ return "uniform";
+ }
+
+ @Override
+ public GroupAssignment assign(
+ GroupSpec groupSpec,
+ SubscribedTopicDescriber subscribedTopicDescriber
+ ) throws PartitionAssignorException {
+ return null;
+ }
+ }
+
public static class NoDefaultConstructorAssignor implements
ConsumerGroupPartitionAssignor, ShareGroupPartitionAssignor {
public NoDefaultConstructorAssignor(String unused) {
}
@@ -143,6 +175,59 @@ public class GroupCoordinatorConfigTest {
assertInstanceOf(CustomAssignor.class, assignors.get(1));
}
+ @Test
+ public void testConsumerGroupAssignorsWithDuplicateNamesFails() {
+ // Two custom assignors resolving to the same name must fail startup.
+ Map<String, Object> configs = new HashMap<>();
+ configs.put(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
+ List.of(CustomAssignor.class.getName(),
DuplicateNameAssignor.class.getName()));
+ assertEquals("Invalid value " + DuplicateNameAssignor.class.getName() +
+ " for configuration group.consumer.assignors: Assignor name
'CustomAssignor' is already " +
+ "registered by another configured assignor. Assignor names,
whether built-in or custom, must be unique",
+ assertThrows(ConfigException.class, () ->
createConfig(configs)).getMessage());
+
+ // Configuring the same built-in twice, once by name and once by class
name, must fail startup.
+ configs.put(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
+ List.of("uniform", UniformAssignor.class.getName()));
+ assertEquals("Invalid value " + UniformAssignor.class.getName() +
+ " for configuration group.consumer.assignors: Assignor name
'uniform' is already " +
+ "registered by another configured assignor. Assignor names,
whether built-in or custom, must be unique",
+ assertThrows(ConfigException.class, () ->
createConfig(configs)).getMessage());
+ }
+
+ @Test
+ public void testConsumerGroupAssignorsWithReservedBuiltinNameFails() {
+ // A custom assignor must not take the name of a built-in, whether or
not the built-in is
+ // itself configured: a member selecting that name would otherwise
silently get the custom one.
+ Map<String, Object> configs = new HashMap<>();
+ configs.put(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
UniformNamedAssignor.class.getName());
+ assertEquals("Invalid value " + UniformNamedAssignor.class.getName() +
+ " for configuration group.consumer.assignors: Assignor name
'uniform' is reserved by a " +
+ "built-in assignor. A custom assignor must not reuse the name
of a built-in assignor",
+ assertThrows(ConfigException.class, () ->
createConfig(configs)).getMessage());
+
+ configs.put(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
+ List.of("uniform", UniformNamedAssignor.class.getName()));
+ assertEquals("Invalid value " + UniformNamedAssignor.class.getName() +
+ " for configuration group.consumer.assignors: Assignor name
'uniform' is reserved by a " +
+ "built-in assignor. A custom assignor must not reuse the name
of a built-in assignor",
+ assertThrows(ConfigException.class, () ->
createConfig(configs)).getMessage());
+ }
+
+ @Test
+ public void testConsumerGroupAssignorsBuiltinByClassName() {
+ // A built-in may also be configured by its class name, so the
reserved-name check must
+ // recognise it by class rather than by name.
+ Map<String, Object> configs = new HashMap<>();
+ configs.put(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
+ List.of(UniformAssignor.class.getName(),
RangeAssignor.class.getName()));
+ GroupCoordinatorConfig config = createConfig(configs);
+ List<ConsumerGroupPartitionAssignor> assignors =
config.consumerGroupAssignors();
+ assertEquals(2, assignors.size());
+ assertInstanceOf(UniformAssignor.class, assignors.get(0));
+ assertInstanceOf(RangeAssignor.class, assignors.get(1));
+ }
+
@Test
public void testShareGroupAssignorFullClassNames() {
// The full class name of the assignors is part of our public api.
Hence,