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


##########
tools/src/test/java/org/apache/kafka/tools/streams/DeleteStreamsGroupOffsetTest.java:
##########
@@ -57,78 +51,65 @@
 import java.util.Objects;
 import java.util.Optional;
 import java.util.Properties;
-import java.util.Set;
-import java.util.concurrent.ExecutionException;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
 
 import static org.apache.kafka.common.GroupState.EMPTY;
+import static 
org.apache.kafka.coordinator.group.GroupCoordinatorConfig.GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG;
+import static 
org.apache.kafka.coordinator.group.GroupCoordinatorConfig.OFFSETS_TOPIC_PARTITIONS_CONFIG;
+import static 
org.apache.kafka.coordinator.group.GroupCoordinatorConfig.OFFSETS_TOPIC_REPLICATION_FACTOR_CONFIG;
+import static 
org.apache.kafka.coordinator.group.GroupCoordinatorConfig.STREAMS_GROUP_MIN_HEARTBEAT_INTERVAL_MS_CONFIG;
+import static 
org.apache.kafka.coordinator.group.GroupCoordinatorConfig.STREAMS_GROUP_MIN_SESSION_TIMEOUT_MS_CONFIG;
+import static 
org.apache.kafka.coordinator.transaction.TransactionLogConfig.TRANSACTIONS_TOPIC_MIN_ISR_CONFIG;
+import static 
org.apache.kafka.coordinator.transaction.TransactionLogConfig.TRANSACTIONS_TOPIC_PARTITIONS_CONFIG;
+import static 
org.apache.kafka.coordinator.transaction.TransactionLogConfig.TRANSACTIONS_TOPIC_REPLICATION_FACTOR_CONFIG;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
-@Timeout(600)
-@Tag("integration")
+@ClusterTestDefaults(
+    types = {Type.CO_KRAFT},
+    brokers = 2,
+    serverProperties = {
+        @ClusterConfigProperty(key = OFFSETS_TOPIC_PARTITIONS_CONFIG, value = 
"1"),
+        @ClusterConfigProperty(key = OFFSETS_TOPIC_REPLICATION_FACTOR_CONFIG, 
value = "1"),
+        @ClusterConfigProperty(key = TRANSACTIONS_TOPIC_PARTITIONS_CONFIG, 
value = "1"),
+        @ClusterConfigProperty(key = 
TRANSACTIONS_TOPIC_REPLICATION_FACTOR_CONFIG, value = "1"),
+        @ClusterConfigProperty(key = TRANSACTIONS_TOPIC_MIN_ISR_CONFIG, value 
= "1"),
+        @ClusterConfigProperty(key = GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG, 
value = "0"),
+        @ClusterConfigProperty(key = 
STREAMS_GROUP_MIN_SESSION_TIMEOUT_MS_CONFIG, value = "100"),
+        @ClusterConfigProperty(key = 
STREAMS_GROUP_MIN_HEARTBEAT_INTERVAL_MS_CONFIG, value = "100"),
+    }
+)
 public class DeleteStreamsGroupOffsetTest {
     private static final String TOPIC_PREFIX = "foo-";
     private static final String APP_ID_PREFIX = "streams-group-command-test";
 
     private static final int RECORD_TOTAL = 5;
-    public static EmbeddedKafkaCluster cluster;
-    private static String bootstrapServers;
     private static final String OUTPUT_TOPIC_PREFIX = "output-topic-";
 
-    @BeforeAll
-    public static void startCluster() {
-        final Properties props = new Properties();
-        cluster = new EmbeddedKafkaCluster(2, props);
-        cluster.start();
-
-        bootstrapServers = cluster.bootstrapServers();
-    }
-
-    @AfterEach
-    public void deleteTopicsAndGroups() {
-        try (final Admin adminClient = cluster.createAdminClient()) {
-            // delete all topics
-            final Set<String> topics = adminClient.listTopics().names().get();
-            adminClient.deleteTopics(topics).all().get();
-            // delete all groups
-            List<String> groupIds =
-                
adminClient.listGroups(ListGroupsOptions.forStreamsGroups().timeoutMs(1000)).all().get()
-                    .stream().map(GroupListing::groupId).toList();
-            adminClient.deleteStreamsGroups(groupIds).all().get();
-        } catch (final UnknownTopicOrPartitionException ignored) {
-        } catch (final ExecutionException | InterruptedException e) {
-            if (!(e.getCause() instanceof UnknownTopicOrPartitionException)) {
-                throw new RuntimeException(e);
-            }
-        }
-    }
-
-    private Properties createStreamsConfig(String bootstrapServers, String 
appId) {
+    private static Properties createStreamsConfig(String bootstrapServers, 
String appId) {
         final Properties configs = new Properties();
         configs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
         configs.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
         configs.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, 
Serdes.StringSerde.class);
         configs.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, 
Serdes.StringSerde.class);
         configs.put(StreamsConfig.GROUP_PROTOCOL_CONFIG, 
GroupProtocol.STREAMS.name().toLowerCase(Locale.getDefault()));
         configs.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, 
StreamsConfig.EXACTLY_ONCE_V2);
+        // Use a unique state directory per instance. TestUtils.randomString 
uses a fixed seed, so app ids are

Review Comment:
   spot on!



-- 
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