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,

Reply via email to