chia7712 commented on code in PR #20384:
URL: https://github.com/apache/kafka/pull/20384#discussion_r3667201208


##########
connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaTopicBasedBackingStore.java:
##########
@@ -39,17 +42,35 @@ public abstract class KafkaTopicBasedBackingStore {
 
     Consumer<TopicAdmin> topicInitializer(String topic, NewTopic 
topicDescription, WorkerConfig config, Time time) {
         return admin -> {
-            log.debug("Creating Connect internal topic for {}", 
getTopicPurpose());
-            // Create the topic if it doesn't exist
-            Set<String> newTopics = createTopics(topicDescription, admin, 
config, time);
-            if (!newTopics.contains(topic)) {
-                // It already existed, so check that the topic cleanup policy 
is compact only and not delete
-                log.debug("Using admin client to check cleanup policy of '{}' 
topic is '{}'", topic, TopicConfig.CLEANUP_POLICY_COMPACT);
-                admin.verifyTopicCleanupPolicyOnlyCompact(topic, 
getTopicConfig(), getTopicPurpose());
+            if (config.internalTopicsCreationEnabled()) {
+                log.debug("Creating Connect internal topic for {}", 
getTopicPurpose());
+                // Create the topic if it doesn't exist
+                Set<String> newTopics = createTopics(topicDescription, admin, 
config, time);
+                if (!newTopics.contains(topic)) {
+                    verifyTopicConfig(topic, admin);
+                }
+            } else {
+                log.debug("Skipping creation of Connect internal topic for {} 
because automatic topic creation is disabled", getTopicPurpose());
+                Map<String, TopicDescription> existing = 
admin.describeTopics(topic);
+                if (existing.isEmpty()) {
+                    String msg = String.format("Topic '%s' specified via the 
'%s' property is missing." +

Review Comment:
   Is this error message correct if it's a connector-specific offset topic? 
Also, what `getTopicConfig()` returns wouldn't be the correct config key in 
that case, right?



##########
connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedConfig.java:
##########
@@ -187,6 +188,13 @@ public final class DistributedConfig extends WorkerConfig {
     public static final String CONNECT_PROTOCOL_DOC = "Compatibility mode for 
Kafka Connect Protocol";
     public static final String CONNECT_PROTOCOL_DEFAULT = 
ConnectProtocolCompatibility.SESSIONED.toString();
 
+
+    public static final String 
INTERNAL_TOPICS_AUTOMATIC_CREATION_ENABLE_CONFIG = 
"internal.topics.automatic.creation.enable";
+    private static final String INTERNAL_TOPICS_AUTOMATIC_CREATION_ENABLE_DOC 
= "Whether to automatically create internal topics used by Connect. "
+            + "This includes the offset, config, and status topics, as well as 
connector-specific offset topics "

Review Comment:
   > connector-specific offset topics
   
   Is this included by the KIP? If not, would you mind adding it?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to