diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java index ffa658fa50ba7..4e2d52fe16de2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java @@ -219,8 +219,7 @@ private void doValidateTopic(final ValidationResult validationResult, } catch (final ExecutionException executionException) { final Throwable cause = executionException.getCause(); if (cause instanceof UnknownTopicOrPartitionException) { - log.info("Internal topic {} is missing", - topicName); + log.info("Internal topic {} is missing", topicName); validationResult.addMissingTopic(topicName); topicsStillToValidate.remove(topicName); } else if (cause instanceof LeaderNotAvailableException) { @@ -296,6 +295,10 @@ private void validateCleanupPolicy(final ValidationResult validationResult, validateCleanupPolicyForUnwindowedChangelogs(validationResult, topicConfig, brokerSideTopicConfig); } else if (topicConfig instanceof WindowedChangelogTopicConfig) { validateCleanupPolicyForWindowedChangelogs(validationResult, topicConfig, brokerSideTopicConfig); + } else if (topicConfig instanceof RepartitionTopicConfig) { + validateCleanupPolicyForRepartitionTopic(validationResult, topicConfig, brokerSideTopicConfig); + } else { + throw new IllegalStateException("Internal topic " + topicConfig.name() + " has unknown type."); } } @@ -307,7 +310,8 @@ private void validateCleanupPolicyForUnwindowedChangelogs(final ValidationResult if (cleanupPolicy.contains(TopicConfig.CLEANUP_POLICY_DELETE)) { validationResult.addMisconfiguration( topicName, - "Cleanup policy of existing internal topic " + topicName + " should not contain \"" + "Cleanup policy (" + TopicConfig.CLEANUP_POLICY_CONFIG + ") of existing internal topic " + + topicName + " should not contain \"" + TopicConfig.CLEANUP_POLICY_DELETE + "\"." ); } @@ -327,8 +331,8 @@ private void validateCleanupPolicyForWindowedChangelogs(final ValidationResult v if (brokerSideRetentionMs < streamsSideRetentionMs) { validationResult.addMisconfiguration( topicName, - "Retention time of existing internal topic " + topicName + " is " + brokerSideRetentionMs + - " but should be " + streamsSideRetentionMs + " or larger." + "Retention time (" + TopicConfig.RETENTION_MS_CONFIG + ") of existing internal topic " + + topicName + " is " + brokerSideRetentionMs + " but should be " + streamsSideRetentionMs + " or larger." ); } final String brokerSideRetentionBytes = @@ -336,7 +340,41 @@ private void validateCleanupPolicyForWindowedChangelogs(final ValidationResult v if (brokerSideRetentionBytes != null) { validationResult.addMisconfiguration( topicName, - "Retention byte of existing internal topic " + topicName + " is set but it should be unset." + "Retention byte (" + TopicConfig.RETENTION_BYTES_CONFIG + ") of existing internal topic " + + topicName + " is set but it should be unset." + ); + } + } + } + + private void validateCleanupPolicyForRepartitionTopic(final ValidationResult validationResult, + final InternalTopicConfig topicConfig, + final Config brokerSideTopicConfig) { + final String topicName = topicConfig.name(); + final String cleanupPolicy = getBrokerSideConfigValue(brokerSideTopicConfig, TopicConfig.CLEANUP_POLICY_CONFIG, topicName); + if (cleanupPolicy.contains(TopicConfig.CLEANUP_POLICY_COMPACT)) { + validationResult.addMisconfiguration( + topicName, + "Cleanup policy (" + TopicConfig.CLEANUP_POLICY_CONFIG + ") of existing internal topic " + + topicName + " should not contain \"" + TopicConfig.CLEANUP_POLICY_COMPACT + "\"." + ); + } else if (cleanupPolicy.contains(TopicConfig.CLEANUP_POLICY_DELETE)) { + final long brokerSideRetentionMs = + Long.parseLong(getBrokerSideConfigValue(brokerSideTopicConfig, TopicConfig.RETENTION_MS_CONFIG, topicName)); + if (brokerSideRetentionMs != -1) { + validationResult.addMisconfiguration( + topicName, + "Retention time (" + TopicConfig.RETENTION_MS_CONFIG + ") of existing internal topic " + + topicName + " is " + brokerSideRetentionMs + " but should be -1." + ); + } + final String brokerSideRetentionBytes = + getBrokerSideConfigValue(brokerSideTopicConfig, TopicConfig.RETENTION_BYTES_CONFIG, topicName); + if (brokerSideRetentionBytes != null) { + validationResult.addMisconfiguration( + topicName, + "Retention byte (" + TopicConfig.RETENTION_BYTES_CONFIG + ") of existing internal topic " + + topicName + " is set but it should be unset." ); } } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/InternalTopicManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/InternalTopicManagerTest.java index d5f568a8d8585..16f5a0bb7f628 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/InternalTopicManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/InternalTopicManagerTest.java @@ -56,6 +56,7 @@ import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.stream.Collectors; import static org.apache.kafka.common.utils.Utils.mkEntry; import static org.apache.kafka.common.utils.Utils.mkMap; @@ -468,16 +469,15 @@ public void shouldExhaustRetriesOnMarkedForDeletionTopic() { @Test public void shouldValidateSuccessfully() { - mockAdminClient.addTopic( - false, - topic1, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null - ); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + setupTopicInMockAdminClient(topic1, repartitionTopicConfig()); + setupTopicInMockAdminClient(topic2, repartitionTopicConfig()); + final InternalTopicConfig internalTopicConfig1 = setupRepartitionTopicConfig(topic1, 1); + final InternalTopicConfig internalTopicConfig2 = setupRepartitionTopicConfig(topic2, 1); - final ValidationResult validationResult = internalTopicManager.validate(Collections.singletonMap(topic1, internalTopicConfig)); + final ValidationResult validationResult = internalTopicManager.validate(mkMap( + mkEntry(topic1, internalTopicConfig1), + mkEntry(topic2, internalTopicConfig2) + )); assertThat(validationResult.missingTopics(), empty()); assertThat(validationResult.misconfigurationsForTopics(), anEmptyMap()); @@ -485,12 +485,7 @@ public void shouldValidateSuccessfully() { @Test public void shouldValidateSuccessfullyWithEmptyInternalTopics() { - mockAdminClient.addTopic( - false, - topic1, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null - ); + setupTopicInMockAdminClient(topic1, repartitionTopicConfig()); final ValidationResult validationResult = internalTopicManager.validate(Collections.emptyMap()); @@ -502,18 +497,10 @@ public void shouldValidateSuccessfullyWithEmptyInternalTopics() { public void shouldReportMissingTopics() { final String missingTopic1 = "missingTopic1"; final String missingTopic2 = "missingTopic2"; - mockAdminClient.addTopic( - false, - topic1, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null - ); - final InternalTopicConfig internalTopicConfig1 = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig1.setNumberOfPartitions(1); - final InternalTopicConfig internalTopicConfig2 = new RepartitionTopicConfig(missingTopic1, Collections.emptyMap()); - internalTopicConfig2.setNumberOfPartitions(1); - final InternalTopicConfig internalTopicConfig3 = new RepartitionTopicConfig(missingTopic2, Collections.emptyMap()); - internalTopicConfig3.setNumberOfPartitions(1); + setupTopicInMockAdminClient(topic1, repartitionTopicConfig()); + final InternalTopicConfig internalTopicConfig1 = setupRepartitionTopicConfig(topic1, 1); + final InternalTopicConfig internalTopicConfig2 = setupRepartitionTopicConfig(missingTopic1, 1); + final InternalTopicConfig internalTopicConfig3 = setupRepartitionTopicConfig(missingTopic2, 1); final ValidationResult validationResult = internalTopicManager.validate(mkMap( mkEntry(topic1, internalTopicConfig1), @@ -530,30 +517,12 @@ public void shouldReportMissingTopics() { @Test public void shouldReportMisconfigurationsOfPartitionCount() { - mockAdminClient.addTopic( - false, - topic1, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null - ); - mockAdminClient.addTopic( - false, - topic2, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null - ); - mockAdminClient.addTopic( - false, - topic3, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null - ); - final InternalTopicConfig internalTopicConfig1 = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig1.setNumberOfPartitions(2); - final InternalTopicConfig internalTopicConfig2 = new RepartitionTopicConfig(topic2, Collections.emptyMap()); - internalTopicConfig2.setNumberOfPartitions(3); - final InternalTopicConfig internalTopicConfig3 = new RepartitionTopicConfig(topic3, Collections.emptyMap()); - internalTopicConfig3.setNumberOfPartitions(1); + setupTopicInMockAdminClient(topic1, repartitionTopicConfig()); + setupTopicInMockAdminClient(topic2, repartitionTopicConfig()); + setupTopicInMockAdminClient(topic3, repartitionTopicConfig()); + final InternalTopicConfig internalTopicConfig1 = setupRepartitionTopicConfig(topic1, 2); + final InternalTopicConfig internalTopicConfig2 = setupRepartitionTopicConfig(topic2, 3); + final InternalTopicConfig internalTopicConfig3 = setupRepartitionTopicConfig(topic3, 1); final ValidationResult validationResult = internalTopicManager.validate(mkMap( mkEntry(topic1, internalTopicConfig1), @@ -581,27 +550,19 @@ public void shouldReportMisconfigurationsOfPartitionCount() { @Test public void shouldReportMisconfigurationsOfCleanupPolicyForUnwindowedChangelogTopics() { - mockAdminClient.addTopic( - false, - topic1, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - mkMap(mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE)) - ); - mockAdminClient.addTopic( - false, - topic2, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - mkMap(mkEntry( - TopicConfig.CLEANUP_POLICY_CONFIG, - TopicConfig.CLEANUP_POLICY_COMPACT + "," + TopicConfig.CLEANUP_POLICY_DELETE - )) + final Map unwindowedChangelogConfigWithDeleteCleanupPolicy = unwindowedChangelogConfig(); + unwindowedChangelogConfigWithDeleteCleanupPolicy.put( + TopicConfig.CLEANUP_POLICY_CONFIG, + TopicConfig.CLEANUP_POLICY_DELETE ); - mockAdminClient.addTopic( - false, - topic3, - Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - mkMap(mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT)) + setupTopicInMockAdminClient(topic1, unwindowedChangelogConfigWithDeleteCleanupPolicy); + final Map unwindowedChangelogConfigWithDeleteCompactCleanupPolicy = unwindowedChangelogConfig(); + unwindowedChangelogConfigWithDeleteCompactCleanupPolicy.put( + TopicConfig.CLEANUP_POLICY_CONFIG, + TopicConfig.CLEANUP_POLICY_COMPACT + "," + TopicConfig.CLEANUP_POLICY_DELETE ); + setupTopicInMockAdminClient(topic2, unwindowedChangelogConfigWithDeleteCompactCleanupPolicy); + setupTopicInMockAdminClient(topic3, unwindowedChangelogConfig()); final InternalTopicConfig internalTopicConfig1 = setupUnwindowedChangelogTopicConfig(topic1, 1); final InternalTopicConfig internalTopicConfig2 = setupUnwindowedChangelogTopicConfig(topic2, 1); final InternalTopicConfig internalTopicConfig3 = setupUnwindowedChangelogTopicConfig(topic3, 1); @@ -619,38 +580,86 @@ public void shouldReportMisconfigurationsOfCleanupPolicyForUnwindowedChangelogTo assertThat(misconfigurationsForTopics.get(topic1).size(), is(1)); assertThat( misconfigurationsForTopics.get(topic1).get(0), - is("Cleanup policy of existing internal topic " + topic1 + " should not contain \"" + is("Cleanup policy (" + TopicConfig.CLEANUP_POLICY_CONFIG + ") of existing internal topic " + topic1 + " should not contain \"" + TopicConfig.CLEANUP_POLICY_DELETE + "\".") ); assertThat(misconfigurationsForTopics, hasKey(topic2)); assertThat(misconfigurationsForTopics.get(topic2).size(), is(1)); assertThat( misconfigurationsForTopics.get(topic2).get(0), - is("Cleanup policy of existing internal topic " + topic2 + " should not contain \"" + is("Cleanup policy (" + TopicConfig.CLEANUP_POLICY_CONFIG + ") of existing internal topic " + topic2 + " should not contain \"" + TopicConfig.CLEANUP_POLICY_DELETE + "\".") ); assertThat(misconfigurationsForTopics, not(hasKey(topic3))); } - private InternalTopicConfig setupUnwindowedChangelogTopicConfig(final String topicName, - final int partitionCount) { - final InternalTopicConfig internalTopicConfig = - new UnwindowedChangelogTopicConfig(topicName, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(partitionCount); - return internalTopicConfig; - } - @Test public void shouldReportMisconfigurationsOfCleanupPolicyForWindowedChangelogTopics() { final long retentionMs = 1000; final long shorterRetentionMs = 900; + setupTopicInMockAdminClient(topic1, windowedChangelogConfig(retentionMs)); + setupTopicInMockAdminClient(topic2, windowedChangelogConfig(shorterRetentionMs)); + final Map windowedChangelogConfigOnlyCleanupPolicyCompact = windowedChangelogConfig(retentionMs); + windowedChangelogConfigOnlyCleanupPolicyCompact.put(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT); + setupTopicInMockAdminClient(topic3, windowedChangelogConfigOnlyCleanupPolicyCompact); + final Map windowedChangelogConfigOnlyCleanupPolicyDelete = windowedChangelogConfig(shorterRetentionMs); + windowedChangelogConfigOnlyCleanupPolicyDelete.put(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE); + setupTopicInMockAdminClient(topic4, windowedChangelogConfigOnlyCleanupPolicyDelete); + final Map windowedChangelogConfigWithRetentionBytes = windowedChangelogConfig(retentionMs); + windowedChangelogConfigWithRetentionBytes.put(TopicConfig.RETENTION_BYTES_CONFIG, "1024"); + setupTopicInMockAdminClient(topic5, windowedChangelogConfigWithRetentionBytes); + final InternalTopicConfig internalTopicConfig1 = setupWindowedChangelogTopicConfig(topic1, 1, retentionMs); + final InternalTopicConfig internalTopicConfig2 = setupWindowedChangelogTopicConfig(topic2, 1, retentionMs); + final InternalTopicConfig internalTopicConfig3 = setupWindowedChangelogTopicConfig(topic3, 1, retentionMs); + final InternalTopicConfig internalTopicConfig4 = setupWindowedChangelogTopicConfig(topic4, 1, retentionMs); + final InternalTopicConfig internalTopicConfig5 = setupWindowedChangelogTopicConfig(topic5, 1, retentionMs); + + final ValidationResult validationResult = internalTopicManager.validate(mkMap( + mkEntry(topic1, internalTopicConfig1), + mkEntry(topic2, internalTopicConfig2), + mkEntry(topic3, internalTopicConfig3), + mkEntry(topic4, internalTopicConfig4), + mkEntry(topic5, internalTopicConfig5) + )); + + final Map> misconfigurationsForTopics = validationResult.misconfigurationsForTopics(); + assertThat(validationResult.missingTopics(), empty()); + assertThat(misconfigurationsForTopics.size(), is(3)); + assertThat(misconfigurationsForTopics, hasKey(topic2)); + assertThat(misconfigurationsForTopics.get(topic2).size(), is(1)); + assertThat( + misconfigurationsForTopics.get(topic2).get(0), + is("Retention time (" + TopicConfig.RETENTION_MS_CONFIG + ") of existing internal topic " + + topic2 + " is " + shorterRetentionMs + " but should be " + retentionMs + " or larger.") + ); + assertThat(misconfigurationsForTopics, hasKey(topic4)); + assertThat(misconfigurationsForTopics.get(topic4).size(), is(1)); + assertThat( + misconfigurationsForTopics.get(topic4).get(0), + is("Retention time (" + TopicConfig.RETENTION_MS_CONFIG + ") of existing internal topic " + + topic4 + " is " + shorterRetentionMs + " but should be " + retentionMs + " or larger.") + ); + assertThat(misconfigurationsForTopics, hasKey(topic5)); + assertThat(misconfigurationsForTopics.get(topic5).size(), is(1)); + assertThat( + misconfigurationsForTopics.get(topic5).get(0), + is("Retention byte (" + TopicConfig.RETENTION_BYTES_CONFIG + ") of existing internal topic " + + topic5 + " is set but it should be unset.") + ); + assertThat(misconfigurationsForTopics, not(hasKey(topic1))); + assertThat(misconfigurationsForTopics, not(hasKey(topic3))); + } + + @Test + public void shouldReportMisconfigurationsOfCleanupPolicyForRepartitionTopics() { + final long retentionMs = 1000; mockAdminClient.addTopic( false, topic1, Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), mkMap( - mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT + "," + TopicConfig.CLEANUP_POLICY_DELETE), - mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs)), + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(-1)), mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) ) ); @@ -659,8 +668,8 @@ public void shouldReportMisconfigurationsOfCleanupPolicyForWindowedChangelogTopi topic2, Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), mkMap( - mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT + "," + TopicConfig.CLEANUP_POLICY_DELETE), - mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(shorterRetentionMs)), + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(-1)), mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) ) ); @@ -669,8 +678,8 @@ public void shouldReportMisconfigurationsOfCleanupPolicyForWindowedChangelogTopi topic3, Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), mkMap( - mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT), - mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(shorterRetentionMs)), + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT + "," + TopicConfig.CLEANUP_POLICY_DELETE), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(-1)), mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) ) ); @@ -680,7 +689,7 @@ public void shouldReportMisconfigurationsOfCleanupPolicyForWindowedChangelogTopi Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), mkMap( mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), - mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(shorterRetentionMs)), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs)), mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) ) ); @@ -690,15 +699,15 @@ public void shouldReportMisconfigurationsOfCleanupPolicyForWindowedChangelogTopi Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), mkMap( mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), - mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs)), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(-1)), mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, "1024") ) ); - final InternalTopicConfig internalTopicConfig1 = setupWindowedChangelogTopicConfig(topic1, 1, retentionMs); - final InternalTopicConfig internalTopicConfig2 = setupWindowedChangelogTopicConfig(topic2, 1, retentionMs); - final InternalTopicConfig internalTopicConfig3 = setupWindowedChangelogTopicConfig(topic3, 1, retentionMs); - final InternalTopicConfig internalTopicConfig4 = setupWindowedChangelogTopicConfig(topic4, 1, retentionMs); - final InternalTopicConfig internalTopicConfig5 = setupWindowedChangelogTopicConfig(topic5, 1, retentionMs); + final InternalTopicConfig internalTopicConfig1 = setupRepartitionTopicConfig(topic1, 1); + final InternalTopicConfig internalTopicConfig2 = setupRepartitionTopicConfig(topic2, 1); + final InternalTopicConfig internalTopicConfig3 = setupRepartitionTopicConfig(topic3, 1); + final InternalTopicConfig internalTopicConfig4 = setupRepartitionTopicConfig(topic4, 1); + final InternalTopicConfig internalTopicConfig5 = setupRepartitionTopicConfig(topic5, 1); final ValidationResult validationResult = internalTopicManager.validate(mkMap( mkEntry(topic1, internalTopicConfig1), @@ -710,29 +719,35 @@ public void shouldReportMisconfigurationsOfCleanupPolicyForWindowedChangelogTopi final Map> misconfigurationsForTopics = validationResult.misconfigurationsForTopics(); assertThat(validationResult.missingTopics(), empty()); - assertThat(misconfigurationsForTopics.size(), is(3)); + assertThat(misconfigurationsForTopics.size(), is(4)); assertThat(misconfigurationsForTopics, hasKey(topic2)); assertThat(misconfigurationsForTopics.get(topic2).size(), is(1)); assertThat( misconfigurationsForTopics.get(topic2).get(0), - is("Retention time of existing internal topic " + topic2 + " is " + shorterRetentionMs + - " but should be " + retentionMs + " or larger.") + is("Cleanup policy (" + TopicConfig.CLEANUP_POLICY_CONFIG + ") of existing internal topic " + + topic2 + " should not contain \"" + TopicConfig.CLEANUP_POLICY_COMPACT + "\".") + ); + assertThat(misconfigurationsForTopics, hasKey(topic3)); + assertThat(misconfigurationsForTopics.get(topic3).size(), is(1)); + assertThat( + misconfigurationsForTopics.get(topic3).get(0), + is("Cleanup policy (" + TopicConfig.CLEANUP_POLICY_CONFIG + ") of existing internal topic " + + topic3 + " should not contain \"" + TopicConfig.CLEANUP_POLICY_COMPACT + "\".") ); assertThat(misconfigurationsForTopics, hasKey(topic4)); assertThat(misconfigurationsForTopics.get(topic4).size(), is(1)); assertThat( misconfigurationsForTopics.get(topic4).get(0), - is("Retention time of existing internal topic " + topic4 + " is " + shorterRetentionMs + - " but should be " + retentionMs + " or larger.") + is("Retention time (" + TopicConfig.RETENTION_MS_CONFIG + ") of existing internal topic " + + topic4 + " is " + retentionMs + " but should be -1.") ); assertThat(misconfigurationsForTopics, hasKey(topic5)); assertThat(misconfigurationsForTopics.get(topic5).size(), is(1)); assertThat( misconfigurationsForTopics.get(topic5).get(0), - is("Retention byte of existing internal topic " + topic5 + " is set but it should be unset.") + is("Retention byte (" + TopicConfig.RETENTION_BYTES_CONFIG + ") of existing internal topic " + + topic5 + " is set but it should be unset.") ); - assertThat(misconfigurationsForTopics, not(hasKey(topic1))); - assertThat(misconfigurationsForTopics, not(hasKey(topic3))); } @Test @@ -762,24 +777,14 @@ public void shouldReportMultipleMisconfigurationsForSameTopic() { assertThat(misconfigurationsForTopics.get(topic1).size(), is(2)); assertThat( misconfigurationsForTopics.get(topic1).get(0), - is("Retention time of existing internal topic " + topic1 + " is " + shorterRetentionMs + - " but should be " + retentionMs + " or larger.") + is("Retention time (" + TopicConfig.RETENTION_MS_CONFIG + ") of existing internal topic " + + topic1 + " is " + shorterRetentionMs + " but should be " + retentionMs + " or larger.") ); assertThat( misconfigurationsForTopics.get(topic1).get(1), - is("Retention byte of existing internal topic " + topic1 + " is set but it should be unset.") - ); - } - - private InternalTopicConfig setupWindowedChangelogTopicConfig(final String topicName, - final int partitionCount, - final long retentionMs) { - final InternalTopicConfig internalTopicConfig = new WindowedChangelogTopicConfig( - topicName, - mkMap(mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs))) + is("Retention byte (" + TopicConfig.RETENTION_BYTES_CONFIG + ") of existing internal topic " + + topic1 + " is set but it should be unset.") ); - internalTopicConfig.setNumberOfPartitions(partitionCount); - return internalTopicConfig; } @Test @@ -788,7 +793,7 @@ public void shouldThrowWhenPartitionCountUnknown() { false, topic1, Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null + mkMap(mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE)) ); final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); @@ -804,11 +809,14 @@ public void shouldRetryWhenCallsThrowTimeoutExceptionDuringValidation() { false, topic1, Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), - null + mkMap( + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, "-1"), + mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) + ) ); mockAdminClient.timeoutNextRequest(2); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); final ValidationResult validationResult = internalTopicManager.validate(Collections.singletonMap(topic1, internalTopicConfig)); @@ -836,13 +844,15 @@ public void shouldOnlyRetryDescribeTopicsWhenDescribeTopicsThrowsLeaderNotAvaila .andReturn(new MockDescribeTopicsResult(mkMap(mkEntry(topic1, topicDescriptionFailFuture)))) .andReturn(new MockDescribeTopicsResult(mkMap(mkEntry(topic1, topicDescriptionSuccessfulFuture)))); final KafkaFutureImpl topicConfigSuccessfulFuture = new KafkaFutureImpl<>(); - topicConfigSuccessfulFuture.complete(new Config(Collections.emptySet())); + topicConfigSuccessfulFuture.complete( + new Config(repartitionTopicConfig().entrySet().stream() + .map(entry -> new ConfigEntry(entry.getKey(), entry.getValue())).collect(Collectors.toSet())) + ); final ConfigResource topicResource = new ConfigResource(Type.TOPIC, topic1); EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource))) .andReturn(new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, topicConfigSuccessfulFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); final ValidationResult validationResult = topicManager.validate(Collections.singletonMap(topic1, internalTopicConfig)); @@ -870,14 +880,16 @@ public void shouldOnlyRetryDescribeConfigsWhenDescribeConfigsThrowsLeaderNotAvai final KafkaFutureImpl topicConfigsFailFuture = new KafkaFutureImpl<>(); topicConfigsFailFuture.completeExceptionally(new LeaderNotAvailableException("Leader Not Available!")); final KafkaFutureImpl topicConfigSuccessfulFuture = new KafkaFutureImpl<>(); - topicConfigSuccessfulFuture.complete(new Config(Collections.emptySet())); + topicConfigSuccessfulFuture.complete( + new Config(repartitionTopicConfig().entrySet().stream() + .map(entry -> new ConfigEntry(entry.getKey(), entry.getValue())).collect(Collectors.toSet())) + ); final ConfigResource topicResource = new ConfigResource(Type.TOPIC, topic1); EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource))) .andReturn(new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, topicConfigsFailFuture)))) .andReturn(new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, topicConfigSuccessfulFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); final ValidationResult validationResult = topicManager.validate(Collections.singletonMap(topic1, internalTopicConfig)); @@ -899,8 +911,7 @@ public void shouldThrowWhenDescribeTopicsThrowsUnexpectedExceptionDuringValidati EasyMock.expect(admin.describeTopics(Collections.singleton(topic1))) .andStubAnswer(() -> new MockDescribeTopicsResult(mkMap(mkEntry(topic1, topicDescriptionFailFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); assertThrows(Throwable.class, () -> topicManager.validate(Collections.singletonMap(topic1, internalTopicConfig))); } @@ -919,8 +930,7 @@ public void shouldThrowWhenDescribeConfigsThrowsUnexpectedExceptionDuringValidat EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource))) .andStubAnswer(() -> new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, configDescriptionFailFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); assertThrows(Throwable.class, () -> topicManager.validate(Collections.singletonMap(topic1, internalTopicConfig))); } @@ -947,8 +957,7 @@ public void shouldThrowWhenTopicDescriptionsDoNotContainTopicDuringValidation() EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource))) .andStubAnswer(() -> new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, topicConfigSuccessfulFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); assertThrows( IllegalStateException.class, @@ -979,8 +988,7 @@ public void shouldThrowWhenConfigDescriptionsDoNotContainTopicDuringValidation() EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource1))) .andStubAnswer(() -> new MockDescribeConfigsResult(mkMap(mkEntry(topicResource2, topicConfigSuccessfulFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); assertThrows( IllegalStateException.class, @@ -992,42 +1000,65 @@ public void shouldThrowWhenConfigDescriptionsDoNotContainTopicDuringValidation() public void shouldThrowWhenConfigDescriptionsDoNotCleanupPolicyForUnwindowedConfigDuringValidation() { shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( setupUnwindowedChangelogTopicConfig(topic1, 1), - new Config(Collections.emptySet()) + configWithoutKey(unwindowedChangelogConfig(), TopicConfig.CLEANUP_POLICY_CONFIG) ); } @Test - public void shouldThrowWhenConfigDescriptionsDoNotCleanupPolicyForWindowedConfigDuringValidation() { + public void shouldThrowWhenConfigDescriptionsDoNotContainCleanupPolicyForWindowedConfigDuringValidation() { final long retentionMs = 1000; shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( setupWindowedChangelogTopicConfig(topic1, 1, retentionMs), - new Config(mkSet( - new ConfigEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs)), - new ConfigEntry(TopicConfig.RETENTION_BYTES_CONFIG, "1000")) - ) + configWithoutKey(windowedChangelogConfig(retentionMs), TopicConfig.CLEANUP_POLICY_CONFIG) ); } @Test - public void shouldThrowWhenConfigDescriptionsDoNotRetentionMsForWindowedConfigDuringValidation() { + public void shouldThrowWhenConfigDescriptionsDoNotContainRetentionMsForWindowedConfigDuringValidation() { + final long retentionMs = 1000; shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( - setupWindowedChangelogTopicConfig(topic1, 1, 1000), - new Config(mkSet( - new ConfigEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), - new ConfigEntry(TopicConfig.RETENTION_BYTES_CONFIG, "1000")) - ) + setupWindowedChangelogTopicConfig(topic1, 1, retentionMs), + configWithoutKey(windowedChangelogConfig(retentionMs), TopicConfig.RETENTION_MS_CONFIG) ); } @Test - public void shouldThrowWhenConfigDescriptionsDoNotRetentionBytesForWindowedConfigDuringValidation() { + public void shouldThrowWhenConfigDescriptionsDoNotContainRetentionBytesForWindowedConfigDuringValidation() { final long retentionMs = 1000; shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( setupWindowedChangelogTopicConfig(topic1, 1, retentionMs), - new Config(mkSet( - new ConfigEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), - new ConfigEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs))) - ) + configWithoutKey(windowedChangelogConfig(retentionMs), TopicConfig.RETENTION_BYTES_CONFIG) + ); + } + + @Test + public void shouldThrowWhenConfigDescriptionsDoNotContainCleanupPolicyForRepartitionConfigDuringValidation() { + shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( + setupRepartitionTopicConfig(topic1, 1), + configWithoutKey(repartitionTopicConfig(), TopicConfig.CLEANUP_POLICY_CONFIG) + ); + } + + @Test + public void shouldThrowWhenConfigDescriptionsDoNotContainRetentionMsForRepartitionConfigDuringValidation() { + shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( + setupRepartitionTopicConfig(topic1, 1), + configWithoutKey(repartitionTopicConfig(), TopicConfig.RETENTION_MS_CONFIG) + ); + } + + @Test + public void shouldThrowWhenConfigDescriptionsDoNotContainRetentionBytesForRepartitionConfigDuringValidation() { + shouldThrowWhenConfigDescriptionsDoNotContainConfigDuringValidation( + setupRepartitionTopicConfig(topic1, 1), + configWithoutKey(repartitionTopicConfig(), TopicConfig.RETENTION_BYTES_CONFIG) + ); + } + + private Config configWithoutKey(final Map config, final String key) { + return new Config(config.entrySet().stream() + .filter(entry -> !entry.getKey().equals(key)) + .map(entry -> new ConfigEntry(entry.getKey(), entry.getValue())).collect(Collectors.toSet()) ); } @@ -1076,13 +1107,15 @@ public void shouldThrowTimeoutExceptionWhenTimeoutIsExceededDuringValidation() { EasyMock.expect(admin.describeTopics(Collections.singleton(topic1))) .andStubAnswer(() -> new MockDescribeTopicsResult(mkMap(mkEntry(topic1, topicDescriptionFailFuture)))); final KafkaFutureImpl topicConfigSuccessfulFuture = new KafkaFutureImpl<>(); - topicConfigSuccessfulFuture.complete(new Config(Collections.emptySet())); + topicConfigSuccessfulFuture.complete( + new Config(repartitionTopicConfig().entrySet().stream() + .map(entry -> new ConfigEntry(entry.getKey(), entry.getValue())).collect(Collectors.toSet())) + ); final ConfigResource topicResource = new ConfigResource(Type.TOPIC, topic1); EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource))) .andStubAnswer(() -> new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, topicConfigSuccessfulFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); assertThrows( TimeoutException.class, @@ -1105,13 +1138,15 @@ public void shouldThrowTimeoutExceptionWhenFuturesNeverCompleteDuringValidation( EasyMock.expect(admin.describeTopics(Collections.singleton(topic1))) .andStubAnswer(() -> new MockDescribeTopicsResult(mkMap(mkEntry(topic1, topicDescriptionFutureThatNeverCompletes)))); final KafkaFutureImpl topicConfigSuccessfulFuture = new KafkaFutureImpl<>(); - topicConfigSuccessfulFuture.complete(new Config(Collections.emptySet())); + topicConfigSuccessfulFuture.complete( + new Config(repartitionTopicConfig().entrySet().stream() + .map(entry -> new ConfigEntry(entry.getKey(), entry.getValue())).collect(Collectors.toSet())) + ); final ConfigResource topicResource = new ConfigResource(Type.TOPIC, topic1); EasyMock.expect(admin.describeConfigs(Collections.singleton(topicResource))) .andStubAnswer(() -> new MockDescribeConfigsResult(mkMap(mkEntry(topicResource, topicConfigSuccessfulFuture)))); EasyMock.replay(admin); - final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topic1, Collections.emptyMap()); - internalTopicConfig.setNumberOfPartitions(1); + final InternalTopicConfig internalTopicConfig = setupRepartitionTopicConfig(topic1, 1); assertThrows( TimeoutException.class, @@ -1119,6 +1154,63 @@ public void shouldThrowTimeoutExceptionWhenFuturesNeverCompleteDuringValidation( ); } + private Map repartitionTopicConfig() { + return mkMap( + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_DELETE), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, "-1"), + mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) + ); + } + + private Map unwindowedChangelogConfig() { + return mkMap( + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT) + ); + } + + private Map windowedChangelogConfig(final long retentionMs) { + return mkMap( + mkEntry(TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT + "," + TopicConfig.CLEANUP_POLICY_DELETE), + mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs)), + mkEntry(TopicConfig.RETENTION_BYTES_CONFIG, null) + ); + } + + private void setupTopicInMockAdminClient(final String topic, final Map topicConfig) { + mockAdminClient.addTopic( + false, + topic, + Collections.singletonList(new TopicPartitionInfo(0, broker1, cluster, Collections.emptyList())), + topicConfig + ); + } + + private InternalTopicConfig setupUnwindowedChangelogTopicConfig(final String topicName, + final int partitionCount) { + final InternalTopicConfig internalTopicConfig = + new UnwindowedChangelogTopicConfig(topicName, Collections.emptyMap()); + internalTopicConfig.setNumberOfPartitions(partitionCount); + return internalTopicConfig; + } + + private InternalTopicConfig setupWindowedChangelogTopicConfig(final String topicName, + final int partitionCount, + final long retentionMs) { + final InternalTopicConfig internalTopicConfig = new WindowedChangelogTopicConfig( + topicName, + mkMap(mkEntry(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(retentionMs))) + ); + internalTopicConfig.setNumberOfPartitions(partitionCount); + return internalTopicConfig; + } + + private InternalTopicConfig setupRepartitionTopicConfig(final String topicName, + final int partitionCount) { + final InternalTopicConfig internalTopicConfig = new RepartitionTopicConfig(topicName, Collections.emptyMap()); + internalTopicConfig.setNumberOfPartitions(partitionCount); + return internalTopicConfig; + } + private static class MockCreateTopicsResult extends CreateTopicsResult { MockCreateTopicsResult(final Map> futures) {