MINOR: Prioritise actionable share group DLQ topic validation errors. (#23656)
ShareGroupDLQStateManager.validateDlqTopic() checked topic existence
(and whether auto-create is enabled) only after delegating to
ShareGroupDLQValidator.validateDlqTopicConfig(). Since the
prefix-compliance check inside validateDlqTopicConfig has no existence
guard, a DLQ topic that was both missing and non-compliant with the
required prefix was reported as a prefix mismatch instead of the more
actionable "does not exist" error.
Move the existence/auto-create check ahead of the delegation, so it is
evaluated first.
In validateDlqTopicConfig the order of checking DLQ prefix and DLQ
enablement on the topic also swapped.
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Reviewers: Andy
Huang (github:andyhuangdev)
Reviewers: Andrew Schofield <aschofield@confluent.io>
---------
Co-authored-by: Claude Sonnet 5 <noreply@anthropic.com>
diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index 357972a..2e6ac0b 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -46,6 +46,7 @@
<!-- server tests -->
<suppress checks="MethodLength|JavaNCSS|NPathComplexity" files="DescribeTopicPartitionsRequestHandlerTest.java"/>
<suppress checks="CyclomaticComplexity" files="ListConsumerGroupTest.java"/>
+ <suppress checks="ClassFanOutComplexity" files="ShareGroupDLQStateManagerTest.java"/>
<suppress checks="ClassFanOutComplexity|CyclomaticComplexity|MethodLength|ParameterNumber|JavaNCSS|ImportControl" files="RequestConvertToJson.java"/>
<suppress checks="ImportControl" files="BrokerRegistrationRequestTest.java"/>
<suppress checks="ImportControl" files="MetadataRequestTest.java"/>
diff --git a/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManager.java b/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManager.java
index cc7be52..9178ac3 100644
--- a/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManager.java
+++ b/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManager.java
@@ -451,25 +451,21 @@
public Optional<Throwable> validateDlqTopic() {
Optional<String> topicNameOpt = cacheHelper.shareGroupDlqTopic(param.groupId());
- // Verify that DLQ topic for the share group is set and is correctly named.
+ // Verify that DLQ topic for the share group is set.
if (topicNameOpt.isEmpty()) {
return Optional.of(new ConfigException(String.format("Configured DLQ topic name in share group: %s is empty.", param.groupId())));
}
String topicName = topicNameOpt.get();
- Optional<Throwable> sharedError = ShareGroupDLQValidator.validateDlqTopicConfig(
- param.groupId(), topicName, topicName, cacheHelper);
- if (sharedError.isPresent()) {
- return sharedError;
- }
-
- // Verify that for a non-existent correctly named DLQ topic, auto create should be enabled.
+ // Verify topic existence (or that auto create is enabled) before the config-only checks in
+ // validateDlqTopicConfig, so a missing topic is reported as missing rather than being masked
+ // by an also-true-but-less-actionable naming/enablement error from that method.
if (!cacheHelper.containsTopic(topicName) && !cacheHelper.isDlqAutoTopicCreateEnabled()) {
return Optional.of(new ConfigException(String.format("DLQ topic does not exist and auto create is disabled on cluster for share group: %s, topic: %s.", param.groupId(), topicName)));
}
- return Optional.empty();
+ return ShareGroupDLQValidator.validateDlqTopicConfig(param.groupId(), topicName, topicName, cacheHelper);
}
public boolean dlqTopicExists() {
diff --git a/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQValidator.java b/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQValidator.java
index 1a5ff10..9a14de1 100644
--- a/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQValidator.java
+++ b/server/src/main/java/org/apache/kafka/server/share/dlq/ShareGroupDLQValidator.java
@@ -62,9 +62,12 @@
}
/**
- * Validates DLQ topic configuration. Checks that the topic name does not start with {@code __},
- * that DLQ is enabled on the topic (if it exists), and that the topic name complies with the
- * configured prefix.
+ * Validates DLQ topic configuration: that the topic name does not start with {@code __}, that
+ * DLQ is enabled on the topic (if it exists), and that the topic name complies with the
+ * configured prefix. This is a config-only check; it does not validate whether the topic
+ * exists or can be created - callers that care about that must check it themselves, and should
+ * do so before calling this method so a missing topic is reported as missing rather than being
+ * masked by an also-true-but-less-actionable naming error from this method.
*
* <p>Callers are responsible for checking that the topic name is present in the config (non-empty)
* before calling this method, and for any implementation-specific checks (e.g., auto-create).
@@ -86,15 +89,11 @@
"Configured DLQ topic name in share group: %s cannot start with __, topic: %s.", groupId, userTopicName)));
}
- // Verify that DLQ is enabled on a correctly named topic, configured on a share group.
- if (cacheHelper.containsTopic(resolvedTopicName) && !cacheHelper.isDlqEnabledOnTopic(resolvedTopicName)) {
- return Optional.of(new ConfigException(
- "DLQ is not enabled on configured DLQ topic for share group: "
- + groupId + ", topic: " + userTopicName));
- }
-
+ // Verify the topic name complies with the configured prefix before the enablement check
+ // below, since this - like the "__" check above - only depends on the configured name
+ // itself, not on the topic's actual (live) state, which the enablement check needs.
Optional<String> topicPrefix = cacheHelper.shareGroupDlqTopicPrefix();
- return topicPrefix.map(prefix -> {
+ Optional<Throwable> prefixError = topicPrefix.map(prefix -> {
if (!prefix.isEmpty() && !userTopicName.startsWith(prefix)) {
return new ConfigException(
"Configured DLQ topic name does not comply with the DLQ topic prefix in share group: "
@@ -102,5 +101,17 @@
}
return null;
});
+ if (prefixError.isPresent()) {
+ return prefixError;
+ }
+
+ // Verify that DLQ is enabled on a correctly named topic, configured on a share group.
+ if (cacheHelper.containsTopic(resolvedTopicName) && !cacheHelper.isDlqEnabledOnTopic(resolvedTopicName)) {
+ return Optional.of(new ConfigException(
+ "DLQ is not enabled on configured DLQ topic for share group: "
+ + groupId + ", topic: " + userTopicName));
+ }
+
+ return Optional.empty();
}
}
diff --git a/server/src/test/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManagerTest.java b/server/src/test/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManagerTest.java
index 441d0b9..b6b2623 100644
--- a/server/src/test/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManagerTest.java
+++ b/server/src/test/java/org/apache/kafka/server/share/dlq/ShareGroupDLQStateManagerTest.java
@@ -54,6 +54,8 @@
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
import org.mockito.ArgumentMatcher;
import org.mockito.Mockito;
@@ -442,6 +444,8 @@
ShareGroupDLQMetadataCacheHelper cacheHelper = mock(ShareGroupDLQMetadataCacheHelper.class);
when(cacheHelper.shareGroupDlqTopic(GROUP_ID)).thenReturn(Optional.of("__internal_dlq"));
when(cacheHelper.shareGroupDlqTopicPrefix()).thenReturn(Optional.empty());
+ // Topic exists, so the existence check passes through to the "__" naming check being tested here.
+ when(cacheHelper.containsTopic("__internal_dlq")).thenReturn(true);
stateManager = builder().withCacheHelper(cacheHelper).build();
stateManager.start();
@@ -468,6 +472,23 @@
}
@Test
+ public void testDlqExistingTopicWithoutDlqConfigAndPrefixMismatchReportsPrefixMismatch() throws Exception {
+ ShareGroupDLQMetadataCacheHelper cacheHelper = mock(ShareGroupDLQMetadataCacheHelper.class);
+ when(cacheHelper.shareGroupDlqTopic(GROUP_ID)).thenReturn(Optional.of(DLQ_TOPIC));
+ when(cacheHelper.shareGroupDlqTopicPrefix()).thenReturn(Optional.of("required-prefix-"));
+ when(cacheHelper.containsTopic(DLQ_TOPIC)).thenReturn(true);
+ when(cacheHelper.isDlqEnabledOnTopic(DLQ_TOPIC)).thenReturn(false);
+
+ stateManager = builder().withCacheHelper(cacheHelper).build();
+ stateManager.start();
+ Throwable cause = getCause(stateManager.dlq(param()));
+ assertInstanceOf(ConfigException.class, cause);
+ assertTrue(cause.getMessage().contains("does not comply with the DLQ topic prefix"));
+ assertFalse(cause.getMessage().contains("DLQ is not enabled"));
+ verifyNoInteractions(mockMetrics);
+ }
+
+ @Test
public void testDlqTopicMissingAndAutoCreateDisabledFailsValidation() throws Exception {
ShareGroupDLQMetadataCacheHelper cacheHelper = mock(ShareGroupDLQMetadataCacheHelper.class);
when(cacheHelper.shareGroupDlqTopic(GROUP_ID)).thenReturn(Optional.of(DLQ_TOPIC));
@@ -483,6 +504,30 @@
verifyNoInteractions(mockMetrics);
}
+ @ParameterizedTest
+ @CsvSource({
+ "false, DLQ topic does not exist, does not comply with the DLQ topic prefix",
+ "true, does not comply with the DLQ topic prefix, DLQ topic does not exist"
+ })
+ public void testDlqTopicMissingAndPrefixMismatchFailsValidation(
+ boolean autoCreateEnabled, String expectedMessageFragment, String unexpectedMessageFragment) throws Exception {
+ ShareGroupDLQMetadataCacheHelper cacheHelper = mock(ShareGroupDLQMetadataCacheHelper.class);
+ when(cacheHelper.shareGroupDlqTopic(GROUP_ID)).thenReturn(Optional.of(DLQ_TOPIC));
+ when(cacheHelper.shareGroupDlqTopicPrefix()).thenReturn(Optional.of("required-prefix-"));
+ when(cacheHelper.containsTopic(DLQ_TOPIC)).thenReturn(false);
+ when(cacheHelper.isDlqAutoTopicCreateEnabled()).thenReturn(autoCreateEnabled);
+
+ stateManager = builder().withCacheHelper(cacheHelper).build();
+ stateManager.start();
+ Throwable cause = getCause(stateManager.dlq(param()));
+ assertInstanceOf(ConfigException.class, cause);
+ assertTrue(cause.getMessage().contains(expectedMessageFragment),
+ "Expected message to contain '" + expectedMessageFragment + "', got: " + cause.getMessage());
+ assertFalse(cause.getMessage().contains(unexpectedMessageFragment),
+ "Did not expect message to contain '" + unexpectedMessageFragment + "', got: " + cause.getMessage());
+ verifyNoInteractions(mockMetrics);
+ }
+
@Test
public void testDlqTopicPrefixMismatchFailsValidation() throws Exception {
ShareGroupDLQMetadataCacheHelper cacheHelper = mock(ShareGroupDLQMetadataCacheHelper.class);