Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

package org.apache.kafka.connect.mirror.clients.admin;

import org.apache.kafka.clients.admin.AlterConfigOp;
import org.apache.kafka.clients.admin.AlterConfigsOptions;
import org.apache.kafka.clients.admin.AlterConfigsResult;
import org.apache.kafka.clients.admin.Config;
Expand All @@ -38,6 +39,7 @@

import java.util.Collection;
import java.util.Map;
import java.util.stream.Collectors;

/** Customised ForwardingAdmin for testing only.
* The class create/alter topics, partitions and ACLs in Kafka then store metadata in {@link FakeLocalMetadataStore}.
Expand Down Expand Up @@ -79,6 +81,24 @@ public CreatePartitionsResult createPartitions(Map<String, NewPartitions> newPar
return createPartitionsResult;
}

@Override
public AlterConfigsResult incrementalAlterConfigs(Map<ConfigResource, Collection<AlterConfigOp>> configs, AlterConfigsOptions options) {
AlterConfigsResult alterConfigsResult = super.incrementalAlterConfigs(configs, options);
configs.forEach((configResource, newConfigs) -> alterConfigsResult.values().get(configResource).whenComplete((ignored, error) -> {
if (error == null) {
if (configResource.type() == ConfigResource.Type.TOPIC) {
FakeLocalMetadataStore.updateTopicConfig(
configResource.name(),
new Config(newConfigs.stream().map(AlterConfigOp::configEntry).collect(Collectors.toList()))
);
}
} else {
log.error("Unable to intercept admin client operation", error);
}
}));
return alterConfigsResult;
}

@Deprecated
@Override
public AlterConfigsResult alterConfigs(Map<ConfigResource, Config> configs, AlterConfigsOptions options) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,10 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.Vector;
import java.util.concurrent.ConcurrentHashMap;

Expand All @@ -35,17 +37,17 @@
public class FakeLocalMetadataStore {
private static final Logger log = LoggerFactory.getLogger(FakeLocalMetadataStore.class);

private static final ConcurrentHashMap<String, ConcurrentHashMap<String, String>> ALL_TOPICS = new ConcurrentHashMap<>();
private static final Set<String> ALL_TOPICS = Collections.newSetFromMap(new ConcurrentHashMap<>());
private static final ConcurrentHashMap<String, ConcurrentHashMap<String, String>> ALL_TOPIC_CONFIGS = new ConcurrentHashMap<>();
private static final ConcurrentHashMap<String, String> ALL_PARTITIONS = new ConcurrentHashMap<>();
private static final ConcurrentHashMap<String, Vector<AclBinding>> ALL_ACLS = new ConcurrentHashMap<>();

/**
* Add topic to allTopics.
* @param newTopic {@link NewTopic}
*/
public static void addTopicToLocalMetadataStore(NewTopic newTopic) {
ConcurrentHashMap<String, String> configs = new ConcurrentHashMap<>(newTopic.configs());
configs.putIfAbsent("partitions", String.valueOf(newTopic.numPartitions()));
ALL_TOPICS.putIfAbsent(newTopic.name(), configs);
ALL_TOPICS.add(newTopic.name());
}

/**
Expand All @@ -54,9 +56,7 @@ public static void addTopicToLocalMetadataStore(NewTopic newTopic) {
* @param newPartitionCount new partition count.
*/
public static void updatePartitionCount(String topic, int newPartitionCount) {
ConcurrentHashMap<String, String> configs = FakeLocalMetadataStore.ALL_TOPICS.getOrDefault(topic, new ConcurrentHashMap<>());
configs.compute("partitions", (key, value) -> String.valueOf(newPartitionCount));
FakeLocalMetadataStore.ALL_TOPICS.putIfAbsent(topic, configs);
FakeLocalMetadataStore.ALL_PARTITIONS.compute(topic, (key, value) -> String.valueOf(newPartitionCount));
}

