This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new ddf1430ffa7 MINOR: Add `TieredStorageTestPlan` as the execution
wrapper created by `TieredStorageTestBuilder.build()` (#22371)
ddf1430ffa7 is described below
commit ddf1430ffa7f3e656a6283d4aa6a2fa2ec6cb81c
Author: Eric Chang <[email protected]>
AuthorDate: Fri Jun 12 19:31:48 2026 +0800
MINOR: Add `TieredStorageTestPlan` as the execution wrapper created by
`TieredStorageTestBuilder.build()` (#22371)
### Summary
MINOR: Add `TieredStorageTestPlan` as the execution wrapper created by
`TieredStorageTestBuilder.build()`.
Migrated tiered storage tests now execute with:
```java
builder.build().execute(clusterInstance, groupProtocol);
```
### Validation
```bash
./gradlew :storage:compileTestJava :storage:checkstyleTest
./gradlew :storage:test --tests
org.apache.kafka.tiered.storage.integration.AlterLogDirTest --tests
org.apache.kafka.tiered.storage.integration.DeleteSegmentsDueToLogStartOffsetBreachTest
--tests org.apache.kafka.tiered.storage.integration.DeleteSegmentsTest
--tests org.apache.kafka.tiered.storage.integration.DeleteTopicTest
--tests
org.apache.kafka.tiered.storage.integration.DisableRemoteLogOnTopicTest
--tests
org.apache.kafka.tiered.storage.integration.EnableRemoteLogOnTopicTest
--tests
org.apache.kafka.tiered.storage.integration.FetchFromLeaderWithCorruptedCheckpointTest
--tests org.apache.kafka.tiered.storage.integration.ListOffsetsTest
--tests
org.apache.kafka.tiered.storage.integration.OffloadAndConsumeFromLeaderTest
--tests
org.apache.kafka.tiered.storage.integration.OffloadAndTxnConsumeFromLeaderTest
--tests org.apache.kafka.tiered.storage.integration.PartitionsExpandTest
--tests
org.apache.kafka.tiered.storage.integration.ReassignReplicaMoveAndExpandTest
--tests
org.apache.kafka.tiered.storage.integration.ReassignReplicaShrinkTest
--tests
org.apache.kafka.tiered.storage.integration.RollAndOffloadActiveSegmentTest
```
Reviewers: TaiJuWu <[email protected]>, Chia-Ping Tsai
<[email protected]>, PoAn Yang <[email protected]>, Ken Huang
<[email protected]>
---
.../tiered/storage/TieredStorageTestBuilder.java | 39 ++++++++++++++++++++--
.../storage/integration/AlterLogDirTest.java | 17 +---------
...eleteSegmentsDueToLogStartOffsetBreachTest.java | 16 +--------
.../storage/integration/DeleteSegmentsTest.java | 16 +--------
.../storage/integration/DeleteTopicTest.java | 16 +--------
.../integration/DisableRemoteLogOnTopicTest.java | 16 +--------
.../integration/EnableRemoteLogOnTopicTest.java | 16 +--------
...FetchFromLeaderWithCorruptedCheckpointTest.java | 16 +--------
.../storage/integration/ListOffsetsTest.java | 16 +--------
.../OffloadAndConsumeFromLeaderTest.java | 16 +--------
.../OffloadAndTxnConsumeFromLeaderTest.java | 13 ++------
.../storage/integration/PartitionsExpandTest.java | 16 +--------
.../ReassignReplicaMoveAndExpandTest.java | 17 +---------
.../integration/ReassignReplicaShrinkTest.java | 16 +--------
.../RollAndOffloadActiveSegmentTest.java | 16 +--------
15 files changed, 52 insertions(+), 210 deletions(-)
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/TieredStorageTestBuilder.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/TieredStorageTestBuilder.java
index db9fd4f9b9b..2334f3d1c3c 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/TieredStorageTestBuilder.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/TieredStorageTestBuilder.java
@@ -17,9 +17,12 @@
package org.apache.kafka.tiered.storage;
import org.apache.kafka.clients.admin.OffsetSpec;
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.config.TopicConfig;
+import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.server.log.remote.storage.LocalTieredStorageEvent;
import org.apache.kafka.storage.internals.log.EpochEntry;
import org.apache.kafka.tiered.storage.actions.AlterLogDirAction;
@@ -58,6 +61,7 @@ import java.io.FilenameFilter;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -314,8 +318,39 @@ public final class TieredStorageTestBuilder {
return this;
}
- public List<TieredStorageTestAction> complete() {
- return actions;
+ /**
+ * Builds an executable test plan from the actions described so far.
+ */
+ public TieredStorageTestPlan build() {
+ return new TieredStorageTestPlan(actions);
+ }
+
+ public static final class TieredStorageTestPlan {
+
+ private final List<TieredStorageTestAction> actions;
+
+ private TieredStorageTestPlan(List<TieredStorageTestAction> actions) {
+ this.actions = List.copyOf(actions);
+ }
+
+ public void execute(ClusterInstance clusterInstance, GroupProtocol
groupProtocol) throws Exception {
+ execute(clusterInstance, Map.of(
+ ConsumerConfig.GROUP_PROTOCOL_CONFIG,
+ groupProtocol.name().toLowerCase(Locale.ROOT)
+ ));
+ }
+
+ public void execute(ClusterInstance clusterInstance, Map<String,
Object> extraConsumerProps) throws Exception {
+ try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
+ try {
+ for (TieredStorageTestAction action : actions) {
+ action.execute(context);
+ }
+ } finally {
+ context.printReport(System.out);
+ }
+ }
+ }
}
private void createProduceAction() {
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/AlterLogDirTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/AlterLogDirTest.java
index d406d8e2eba..1a4417088fc 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/AlterLogDirTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/AlterLogDirTest.java
@@ -16,20 +16,16 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
import java.util.Locale;
-import java.util.Map;
import java.util.Set;
import static org.apache.kafka.common.utils.Utils.mkEntry;
@@ -94,17 +90,6 @@ public final class AlterLogDirTest {
.expectFetchFromTieredStorage(broker0, topicB, p0, 3)
.consume(topicB, p0, 0L, 4, 3);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsDueToLogStartOffsetBreachTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsDueToLogStartOffsetBreachTest.java
index 4fc69da6c29..16585c733b5 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsDueToLogStartOffsetBreachTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsDueToLogStartOffsetBreachTest.java
@@ -16,15 +16,12 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -118,17 +115,6 @@ public final class
DeleteSegmentsDueToLogStartOffsetBreachTest {
.expectFetchFromTieredStorage(broker1, topicA, p0, 1)
.consume(topicA, p0, 7L, 3, 1);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsTest.java
index e416e6be458..e0ad94af782 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteSegmentsTest.java
@@ -16,16 +16,13 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.config.TopicConfig;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -121,17 +118,6 @@ public final class DeleteSegmentsTest {
.expectFetchFromTieredStorage(broker0, topicA, p0, 0)
.consume(topicA, p0, 0L, 1, 0);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteTopicTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteTopicTest.java
index a7c58790f70..839d8330e3c 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteTopicTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DeleteTopicTest.java
@@ -16,15 +16,12 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -102,17 +99,6 @@ public final class DeleteTopicTest {
.expectEmptyRemoteStorage(topicA, p0)
.expectEmptyRemoteStorage(topicA, p1);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
\ No newline at end of file
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DisableRemoteLogOnTopicTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DisableRemoteLogOnTopicTest.java
index 88963085d45..e3630fecb76 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DisableRemoteLogOnTopicTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/DisableRemoteLogOnTopicTest.java
@@ -16,16 +16,13 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.config.TopicConfig;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.HashMap;
@@ -144,17 +141,6 @@ public final class DisableRemoteLogOnTopicTest {
// verify the local log is still consumable
.consume(topicA, p0, 3L, 1, 0);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/EnableRemoteLogOnTopicTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/EnableRemoteLogOnTopicTest.java
index 598480216b8..eed654ea776 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/EnableRemoteLogOnTopicTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/EnableRemoteLogOnTopicTest.java
@@ -16,16 +16,13 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.config.TopicConfig;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -114,17 +111,6 @@ public final class EnableRemoteLogOnTopicTest {
.expectFetchFromTieredStorage(broker1, topicA, p1, 4)
.consume(topicA, p1, 0L, 5, 4);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/FetchFromLeaderWithCorruptedCheckpointTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/FetchFromLeaderWithCorruptedCheckpointTest.java
index 07af185b2e3..4a1e6db412a 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/FetchFromLeaderWithCorruptedCheckpointTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/FetchFromLeaderWithCorruptedCheckpointTest.java
@@ -18,7 +18,6 @@ package org.apache.kafka.tiered.storage.integration;
import kafka.server.ReplicaManager;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
@@ -26,9 +25,7 @@ import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
import org.apache.kafka.storage.internals.checkpoint.CleanShutdownFileHandler;
import org.apache.kafka.storage.internals.log.LogManager;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -109,17 +106,6 @@ public final class
FetchFromLeaderWithCorruptedCheckpointTest {
.expectFetchFromTieredStorage(broker0, topicA, p0, 4)
.consume(topicA, p0, 0L, 5, 4);
- final Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ListOffsetsTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ListOffsetsTest.java
index c2ca43abb74..3662740f062 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ListOffsetsTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ListOffsetsTest.java
@@ -17,7 +17,6 @@
package org.apache.kafka.tiered.storage.integration;
import org.apache.kafka.clients.admin.OffsetSpec;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
@@ -26,9 +25,7 @@ import org.apache.kafka.common.test.api.Type;
import org.apache.kafka.common.utils.MockTime;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.storage.internals.log.EpochEntry;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -137,17 +134,6 @@ public final class ListOffsetsTest {
.expectListOffsets(topicA, p0, OffsetSpec.latestTiered(), new
EpochEntry(NO_PARTITION_LEADER_EPOCH, 3))
.expectListOffsets(topicA, p0, OffsetSpec.latest(), new
EpochEntry(1, 6));
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndConsumeFromLeaderTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndConsumeFromLeaderTest.java
index 539cd008168..5a41324aaf8 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndConsumeFromLeaderTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndConsumeFromLeaderTest.java
@@ -16,15 +16,12 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -160,17 +157,6 @@ public final class OffloadAndConsumeFromLeaderTest {
.expectFetchFromTieredStorage(broker0, topicB, p0, 2)
.consume(topicB, p0, 1L, 4, 3);
- final Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndTxnConsumeFromLeaderTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndTxnConsumeFromLeaderTest.java
index b5f99785e0f..38334e1cfb8 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndTxnConsumeFromLeaderTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/OffloadAndTxnConsumeFromLeaderTest.java
@@ -24,9 +24,7 @@ import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
import org.apache.kafka.server.log.remote.storage.RemoteLogManagerConfig;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.FetchCountAndOp;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import org.apache.kafka.tiered.storage.specs.RemoteFetchCount;
@@ -111,15 +109,8 @@ public final class OffloadAndTxnConsumeFromLeaderTest {
ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT),
ConsumerConfig.ISOLATION_LEVEL_CONFIG,
IsolationLevel.READ_COMMITTED.toString()
);
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+
+ builder.build().execute(clusterInstance, extraConsumerProps);
}
private static RemoteFetchCount getRemoteFetchCount() {
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/PartitionsExpandTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/PartitionsExpandTest.java
index f7e1ae2839d..2bac0d5d7c8 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/PartitionsExpandTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/PartitionsExpandTest.java
@@ -16,15 +16,12 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -131,17 +128,6 @@ public final class PartitionsExpandTest {
.expectFetchFromTieredStorage(broker1, topicA, p2, 1)
.consume(topicA, p2, 1L, 2, 1);
- final Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaMoveAndExpandTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaMoveAndExpandTest.java
index 3c48f3cdc09..38466e1b621 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaMoveAndExpandTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaMoveAndExpandTest.java
@@ -16,21 +16,17 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.ArrayList;
import java.util.List;
import java.util.Locale;
-import java.util.Map;
import java.util.Set;
import static org.apache.kafka.common.utils.Utils.mkEntry;
@@ -117,17 +113,6 @@ public final class ReassignReplicaMoveAndExpandTest {
.expectFetchFromTieredStorage(BROKER_1, topicB, p0, 3)
.consume(topicB, p0, 0L, 4, 3);
- Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaShrinkTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaShrinkTest.java
index 10604a92b46..3be02ef7106 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaShrinkTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/ReassignReplicaShrinkTest.java
@@ -16,15 +16,12 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.List;
@@ -116,17 +113,6 @@ public final class ReassignReplicaShrinkTest {
.expectFetchFromTieredStorage(broker0, topicA, p1, 3)
.consume(topicA, p1, 0L, 4, 3);
- final Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
}
diff --git
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/RollAndOffloadActiveSegmentTest.java
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/RollAndOffloadActiveSegmentTest.java
index e960c234760..43d57c7a000 100644
---
a/storage/src/test/java/org/apache/kafka/tiered/storage/integration/RollAndOffloadActiveSegmentTest.java
+++
b/storage/src/test/java/org/apache/kafka/tiered/storage/integration/RollAndOffloadActiveSegmentTest.java
@@ -16,16 +16,13 @@
*/
package org.apache.kafka.tiered.storage.integration;
-import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.common.config.TopicConfig;
import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfig;
import org.apache.kafka.common.test.api.ClusterTemplate;
import org.apache.kafka.common.test.api.Type;
-import org.apache.kafka.tiered.storage.TieredStorageTestAction;
import org.apache.kafka.tiered.storage.TieredStorageTestBuilder;
-import org.apache.kafka.tiered.storage.TieredStorageTestContext;
import org.apache.kafka.tiered.storage.specs.KeyValueSpec;
import java.util.HashMap;
@@ -96,18 +93,7 @@ public final class RollAndOffloadActiveSegmentTest {
.expectFetchFromTieredStorage(broker0, topicA, p0, 4)
.consume(topicA, p0, 0L, 4, 4);
- final Map<String, Object> extraConsumerProps = Map.of(
- ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol.name().toLowerCase(Locale.ROOT)
- );
- try (TieredStorageTestContext context = new
TieredStorageTestContext(clusterInstance, extraConsumerProps)) {
- try {
- for (TieredStorageTestAction action : builder.complete()) {
- action.execute(context);
- }
- } finally {
- context.printReport(System.out);
- }
- }
+ builder.build().execute(clusterInstance, groupProtocol);
}
private static Map<String, String> configsToBeAdded() {