MINOR: Use ClusterInstance#createTopic in various integration tests - #23513
Conversation
|
Thanks for the cleanup! The follow-up in apache/kafka#22874 (comment) was meant to clean up the whole codebase, and there are still many places using the same pattern, such as |
Thank you for the review! Sorry for missing other files in the first place. I have added them in the new commits. Please let me know if I missed any. |
|
A label of 'needs-attention' was automatically added to this PR in order to raise the |
5d03882 to
1c473c1
Compare
m1a2st
left a comment
There was a problem hiding this comment.
Could you also take a look at the following tests? They follow the same pattern:
DeleteTopicTestConnectInternalTopicsTestBootstrapControllersIntegrationTest
There are also several tests that create a topic using only createTopics(...).all().get() (or .topicId(...).get()) without an explicit wait:
ConfigCommandIntegrationTestDescribeConsumerGroupTestDeleteOffsetsConsumerGroupCommandIntegrationTestLeaderElectionCommandTestResetConsumerGroupOffsetTestDeleteRecordsCommandTestConsumerIntegrationTestShareConsumerDLQTestCordonedLogDirsIntegrationTest
1c473c1 to
cf7a92d
Compare
|
@m1a2st Thanks for the review! I've converted the tests you listed, plus a few more with the same pattern, and updated the PR description with the criteria I used. Also few still create topics manually because they need a replica assignment together with configs. Would it make sense to add a createTopicWithAssignment(name, assignment, configs) overload? |
m1a2st
left a comment
There was a problem hiding this comment.
Thanks @lytt925 for the update.
Also few still create topics manually because they need a replica assignment together with configs. Would it make sense to add a createTopicWithAssignment(name, assignment, configs) overload?
Yes, I think adding an overload makes sense. It would allow us to convert the remaining call sites that require both a replica assignment and topic configs.
| @@ -219,15 +219,15 @@ public void testStaticCordonUncordonLogDirs() throws Exception { | |||
| setCordonedLogDirs(admin, List.of(), BROKER_0); | |||
|
|
|||
| // We can't create topics again | |||
There was a problem hiding this comment.
This comment seems incorrect. Could you update it?
| default void createTopicWithAssignment(String topicName, Map<Integer, List<Integer>> replicaAssignment, Map<String, String> props) throws InterruptedException { | ||
| try (Admin admin = admin()) { | ||
| admin.createTopics(List.of(new NewTopic(topicName, replicaAssignment))); | ||
| admin.createTopics(List.of(new NewTopic(topicName, replicaAssignment).configs(props))); |
There was a problem hiding this comment.
We need to call get to have the correct error propagated.
| void testEnableRemoteLogOnExistingTopic() throws Exception { | ||
| try (var admin = cluster.admin()) { | ||
| admin.createTopics(List.of(new NewTopic(testTopicName, numPartitions, numReplicationFactor).configs(Map.of()))).all().get(); | ||
| cluster.createTopic(testTopicName, numPartitions, numReplicationFactor, Map.of()); |
There was a problem hiding this comment.
cluster.createTopic(testTopicName, numPartitions, numReplicationFactor)
| } | ||
| String topic1 = "kafka.testTopic1"; | ||
| String hiddenConsumerTopic = Topic.GROUP_METADATA_TOPIC_NAME; | ||
| int partition = 2; |
There was a problem hiding this comment.
Those local variables were useful before, but now they look a bit verbose, right?
| @ClusterTest(brokers = 3) | ||
| public void testAlterPartitionCount(ClusterInstance clusterInstance) throws Exception { | ||
| String testTopicName = TestUtils.randomString(10); | ||
| int partition = 2; |
| @ClusterTemplate("generate") | ||
| public void testAlterAssignment(ClusterInstance clusterInstance) throws Exception { | ||
| String testTopicName = TestUtils.randomString(10); | ||
| int partition = 2; |
| @ClusterTest(brokers = 3) | ||
| public void testAlterAssignmentWithMoreAssignmentThanPartitions(ClusterInstance clusterInstance) throws Exception { | ||
| String testTopicName = TestUtils.randomString(10); | ||
| int partition = 2; |
| @ClusterTemplate("generate") | ||
| public void testAlterAssignmentWithMorePartitionsThanAssignment(ClusterInstance clusterInstance) throws Exception { | ||
| String testTopicName = TestUtils.randomString(10); | ||
| int partition = 2; |
| @@ -645,7 +640,7 @@ private void updateAndCheckInvalidBrokerConfig(Optional<String> brokerIdOrDefaul | |||
| public void testUpdateInvalidTopicConfigs() throws ExecutionException, InterruptedException { | |||
There was a problem hiding this comment.
We can scrape ExecutionException out
| @@ -218,9 +217,6 @@ private String execute(LogDirsCommand.LogDirsCommandOptions options, Admin admin | |||
| } | |||
|
|
|||
| private void createTopic(ClusterInstance clusterInstance, String topic) { | |||
There was a problem hiding this comment.
We could use clusterInstance.createTopicWithAssignment instead of this helper, right?
| @ClusterTest(brokers = 3) | ||
| public void testCreateWithTopicNameCollision(ClusterInstance clusterInstance) throws Exception { | ||
| String topic = "foo_bar"; | ||
| int partitions = 1; |
Follow-up to apache/kafka#22874
(comment)
Currently many tests running on
ClusterInstancecreate topics manuallyusing
Admin#createTopics. The post-creation waiting behavior is notinconsistent. Some tests implement custom waits, some call
.all().get(), and others don't wait.Changes
This PR switch these calls to
ClusterInstance#createTopicorClusterInstance#createTopicWithAssignment.when the topic is for test setup and maps to one of the
helpers.
Also add a new overload of createTopicWithAssignment to convert more
tests.
I kept the original createTopics calls where the test needs the returned
topic ID
Reviewers: Ming-Yen Chung mingyen066@gmail.com
Reviewers: Ken Huang s7133700@gmail.com