/**
Expand All @@ -65,7 +65,7 @@ public static void updatePartitionCount(String topic, int newPartitionCount) {
* @param newConfig topic config
*/
public static void updateTopicConfig(String topic, Config newConfig) {
ConcurrentHashMap<String, String> topicConfigs = FakeLocalMetadataStore.ALL_TOPICS.getOrDefault(topic, new ConcurrentHashMap<>());
ConcurrentHashMap<String, String> topicConfigs = FakeLocalMetadataStore.ALL_TOPIC_CONFIGS.getOrDefault(topic, new ConcurrentHashMap<>());
newConfig.entries().stream().forEach(configEntry -> {
if (configEntry.name() != null) {
if (configEntry.value() != null) {
Expand All @@ -76,7 +76,7 @@ public static void updateTopicConfig(String topic, Config newConfig) {
}
}
});
FakeLocalMetadataStore.ALL_TOPICS.putIfAbsent(topic, topicConfigs);
FakeLocalMetadataStore.ALL_TOPIC_CONFIGS.putIfAbsent(topic, topicConfigs);
}

/**
Expand All @@ -85,7 +85,7 @@ public static void updateTopicConfig(String topic, Config newConfig) {
* @return true if topic name is a key in allTopics
*/
public static Boolean containsTopic(String topic) {
return ALL_TOPICS.containsKey(topic);
return ALL_TOPICS.contains(topic);
}

/**
Expand All @@ -94,7 +94,11 @@ public static Boolean containsTopic(String topic) {
* @return topic configurations.
*/
public static Map<String, String> topicConfig(String topic) {
return ALL_TOPICS.getOrDefault(topic, new ConcurrentHashMap<>());
return ALL_TOPIC_CONFIGS.getOrDefault(topic, new ConcurrentHashMap<>());
}

public static String partitions(String topic) {
return ALL_PARTITIONS.get(topic);
}

/**
Expand Down Expand Up @@ -122,6 +126,8 @@ public static void addACLs(String principal, AclBinding aclBinding) {
*/
public static void clear() {
ALL_TOPICS.clear();
ALL_TOPIC_CONFIGS.clear();
ALL_PARTITIONS.clear();
ALL_ACLS.clear();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,14 +17,16 @@
package org.apache.kafka.connect.mirror.integration;

import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.AlterConfigOp;
import org.apache.kafka.clients.admin.ConfigEntry;
import org.apache.kafka.clients.admin.NewPartitions;
import org.apache.kafka.common.acl.AccessControlEntry;
import org.apache.kafka.common.acl.AccessControlEntryFilter;
import org.apache.kafka.common.acl.AclBinding;
import org.apache.kafka.common.acl.AclBindingFilter;
import org.apache.kafka.common.acl.AclOperation;
import org.apache.kafka.common.acl.AclPermissionType;
import org.apache.kafka.common.config.TopicConfig;
import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.config.internals.BrokerSecurityConfigs;
import org.apache.kafka.common.resource.PatternType;
import org.apache.kafka.common.resource.ResourcePattern;
Expand Down Expand Up @@ -58,7 +60,6 @@
import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX;
import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLIENT_PRODUCER_OVERRIDES_PREFIX;
import static org.apache.kafka.test.TestUtils.waitForCondition;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;

/**
Expand Down Expand Up @@ -153,7 +154,7 @@ public void startClusters() throws Exception {
additionalBackupClusterClientsConfigs.putAll(superUserConfig());
backupWorkerProps.putAll(superUserConfig());

HashMap<String, String> additionalConfig = new HashMap<String, String>(superUserConfig()) {{
Map<String, String> additionalConfig = new HashMap<String, String>(superUserConfig()) {{
put(FORWARDING_ADMIN_CLASS, FakeForwardingAdminWithLocalMetadata.class.getName());
}};

Expand Down Expand Up @@ -199,6 +200,9 @@ public void shutdownClusters() throws Exception {

@Test
public void testReplicationIsCreatingTopicsUsingProvidedForwardingAdmin() throws Exception {
// Disable topic refreshing since org.apache.kafka.connect.mirror.MirrorSourceConnector.refreshTopicPartitions can replicate the topics
// instead of org.apache.kafka.connect.mirror.MirrorSourceConnector.computeAndCreateTopicPartitions that we are checking in this test
mm2Props.put("refresh.topics.enabled", "false");
produceMessages(primaryProducer, "test-topic-1");
produceMessages(backupProducer, "test-topic-1");
String consumerGroupName = "consumer-group-testReplication";
Expand Down Expand Up @@ -235,6 +239,7 @@ public void testReplicationIsCreatingTopicsUsingProvidedForwardingAdmin() throws

@Test
public void testCreatePartitionsUseProvidedForwardingAdmin() throws Exception {
mm2Props.put("refresh.topics.enabled", "true");
mm2Config = new MirrorMakerConfig(mm2Props);
produceMessages(backupProducer, "test-topic-1");
produceMessages(primaryProducer, "test-topic-1");
Expand Down Expand Up @@ -264,7 +269,7 @@ public void testCreatePartitionsUseProvidedForwardingAdmin() throws Exception {
waitForTopicPartitionCreated(backup, "primary.test-topic-1", NUM_PARTITIONS + 1);

// expect to use FakeForwardingAdminWithLocalMetadata to update number of partitions in local store
waitForTopicConfigPersistInFakeLocalMetaDataStore("primary.test-topic-1", "partitions", String.valueOf(NUM_PARTITIONS + 1));
waitForPartitionsPersistInFakeLocalMetaDataStore("primary.test-topic-1", String.valueOf(NUM_PARTITIONS + 1));
}

@Test
Expand All @@ -285,16 +290,21 @@ public void testSyncTopicConfigUseProvidedForwardingAdmin() throws Exception {
waitForTopicCreated(primary, "backup.test-topic-1");
waitForTopicCreated(backup, "primary.test-topic-1");

// make sure the topic config is synced into the other cluster
assertEquals(TopicConfig.CLEANUP_POLICY_COMPACT, getTopicConfig(backup.kafka(), "primary.test-topic-1", TopicConfig.CLEANUP_POLICY_CONFIG),
"topic config was synced");

// expect to use FakeForwardingAdminWithLocalMetadata to create remote topics into local store
waitForTopicToPersistInFakeLocalMetadataStore("backup.test-topic-1");
waitForTopicToPersistInFakeLocalMetadataStore("primary.test-topic-1");

// expect to use FakeForwardingAdminWithLocalMetadata to update topic config in local store
waitForTopicConfigPersistInFakeLocalMetaDataStore("primary.test-topic-1", TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT);
ConfigResource backupConfigResource = new ConfigResource(ConfigResource.Type.TOPIC, "primary.test-topic-1");
Collection<AlterConfigOp> backupRetentionBytesOp = Collections.singletonList(new AlterConfigOp(new ConfigEntry("retention.bytes", "2000"),
AlterConfigOp.OpType.SET
));
Map<ConfigResource, Collection<AlterConfigOp>> backupConfigOps = Collections.singletonMap(backupConfigResource, backupRetentionBytesOp);
backup.kafka().incrementalAlterConfigs(backupConfigOps);

ConfigResource primaryConfigResource = new ConfigResource(ConfigResource.Type.TOPIC, "test-topic-1");
Collection<AlterConfigOp> primaryRetentionBytesOp = Collections.singletonList(new AlterConfigOp(new ConfigEntry("retention.bytes", "3000"),
AlterConfigOp.OpType.SET
));
Map<ConfigResource, Collection<AlterConfigOp>> primaryConfigOps = Collections.singletonMap(primaryConfigResource, primaryRetentionBytesOp);
primary.kafka().incrementalAlterConfigs(primaryConfigOps);

waitForTopicConfigPersistInFakeLocalMetaDataStore("primary.test-topic-1", "retention.bytes", "3000");
}

@Test
Expand Down Expand Up @@ -358,4 +368,10 @@ void waitForTopicConfigPersistInFakeLocalMetaDataStore(String topicName, String
"Topic: " + topicName + "'s configs don't have " + configName + ":" + expectedConfigValue
);
}

void waitForPartitionsPersistInFakeLocalMetaDataStore(String topicName, String expectedConfigValue) throws InterruptedException {
waitForCondition(() -> expectedConfigValue.equals(FakeLocalMetadataStore.partitions(topicName)), FAKE_LOCAL_METADATA_STORE_SYNC_DURATION_MS,
"Topic: " + topicName + " doesn't have partitions:" + expectedConfigValue
);
}
}