This is an automated email from the ASF dual-hosted git repository.

lzljs3620320 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink-table-store.git


The following commit(s) were added to refs/heads/master by this push:
     new 74cc7ce1 [FLINK-30997] Refactor tests in connector to extends 
AbstractTestBase
74cc7ce1 is described below

commit 74cc7ce1e516e96349021fa7a4c96f2e992d1b28
Author: Shammon FY <[email protected]>
AuthorDate: Tue Feb 14 09:30:01 2023 +0800

    [FLINK-30997] Refactor tests in connector to extends AbstractTestBase
    
    This closes #512
---
 .../store/connector/AppendOnlyTableITCase.java     |   2 +-
 .../store/connector/BatchFileStoreITCase.java      |   2 +-
 .../table/store/connector/CatalogITCaseBase.java   |  17 +-
 .../table/store/connector/CatalogTableITCase.java  |   7 +-
 ...AndMultiPartitionedTableWIthKafkaLogITCase.java |   6 +-
 .../ComputedColumnAndWatermarkTableITCase.java     |   6 +-
 .../store/connector/ContinuousFileStoreITCase.java |  35 ++--
 .../table/store/connector/FileStoreITCase.java     |  50 ++---
 .../store/connector/FileSystemCatalogITCase.java   |  16 +-
 .../store/connector/ForceCompactionITCase.java     |   2 +-
 .../connector/FullCompactionFileStoreITCase.java   |   6 +-
 .../table/store/connector/LargeDataITCase.java     |   2 +-
 .../table/store/connector/LogSystemITCase.java     |   8 +-
 .../table/store/connector/LookupJoinITCase.java    |  10 +-
 .../table/store/connector/MappingTableITCase.java  |  10 +-
 .../table/store/connector/PartialUpdateITCase.java |   2 +-
 .../store/connector/PreAggregationITCase.java      |   2 +-
 .../table/store/connector/PredicateITCase.java     |   2 +-
 .../table/store/connector/RescaleBucketITCase.java |  28 ++-
 .../table/store/connector/SchemaChangeITCase.java  |   2 +-
 .../StreamingReadWriteTableWithKafkaLogITCase.java |  10 +-
 .../store/connector/StreamingWarehouseITCase.java  |   2 +-
 .../store/connector/sink/CompactorSinkITCase.java  |  27 +--
 .../store/connector/sink/SinkSavepointITCase.java  |  53 +++--
 .../connector/source/CompactorSourceITCase.java    |  10 +-
 .../store/connector/util/AbstractTestBase.java     |  22 +-
 .../util/MiniClusterWithClientExtension.java       | 225 +++++++++++++++++++++
 .../table/store/kafka/KafkaTableTestBase.java      |  71 ++++---
 28 files changed, 449 insertions(+), 186 deletions(-)

diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/AppendOnlyTableITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/AppendOnlyTableITCase.java
index 757934eb..7bc17944 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/AppendOnlyTableITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/AppendOnlyTableITCase.java
@@ -23,7 +23,7 @@ import org.apache.flink.table.store.file.Snapshot;
 import org.apache.flink.types.Row;
 import org.apache.flink.types.RowKind;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.util.Arrays;
 import java.util.Collections;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/BatchFileStoreITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/BatchFileStoreITCase.java
index 7fdd2278..b6ab0591 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/BatchFileStoreITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/BatchFileStoreITCase.java
@@ -21,7 +21,7 @@ package org.apache.flink.table.store.connector;
 import org.apache.flink.table.store.CoreOptions;
 import org.apache.flink.types.Row;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.util.Collections;
 import java.util.List;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogITCaseBase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogITCaseBase.java
index cc916ddb..6b03d548 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogITCaseBase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogITCaseBase.java
@@ -30,22 +30,22 @@ import 
org.apache.flink.table.catalog.exceptions.TableNotExistException;
 import org.apache.flink.table.delegation.Parser;
 import org.apache.flink.table.operations.Operation;
 import org.apache.flink.table.operations.ddl.CreateCatalogOperation;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.table.store.file.Snapshot;
 import org.apache.flink.table.store.file.utils.SnapshotManager;
 import org.apache.flink.table.store.fs.Path;
 import org.apache.flink.table.store.fs.local.LocalFileIO;
-import org.apache.flink.test.util.AbstractTestBase;
 import org.apache.flink.types.Row;
 import org.apache.flink.util.CloseableIterator;
 
 import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableList;
 
-import org.junit.Before;
+import org.junit.jupiter.api.BeforeEach;
 
 import javax.annotation.Nullable;
 
+import java.io.File;
 import java.io.IOException;
-import java.net.URI;
 import java.time.Duration;
 import java.util.Collections;
 import java.util.List;
@@ -59,16 +59,15 @@ public abstract class CatalogITCaseBase extends 
AbstractTestBase {
     protected TableEnvironment sEnv;
     protected String path;
 
-    @Before
+    @BeforeEach
     public void before() throws IOException {
         tEnv = 
TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());
         String catalog = "TABLE_STORE";
-        URI uri = TEMPORARY_FOLDER.newFolder().toURI();
-        path = uri.toString();
+        path = getTempDirPath();
         tEnv.executeSql(
                 String.format(
                         "CREATE CATALOG %s WITH (" + "'type'='table-store', 
'warehouse'='%s')",
-                        catalog, uri));
+                        catalog, path));
         tEnv.useCatalog(catalog);
 
         sEnv = 
TableEnvironment.create(EnvironmentSettings.newInstance().inStreamingMode().build());
@@ -134,7 +133,9 @@ public abstract class CatalogITCaseBase extends 
AbstractTestBase {
     }
 
     protected Path getTableDirectory(String tableName) {
-        return new Path(path + String.format("%s.db/%s", 
tEnv.getCurrentDatabase(), tableName));
+        return new Path(
+                new File(path, String.format("%s.db/%s", 
tEnv.getCurrentDatabase(), tableName))
+                        .toString());
     }
 
     @Nullable
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogTableITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogTableITCase.java
index cc693e2a..30a40ec1 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogTableITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CatalogTableITCase.java
@@ -26,8 +26,9 @@ import org.apache.flink.table.store.types.IntType;
 import org.apache.flink.types.Row;
 
 import org.apache.commons.lang3.StringUtils;
-import org.jetbrains.annotations.NotNull;
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
+
+import javax.annotation.Nonnull;
 
 import java.util.Arrays;
 import java.util.List;
@@ -273,7 +274,7 @@ public class CatalogTableITCase extends CatalogITCaseBase {
                                                 : "[5],[10]")));
     }
 
-    @NotNull
+    @Nonnull
     private List<String> getRowStringList(List<Row> rows) {
         return rows.stream()
                 .map(
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CompositePkAndMultiPartitionedTableWIthKafkaLogITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CompositePkAndMultiPartitionedTableWIthKafkaLogITCase.java
index 0192e4ec..a6d2f4bc 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CompositePkAndMultiPartitionedTableWIthKafkaLogITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/CompositePkAndMultiPartitionedTableWIthKafkaLogITCase.java
@@ -22,8 +22,8 @@ import 
org.apache.flink.table.store.file.utils.BlockingIterator;
 import org.apache.flink.table.store.kafka.KafkaTableTestBase;
 import org.apache.flink.types.Row;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -55,7 +55,7 @@ import static 
org.apache.flink.table.store.connector.util.ReadWriteTableTestUtil
  */
 public class CompositePkAndMultiPartitionedTableWIthKafkaLogITCase extends 
KafkaTableTestBase {
 
-    @Before
+    @BeforeEach
     public void setUp() throws Exception {
         init(createAndRegisterTempFile("").toString());
     }
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ComputedColumnAndWatermarkTableITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ComputedColumnAndWatermarkTableITCase.java
index a72ae55d..f05623e5 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ComputedColumnAndWatermarkTableITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ComputedColumnAndWatermarkTableITCase.java
@@ -21,8 +21,8 @@ package org.apache.flink.table.store.connector;
 import org.apache.flink.table.store.kafka.KafkaTableTestBase;
 import org.apache.flink.types.Row;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.time.LocalDateTime;
 import java.util.Arrays;
@@ -47,7 +47,7 @@ import static 
org.apache.flink.table.store.connector.util.ReadWriteTableTestUtil
 /** Table store IT case when the table has computed column and watermark spec. 
*/
 public class ComputedColumnAndWatermarkTableITCase extends KafkaTableTestBase {
 
-    @Before
+    @BeforeEach
     public void setUp() throws Exception {
         init(createAndRegisterTempFile("").toString());
     }
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ContinuousFileStoreITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ContinuousFileStoreITCase.java
index e5288864..57c5a757 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ContinuousFileStoreITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ContinuousFileStoreITCase.java
@@ -22,13 +22,14 @@ import org.apache.flink.table.store.file.Snapshot;
 import org.apache.flink.table.store.file.utils.BlockingIterator;
 import org.apache.flink.table.store.file.utils.SnapshotManager;
 import org.apache.flink.table.store.fs.local.LocalFileIO;
+import 
org.apache.flink.testutils.junit.extensions.parameterized.ParameterizedTestExtension;
+import org.apache.flink.testutils.junit.extensions.parameterized.Parameters;
 import org.apache.flink.types.Row;
 
 import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableList;
 
-import org.junit.Test;
-import org.junit.runner.RunWith;
-import org.junit.runners.Parameterized;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.extension.ExtendWith;
 
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -41,7 +42,7 @@ import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 /** SQL ITCase for continuous file store. */
-@RunWith(Parameterized.class)
+@ExtendWith(ParameterizedTestExtension.class)
 public class ContinuousFileStoreITCase extends CatalogITCaseBase {
 
     private final boolean changelogFile;
@@ -50,7 +51,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
         this.changelogFile = changelogFile;
     }
 
-    @Parameterized.Parameters(name = "changelogFile-{0}")
+    @Parameters(name = "changelogFile-{0}")
     public static Collection<Boolean> parameters() {
         return Arrays.asList(true, false);
     }
@@ -64,22 +65,22 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
                         + options);
     }
 
-    @Test
+    @TestTemplate
     public void testWithoutPrimaryKey() throws Exception {
         testSimple("T1");
     }
 
-    @Test
+    @TestTemplate
     public void testWithPrimaryKey() throws Exception {
         testSimple("T2");
     }
 
-    @Test
+    @TestTemplate
     public void testProjectionWithoutPrimaryKey() throws Exception {
         testProjection("T1");
     }
 
-    @Test
+    @TestTemplate
     public void testProjectionWithPrimaryKey() throws Exception {
         testProjection("T2");
     }
@@ -108,7 +109,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
         assertThat(iterator.collect(1)).containsExactlyInAnyOrder(Row.of("8", 
"9"));
     }
 
-    @Test
+    @TestTemplate
     public void testContinuousLatest() throws TimeoutException {
         batchSql("INSERT INTO T1 VALUES ('1', '2', '3'), ('4', '5', '6')");
 
@@ -121,7 +122,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
                 .containsExactlyInAnyOrder(Row.of("7", "8", "9"), Row.of("10", 
"11", "12"));
     }
 
-    @Test
+    @TestTemplate
     public void testContinuousFromTimestamp() throws Exception {
         String sql =
                 "SELECT * FROM T1 /*+ OPTIONS('log.scan'='from-timestamp', 
'log.scan.timestamp-millis'='%s') */";
@@ -182,7 +183,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
         iterator.close();
     }
 
-    @Test
+    @TestTemplate
     public void testLackStartupTimestamp() {
         assertThatThrownBy(
                         () ->
@@ -191,7 +192,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
                 .hasMessageContaining("Unable to create a source for reading 
table");
     }
 
-    @Test
+    @TestTemplate
     public void testConfigureStartupTimestamp() throws Exception {
         // Configure 'log.scan.timestamp-millis' without 'log.scan'.
         BlockingIterator<Row, Row> iterator =
@@ -214,7 +215,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
                 .hasMessageContaining("Unable to create a source for reading 
table");
     }
 
-    @Test
+    @TestTemplate
     public void testConfigureStartupSnapshot() throws Exception {
         // Configure 'scan.snapshot-id' without 'scan.mode'.
         BlockingIterator<Row, Row> iterator =
@@ -245,7 +246,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
                 .hasMessageContaining("Unable to create a source for reading 
table");
     }
 
-    @Test
+    @TestTemplate
     public void testIgnoreOverwrite() throws TimeoutException {
         BlockingIterator<Row, Row> iterator =
                 BlockingIterator.of(streamSqlIter("SELECT * FROM T1"));
@@ -261,7 +262,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
         assertThat(iterator.collect(1)).containsExactlyInAnyOrder(Row.of("9", 
"10", "11"));
     }
 
-    @Test
+    @TestTemplate
     public void testUnsupportedUpsert() {
         assertThatThrownBy(
                 () ->
@@ -270,7 +271,7 @@ public class ContinuousFileStoreITCase extends 
CatalogITCaseBase {
                 "File store continuous reading dose not support upsert 
changelog mode");
     }
 
-    @Test
+    @TestTemplate
     public void testUnsupportedEventual() {
         assertThatThrownBy(
                 () ->
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileStoreITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileStoreITCase.java
index 5f93e6d1..10aa4590 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileStoreITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileStoreITCase.java
@@ -36,6 +36,7 @@ import 
org.apache.flink.table.store.connector.sink.FlinkSinkBuilder;
 import org.apache.flink.table.store.connector.source.ContinuousFileStoreSource;
 import org.apache.flink.table.store.connector.source.FlinkSourceBuilder;
 import org.apache.flink.table.store.connector.source.StaticFileStoreSource;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.table.store.file.schema.SchemaManager;
 import org.apache.flink.table.store.file.schema.UpdateSchema;
 import org.apache.flink.table.store.file.utils.BlockingIterator;
@@ -49,18 +50,15 @@ import org.apache.flink.table.types.logical.IntType;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.table.types.logical.VarCharType;
 import org.apache.flink.table.types.utils.TypeConversions;
-import org.apache.flink.test.util.AbstractTestBase;
+import 
org.apache.flink.testutils.junit.extensions.parameterized.ParameterizedTestExtension;
+import org.apache.flink.testutils.junit.extensions.parameterized.Parameters;
 import org.apache.flink.types.Row;
 import org.apache.flink.types.RowKind;
 import org.apache.flink.util.CloseableIterator;
 
-import org.junit.Assume;
-import org.junit.Test;
-import org.junit.rules.TemporaryFolder;
-import org.junit.runner.RunWith;
-import org.junit.runners.Parameterized;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.extension.ExtendWith;
 
-import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -79,12 +77,14 @@ import static org.apache.flink.table.store.CoreOptions.PATH;
 import static 
org.apache.flink.table.store.connector.LogicalTypeConversion.toDataType;
 import static 
org.apache.flink.table.store.file.utils.FailingFileIO.retryArtificialException;
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.junit.jupiter.api.Assumptions.assumeFalse;
+import static org.junit.jupiter.api.Assumptions.assumeTrue;
 
 /**
  * ITCase for {@link StaticFileStoreSource}, {@link ContinuousFileStoreSource} 
and {@link
  * FileStoreSink}.
  */
-@RunWith(Parameterized.class)
+@ExtendWith(ParameterizedTestExtension.class)
 public class FileStoreITCase extends AbstractTestBase {
 
     public static final RowType TABLE_TYPE =
@@ -124,7 +124,7 @@ public class FileStoreITCase extends AbstractTestBase {
         this.env = isBatch ? buildBatchEnv() : buildStreamEnv();
     }
 
-    @Parameterized.Parameters(name = "isBatch-{0}")
+    @Parameters(name = "isBatch-{0}")
     public static List<Boolean> getVarSeg() {
         return Arrays.asList(true, false);
     }
@@ -133,7 +133,7 @@ public class FileStoreITCase extends AbstractTestBase {
         return new SerializableRowData(row, 
InternalSerializers.create(TABLE_TYPE));
     }
 
-    @Test
+    @TestTemplate
     public void testPartitioned() throws Exception {
         FileStoreTable table = buildFileStoreTable(new int[] {1}, new int[] 
{1, 2});
 
@@ -153,7 +153,7 @@ public class FileStoreITCase extends AbstractTestBase {
         assertThat(results).containsExactlyInAnyOrder(expected);
     }
 
-    @Test
+    @TestTemplate
     public void testNonPartitioned() throws Exception {
         FileStoreTable table = buildFileStoreTable(new int[0], new int[] {2});
 
@@ -170,9 +170,9 @@ public class FileStoreITCase extends AbstractTestBase {
         assertThat(results).containsExactlyInAnyOrder(expected);
     }
 
-    @Test
+    @TestTemplate
     public void testOverwrite() throws Exception {
-        Assume.assumeTrue(isBatch);
+        assumeTrue(isBatch);
 
         FileStoreTable table = buildFileStoreTable(new int[] {1}, new int[] 
{1, 2});
 
@@ -219,7 +219,7 @@ public class FileStoreITCase extends AbstractTestBase {
         assertThat(results).containsExactlyInAnyOrder(expected);
     }
 
-    @Test
+    @TestTemplate
     public void testPartitionedNonKey() throws Exception {
         FileStoreTable table = buildFileStoreTable(new int[] {1}, new int[0]);
 
@@ -241,12 +241,12 @@ public class FileStoreITCase extends AbstractTestBase {
         assertThat(results).containsExactlyInAnyOrder(expected);
     }
 
-    @Test
+    @TestTemplate
     public void testKeyedProjection() throws Exception {
         testProjection(buildFileStoreTable(new int[0], new int[] {2}));
     }
 
-    @Test
+    @TestTemplate
     public void testNonKeyedProjection() throws Exception {
         testProjection(buildFileStoreTable(new int[0], new int[0]));
     }
@@ -289,18 +289,18 @@ public class FileStoreITCase extends AbstractTestBase {
         assertThat(results).containsExactlyInAnyOrder(expected);
     }
 
-    @Test
+    @TestTemplate
     public void testContinuous() throws Exception {
         innerTestContinuous(buildFileStoreTable(new int[0], new int[] {2}));
     }
 
-    @Test
+    @TestTemplate
     public void testContinuousWithoutPK() throws Exception {
         innerTestContinuous(buildFileStoreTable(new int[0], new int[0]));
     }
 
     private void innerTestContinuous(FileStoreTable table) throws Exception {
-        Assume.assumeFalse(isBatch);
+        assumeFalse(isBatch);
 
         BlockingIterator<RowData, Row> iterator =
                 BlockingIterator.of(
@@ -346,7 +346,7 @@ public class FileStoreITCase extends AbstractTestBase {
     }
 
     public FileStoreTable buildFileStoreTable(int[] partitions, int[] 
primaryKey) throws Exception {
-        return buildFileStoreTable(isBatch, TEMPORARY_FOLDER, partitions, 
primaryKey);
+        return buildFileStoreTable(isBatch, getTempDirPath(), partitions, 
primaryKey);
     }
 
     private static RowData srcRow(RowKind kind, int v, String p, int k) {
@@ -369,9 +369,9 @@ public class FileStoreITCase extends AbstractTestBase {
     }
 
     public static FileStoreTable buildFileStoreTable(
-            boolean noFail, TemporaryFolder temporaryFolder, int[] partitions, 
int[] primaryKey)
+            boolean noFail, String temporaryPath, int[] partitions, int[] 
primaryKey)
             throws Exception {
-        Options options = buildConfiguration(noFail, 
temporaryFolder.newFolder());
+        Options options = buildConfiguration(noFail, temporaryPath);
         Path tablePath = new CoreOptions(options.toMap()).path();
         UpdateSchema updateSchema =
                 new UpdateSchema(
@@ -392,15 +392,15 @@ public class FileStoreITCase extends AbstractTestBase {
                 });
     }
 
-    public static Options buildConfiguration(boolean noFail, File folder) {
+    public static Options buildConfiguration(boolean noFail, String 
temporaryPath) {
         Options options = new Options();
         options.set(BUCKET, NUM_BUCKET);
         if (noFail) {
-            options.set(PATH, folder.toURI().toString());
+            options.set(PATH, temporaryPath);
         } else {
             String failingName = UUID.randomUUID().toString();
             FailingFileIO.reset(failingName, 3, 100);
-            options.set(PATH, FailingFileIO.getFailingPath(failingName, 
folder.getPath()));
+            options.set(PATH, FailingFileIO.getFailingPath(failingName, 
temporaryPath));
         }
         options.set(FILE_FORMAT, "avro");
         return options;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileSystemCatalogITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileSystemCatalogITCase.java
index e600ab9d..78ea64aa 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileSystemCatalogITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FileSystemCatalogITCase.java
@@ -26,9 +26,8 @@ import org.apache.flink.table.store.kafka.KafkaTableTestBase;
 import org.apache.flink.types.Row;
 import org.apache.flink.util.CloseableIterator;
 
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.io.File;
 import java.io.IOException;
@@ -39,6 +38,7 @@ import java.util.UUID;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.AssertionsForClassTypes.assertThatThrownBy;
+import static org.junit.jupiter.api.Assertions.assertEquals;
 
 /** ITCase for {@link FlinkCatalog}. */
 public class FileSystemCatalogITCase extends KafkaTableTestBase {
@@ -46,9 +46,9 @@ public class FileSystemCatalogITCase extends 
KafkaTableTestBase {
     private String path;
     private static final String DB_NAME = "default";
 
-    @Before
+    @BeforeEach
     public void before() throws IOException {
-        path = TEMPORARY_FOLDER.newFolder().toURI().toString();
+        path = getTempDirPath();
         tEnv.executeSql(
                 String.format(
                         "CREATE CATALOG fs WITH ('type'='table-store', 
'warehouse'='%s')", path));
@@ -76,13 +76,15 @@ public class FileSystemCatalogITCase extends 
KafkaTableTestBase {
                 .hasMessage("Could not execute ALTER TABLE fs.default.t1 
RENAME TO fs.default.t2");
 
         tEnv.executeSql("ALTER TABLE t1 RENAME TO t3").await();
-        Assert.assertEquals(Arrays.asList(Row.of("t2"), Row.of("t3")), 
collect("SHOW TABLES"));
+        assertEquals(Arrays.asList(Row.of("t2"), Row.of("t3")), collect("SHOW 
TABLES"));
 
         Identifier identifier = new Identifier(DB_NAME, "t3");
         Catalog catalog =
                 ((FlinkCatalog) 
tEnv.getCatalog(tEnv.getCurrentCatalog()).get()).catalog();
         Path tablePath = catalog.getTableLocation(identifier);
-        Assert.assertEquals(tablePath.toString(), path + DB_NAME + ".db" + 
File.separator + "t3");
+        assertEquals(
+                tablePath.toString(),
+                new File(path, DB_NAME + ".db" + File.separator + 
"t3").toString());
 
         BlockingIterator<Row, Row> iterator =
                 BlockingIterator.of(tEnv.from("t3").execute().collect());
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ForceCompactionITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ForceCompactionITCase.java
index 1e4e308f..cf2d705d 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ForceCompactionITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/ForceCompactionITCase.java
@@ -31,7 +31,7 @@ import org.apache.flink.table.store.types.DataType;
 import org.apache.flink.table.store.types.RowType;
 import org.apache.flink.table.store.types.VarCharType;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.util.Arrays;
 import java.util.List;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FullCompactionFileStoreITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FullCompactionFileStoreITCase.java
index d66ce138..acb004fe 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FullCompactionFileStoreITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/FullCompactionFileStoreITCase.java
@@ -22,8 +22,8 @@ import 
org.apache.flink.table.store.file.utils.BlockingIterator;
 import org.apache.flink.types.Row;
 import org.apache.flink.types.RowKind;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
 
@@ -34,7 +34,7 @@ public class FullCompactionFileStoreITCase extends 
CatalogITCaseBase {
     private final String table = "T";
 
     @Override
-    @Before
+    @BeforeEach
     public void before() throws IOException {
         super.before();
         String options =
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LargeDataITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LargeDataITCase.java
index 5669a332..1b2e3ffc 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LargeDataITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LargeDataITCase.java
@@ -20,7 +20,7 @@ package org.apache.flink.table.store.connector;
 
 import org.apache.flink.types.Row;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.util.List;
 
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LogSystemITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LogSystemITCase.java
index 02b77fac..6ff30ad4 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LogSystemITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LogSystemITCase.java
@@ -23,8 +23,8 @@ import 
org.apache.flink.table.store.file.utils.BlockingIterator;
 import org.apache.flink.table.store.kafka.KafkaTableTestBase;
 import org.apache.flink.types.Row;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
 import java.util.List;
@@ -34,13 +34,13 @@ import static org.assertj.core.api.Assertions.assertThat;
 /** ITCase for table with log system. */
 public class LogSystemITCase extends KafkaTableTestBase {
 
-    @Before
+    @BeforeEach
     public void before() throws IOException {
         tEnv.executeSql(
                 String.format(
                         "CREATE CATALOG TABLE_STORE WITH ("
                                 + "'type'='table-store', 'warehouse'='%s')",
-                        TEMPORARY_FOLDER.newFolder().toURI()));
+                        getTempDirPath()));
         tEnv.useCatalog("TABLE_STORE");
     }
 
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LookupJoinITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LookupJoinITCase.java
index 54c7c34f..44f5e1ee 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LookupJoinITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/LookupJoinITCase.java
@@ -21,12 +21,12 @@ package org.apache.flink.table.store.connector;
 import org.apache.flink.table.api.EnvironmentSettings;
 import org.apache.flink.table.api.TableEnvironment;
 import org.apache.flink.table.api.config.ExecutionConfigOptions;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.table.store.file.utils.BlockingIterator;
-import org.apache.flink.test.util.AbstractTestBase;
 import org.apache.flink.types.Row;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.time.Duration;
 import java.util.List;
@@ -41,14 +41,14 @@ public class LookupJoinITCase extends AbstractTestBase {
 
     private TableEnvironment env;
 
-    @Before
+    @BeforeEach
     public void before() throws Exception {
         env = 
TableEnvironment.create(EnvironmentSettings.newInstance().inStreamingMode().build());
         env.getConfig().getConfiguration().set(CHECKPOINTING_INTERVAL, 
Duration.ofMillis(100));
         env.getConfig()
                 .getConfiguration()
                 
.set(ExecutionConfigOptions.TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM, 1);
-        String path = TEMPORARY_FOLDER.newFolder().toURI().toString();
+        String path = getTempDirPath();
         env.executeSql(
                 String.format(
                         "CREATE CATALOG my_catalog WITH ('type'='table-store', 
'warehouse'='%s')",
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/MappingTableITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/MappingTableITCase.java
index 7f909390..91baf6c6 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/MappingTableITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/MappingTableITCase.java
@@ -21,13 +21,13 @@ package org.apache.flink.table.store.connector;
 import org.apache.flink.table.api.EnvironmentSettings;
 import org.apache.flink.table.api.TableEnvironment;
 import org.apache.flink.table.api.ValidationException;
-import org.apache.flink.test.util.AbstractTestBase;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.types.Row;
 
 import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableList;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
 import java.util.List;
@@ -42,10 +42,10 @@ public class MappingTableITCase extends AbstractTestBase {
     private TableEnvironment tEnv;
     private String path;
 
-    @Before
+    @BeforeEach
     public void before() throws IOException {
         tEnv = 
TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());
-        path = TEMPORARY_FOLDER.newFolder().toURI().toString();
+        path = getTempDirPath();
     }
 
     @Test
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PartialUpdateITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PartialUpdateITCase.java
index 1f9a2f6b..a48b512d 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PartialUpdateITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PartialUpdateITCase.java
@@ -24,7 +24,7 @@ import org.apache.flink.types.RowKind;
 import org.apache.flink.util.CloseableIterator;
 
 import org.awaitility.Awaitility;
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.util.ArrayList;
 import java.util.Arrays;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PreAggregationITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PreAggregationITCase.java
index 732cd9a2..53d5a292 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PreAggregationITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PreAggregationITCase.java
@@ -20,7 +20,7 @@ package org.apache.flink.table.store.connector;
 
 import org.apache.flink.types.Row;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.math.BigDecimal;
 import java.time.LocalDate;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PredicateITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PredicateITCase.java
index 3a73c0bf..81651e83 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PredicateITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/PredicateITCase.java
@@ -20,7 +20,7 @@ package org.apache.flink.table.store.connector;
 
 import org.apache.flink.types.Row;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import static org.assertj.core.api.Assertions.assertThat;
 
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/RescaleBucketITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/RescaleBucketITCase.java
index cb90af51..a61dd8a1 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/RescaleBucketITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/RescaleBucketITCase.java
@@ -18,8 +18,7 @@
 
 package org.apache.flink.table.store.connector;
 
-import org.apache.flink.api.common.JobID;
-import org.apache.flink.client.program.ClusterClient;
+import org.apache.flink.core.execution.JobClient;
 import org.apache.flink.core.execution.SavepointFormatType;
 import org.apache.flink.runtime.jobgraph.SavepointConfigOptions;
 import org.apache.flink.table.store.file.Snapshot;
@@ -29,7 +28,7 @@ import 
org.apache.flink.table.store.file.utils.SnapshotManager;
 import org.apache.flink.table.store.fs.local.LocalFileIO;
 import org.apache.flink.types.Row;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import javax.annotation.Nullable;
 
@@ -82,13 +81,12 @@ public class RescaleBucketITCase extends CatalogITCaseBase {
                         + "END";
 
         
sEnv.getConfig().getConfiguration().set(SavepointConfigOptions.SAVEPOINT_PATH, 
path);
-        ClusterClient<?> client = MINI_CLUSTER_RESOURCE.getClusterClient();
 
         // step1: run streaming insert
-        JobID jobId = startJobAndCommitSnapshot(streamSql, null);
+        JobClient jobClient = startJobAndCommitSnapshot(streamSql, null);
 
         // step2: stop with savepoint
-        stopJobSafely(client, jobId);
+        stopJobSafely(jobClient);
 
         final Snapshot snapshotBeforeRescale = findLatestSnapshot("T3");
         assertThat(snapshotBeforeRescale).isNotNull();
@@ -109,9 +107,10 @@ public class RescaleBucketITCase extends CatalogITCaseBase 
{
         assertThat(batchSql("SELECT * FROM 
T3")).containsExactlyInAnyOrderElementsOf(committedData);
 
         // step5: resume streaming job
-        JobID resumedJobId = startJobAndCommitSnapshot(streamSql, 
snapshotAfterRescale.id());
+        JobClient resumedJobClient =
+                startJobAndCommitSnapshot(streamSql, 
snapshotAfterRescale.id());
         // stop job
-        stopJobSafely(client, resumedJobId);
+        stopJobSafely(resumedJobClient);
 
         // check snapshot and schema
         Snapshot lastSnapshot = findLatestSnapshot("T3");
@@ -137,18 +136,17 @@ public class RescaleBucketITCase extends 
CatalogITCaseBase {
         }
     }
 
-    private JobID startJobAndCommitSnapshot(String sql, @Nullable Long 
initSnapshotId)
+    private JobClient startJobAndCommitSnapshot(String sql, @Nullable Long 
initSnapshotId)
             throws Exception {
-        JobID jobId = sEnv.executeSql(sql).getJobClient().get().getJobID();
+        JobClient jobClient = sEnv.executeSql(sql).getJobClient().get();
         // let job run until the first snapshot is finished
         waitForTheNextSnapshot(initSnapshotId);
-        return jobId;
+        return jobClient;
     }
 
-    private void stopJobSafely(ClusterClient<?> client, JobID jobId)
-            throws ExecutionException, InterruptedException {
-        client.stopWithSavepoint(jobId, true, path, 
SavepointFormatType.DEFAULT);
-        while (!client.getJobStatus(jobId).get().isGloballyTerminalState()) {
+    private void stopJobSafely(JobClient client) throws ExecutionException, 
InterruptedException {
+        client.stopWithSavepoint(true, path, SavepointFormatType.DEFAULT);
+        while (!client.getJobStatus().get().isGloballyTerminalState()) {
             Thread.sleep(2000L);
         }
     }
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/SchemaChangeITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/SchemaChangeITCase.java
index 8be602fa..d275b0ef 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/SchemaChangeITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/SchemaChangeITCase.java
@@ -18,7 +18,7 @@
 
 package org.apache.flink.table.store.connector;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.util.Map;
 
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingReadWriteTableWithKafkaLogITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingReadWriteTableWithKafkaLogITCase.java
index 6327ed58..afdb0e09 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingReadWriteTableWithKafkaLogITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingReadWriteTableWithKafkaLogITCase.java
@@ -23,9 +23,9 @@ import 
org.apache.flink.table.store.file.utils.BlockingIterator;
 import org.apache.flink.table.store.kafka.KafkaTableTestBase;
 import org.apache.flink.types.Row;
 
-import org.junit.Before;
-import org.junit.Ignore;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.Test;
 
 import java.util.Arrays;
 import java.util.Collections;
@@ -57,7 +57,7 @@ import static 
org.apache.flink.table.store.connector.util.ReadWriteTableTestUtil
 /** Streaming reading and writing with Kafka log IT cases. */
 public class StreamingReadWriteTableWithKafkaLogITCase extends 
KafkaTableTestBase {
 
-    @Before
+    @BeforeEach
     public void setUp() throws Exception {
         init(createAndRegisterTempFile("").toString());
     }
@@ -1251,7 +1251,7 @@ public class StreamingReadWriteTableWithKafkaLogITCase 
extends KafkaTableTestBas
      * 
href="https://issues.apache.org/jira/browse/FLINK-28185";>FLINK-28185</a>. This 
bug will be
      * fixed in Flink-1.16.1 and after we update flink version this case can 
work.
      */
-    @Ignore
+    @Disabled
     @Test
     public void testReadInsertOnlyChangelogFromEnormousTimestamp() throws 
Exception {
         List<Row> initialRecords =
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingWarehouseITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingWarehouseITCase.java
index ee38d2f1..7b4b2e4c 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingWarehouseITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/StreamingWarehouseITCase.java
@@ -24,7 +24,7 @@ import 
org.apache.flink.table.store.file.utils.BlockingIterator;
 import org.apache.flink.table.store.kafka.KafkaTableTestBase;
 import org.apache.flink.types.Row;
 
-import org.junit.Test;
+import org.junit.jupiter.api.Test;
 
 import java.time.LocalDateTime;
 import java.util.function.Function;
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/CompactorSinkITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/CompactorSinkITCase.java
index ef5b1009..9ca62a84 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/CompactorSinkITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/CompactorSinkITCase.java
@@ -23,6 +23,7 @@ import 
org.apache.flink.streaming.api.datastream.DataStreamSource;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.store.connector.source.CompactorSourceBuilder;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.table.store.data.BinaryString;
 import org.apache.flink.table.store.data.GenericRow;
 import org.apache.flink.table.store.file.Snapshot;
@@ -41,11 +42,9 @@ import 
org.apache.flink.table.store.table.source.DataTableScan;
 import org.apache.flink.table.store.types.DataType;
 import org.apache.flink.table.store.types.DataTypes;
 import org.apache.flink.table.store.types.RowType;
-import org.apache.flink.test.util.AbstractTestBase;
 
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
 import java.util.Arrays;
@@ -55,6 +54,8 @@ import java.util.List;
 import java.util.Map;
 import java.util.UUID;
 
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
 /** IT cases for {@link CompactorSinkBuilder} and {@link CompactorSink}. */
 public class CompactorSinkITCase extends AbstractTestBase {
 
@@ -68,9 +69,9 @@ public class CompactorSinkITCase extends AbstractTestBase {
     private Path tablePath;
     private String commitUser;
 
-    @Before
+    @BeforeEach
     public void before() throws IOException {
-        tablePath = new Path(TEMPORARY_FOLDER.newFolder().toString());
+        tablePath = new Path(getTempDirPath());
         commitUser = UUID.randomUUID().toString();
     }
 
@@ -92,8 +93,8 @@ public class CompactorSinkITCase extends AbstractTestBase {
         commit.commit(1, write.prepareCommit(true, 1));
 
         Snapshot snapshot = 
snapshotManager.snapshot(snapshotManager.latestSnapshotId());
-        Assert.assertEquals(2, snapshot.id());
-        Assert.assertEquals(Snapshot.CommitKind.APPEND, snapshot.commitKind());
+        assertEquals(2, snapshot.id());
+        assertEquals(Snapshot.CommitKind.APPEND, snapshot.commitKind());
 
         write.close();
         commit.close();
@@ -112,18 +113,18 @@ public class CompactorSinkITCase extends AbstractTestBase 
{
         env.execute();
 
         snapshot = 
snapshotManager.snapshot(snapshotManager.latestSnapshotId());
-        Assert.assertEquals(3, snapshot.id());
-        Assert.assertEquals(Snapshot.CommitKind.COMPACT, 
snapshot.commitKind());
+        assertEquals(3, snapshot.id());
+        assertEquals(Snapshot.CommitKind.COMPACT, snapshot.commitKind());
 
         DataTableScan.DataFilePlan plan = table.newScan().plan();
-        Assert.assertEquals(3, plan.splits().size());
+        assertEquals(3, plan.splits().size());
         for (DataSplit split : plan.splits) {
             if (split.partition().getInt(1) == 15) {
                 // compacted
-                Assert.assertEquals(1, split.files().size());
+                assertEquals(1, split.files().size());
             } else {
                 // not compacted
-                Assert.assertEquals(2, split.files().size());
+                assertEquals(2, split.files().size());
             }
         }
     }
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/SinkSavepointITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/SinkSavepointITCase.java
index 81e026f7..d63c7a85 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/SinkSavepointITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/sink/SinkSavepointITCase.java
@@ -18,12 +18,11 @@
 
 package org.apache.flink.table.store.connector.sink;
 
-import org.apache.flink.api.common.JobID;
 import org.apache.flink.api.common.JobStatus;
-import org.apache.flink.client.program.ClusterClient;
 import org.apache.flink.configuration.CheckpointingOptions;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.configuration.StateBackendOptions;
+import org.apache.flink.core.execution.JobClient;
 import org.apache.flink.core.execution.SavepointFormatType;
 import org.apache.flink.runtime.jobgraph.SavepointRestoreSettings;
 import 
org.apache.flink.runtime.scheduler.stopwithsavepoint.StopWithSavepointStoppingException;
@@ -33,15 +32,15 @@ import org.apache.flink.table.api.EnvironmentSettings;
 import org.apache.flink.table.api.TableEnvironment;
 import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
 import org.apache.flink.table.api.config.ExecutionConfigOptions;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.table.store.file.utils.FailingFileIO;
-import org.apache.flink.test.util.AbstractTestBase;
 import org.apache.flink.types.Row;
 import org.apache.flink.util.CloseableIterator;
 import org.apache.flink.util.ExceptionUtils;
 
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
 
 import java.time.Duration;
 import java.util.ArrayList;
@@ -53,47 +52,46 @@ import java.util.concurrent.ThreadLocalRandom;
 import java.util.stream.Collectors;
 import java.util.stream.IntStream;
 
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
 /** IT cases for {@link FileStoreSink} when writing file store and with 
savepoints. */
 public class SinkSavepointITCase extends AbstractTestBase {
 
     private String path;
     private String failingName;
 
-    @Before
+    @BeforeEach
     public void before() throws Exception {
-        path = TEMPORARY_FOLDER.newFolder().toPath().toString();
+        path = getTempDirPath();
         // for failure tests
         failingName = UUID.randomUUID().toString();
         FailingFileIO.reset(failingName, 100, 500);
     }
 
-    @Test(timeout = 180000)
+    @Test
+    @Timeout(180000)
     public void testRecoverFromSavepoint() throws Exception {
         String failingPath = FailingFileIO.getFailingPath(failingName, path);
         String savepointPath = null;
-        JobID jobId;
-        ClusterClient<?> client = MINI_CLUSTER_RESOURCE.getClusterClient();
         ThreadLocalRandom random = ThreadLocalRandom.current();
 
         OUTER:
         while (true) {
             // start a new job or recover from savepoint
-            jobId = runRecoverFromSavepointJob(failingPath, savepointPath);
+            JobClient jobClient = runRecoverFromSavepointJob(failingPath, 
savepointPath);
             while (true) {
                 // wait for a random number of time before stopping with 
savepoint
                 Thread.sleep(random.nextInt(5000));
-                if (client.getJobStatus(jobId).get() == JobStatus.FINISHED) {
+                if (jobClient.getJobStatus().get() == JobStatus.FINISHED) {
                     // job finished, check for result
                     break OUTER;
                 }
                 try {
                     // try to stop with savepoint
                     savepointPath =
-                            client.stopWithSavepoint(
-                                            jobId,
-                                            false,
-                                            path + "/savepoint",
-                                            SavepointFormatType.DEFAULT)
+                            jobClient
+                                    .stopWithSavepoint(
+                                            false, path + "/savepoint", 
SavepointFormatType.DEFAULT)
                                     .get();
                     break;
                 } catch (Exception e) {
@@ -113,7 +111,7 @@ public class SinkSavepointITCase extends AbstractTestBase {
                 }
             }
             // wait for job to stop
-            while 
(!client.getJobStatus(jobId).get().isGloballyTerminalState()) {
+            while (!jobClient.getJobStatus().get().isGloballyTerminalState()) {
                 Thread.sleep(1000);
             }
             // recover from savepoint in the next round
@@ -122,7 +120,7 @@ public class SinkSavepointITCase extends AbstractTestBase {
         checkRecoverFromSavepointResult(failingPath);
     }
 
-    private JobID runRecoverFromSavepointJob(String failingPath, String 
savepointPath)
+    private JobClient runRecoverFromSavepointJob(String failingPath, String 
savepointPath)
             throws Exception {
         Configuration conf = new Configuration();
         if (savepointPath != null) {
@@ -187,17 +185,15 @@ public class SinkSavepointITCase extends AbstractTestBase 
{
         FailingFileIO.retryArtificialException(() -> 
tEnv.executeSql(createSinkSql));
 
         String insertIntoSql = "INSERT INTO T SELECT * FROM 
default_catalog.default_database.S";
-        JobID jobId =
+        JobClient jobClient =
                 FailingFileIO.retryArtificialException(() -> 
tEnv.executeSql(insertIntoSql))
                         .getJobClient()
-                        .get()
-                        .getJobID();
+                        .get();
 
-        ClusterClient<?> client = MINI_CLUSTER_RESOURCE.getClusterClient();
-        while (client.getJobStatus(jobId).get() == JobStatus.INITIALIZING) {
+        while (jobClient.getJobStatus().get() == JobStatus.INITIALIZING) {
             Thread.sleep(1000);
         }
-        return jobId;
+        return jobClient;
     }
 
     private void checkRecoverFromSavepointResult(String failingPath) throws 
Exception {
@@ -221,12 +217,11 @@ public class SinkSavepointITCase extends AbstractTestBase 
{
         try (CloseableIterator<Row> it = tEnv.executeSql("SELECT * FROM 
T").collect()) {
             while (it.hasNext()) {
                 Row row = it.next();
-                Assert.assertEquals(1, row.getArity());
+                assertEquals(1, row.getArity());
                 actual.add((Integer) row.getField(0));
             }
         }
         Collections.sort(actual);
-        Assert.assertEquals(
-                IntStream.range(0, 
100000).boxed().collect(Collectors.toList()), actual);
+        assertEquals(IntStream.range(0, 
100000).boxed().collect(Collectors.toList()), actual);
     }
 }
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/source/CompactorSourceITCase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/source/CompactorSourceITCase.java
index 70dca8da..54caf419 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/source/CompactorSourceITCase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/source/CompactorSourceITCase.java
@@ -21,6 +21,7 @@ package org.apache.flink.table.store.connector.source;
 import org.apache.flink.streaming.api.datastream.DataStreamSource;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.table.store.data.BinaryRow;
 import org.apache.flink.table.store.data.BinaryRowWriter;
 import org.apache.flink.table.store.data.BinaryString;
@@ -38,11 +39,10 @@ import org.apache.flink.table.store.table.sink.TableWrite;
 import org.apache.flink.table.store.types.DataType;
 import org.apache.flink.table.store.types.DataTypes;
 import org.apache.flink.table.store.types.RowType;
-import org.apache.flink.test.util.AbstractTestBase;
 import org.apache.flink.util.CloseableIterator;
 
-import org.junit.Before;
-import org.junit.Test;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 
 import java.io.IOException;
 import java.io.UncheckedIOException;
@@ -71,9 +71,9 @@ public class CompactorSourceITCase extends AbstractTestBase {
     private Path tablePath;
     private String commitUser;
 
-    @Before
+    @BeforeEach
     public void before() throws IOException {
-        tablePath = new Path(TEMPORARY_FOLDER.newFolder().toString());
+        tablePath = new Path(getTempDirPath());
         commitUser = UUID.randomUUID().toString();
     }
 
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/AbstractTestBase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/AbstractTestBase.java
index 78a187e2..b6f2b396 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/AbstractTestBase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/AbstractTestBase.java
@@ -18,11 +18,13 @@
 
 package org.apache.flink.table.store.connector.util;
 
+import org.apache.flink.client.program.ClusterClient;
+import org.apache.flink.runtime.client.JobStatusMessage;
 import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
 import org.apache.flink.table.store.utils.FileIOUtils;
-import org.apache.flink.test.junit5.MiniClusterExtension;
 import org.apache.flink.util.TestLogger;
 
+import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.extension.RegisterExtension;
 import org.junit.jupiter.api.io.TempDir;
 import org.slf4j.Logger;
@@ -41,8 +43,8 @@ public class AbstractTestBase extends TestLogger {
     private static final int DEFAULT_PARALLELISM = 4;
 
     @RegisterExtension
-    protected static final MiniClusterExtension MINI_CLUSTER_EXTENSION =
-            new MiniClusterExtension(
+    protected static final MiniClusterWithClientExtension 
MINI_CLUSTER_EXTENSION =
+            new MiniClusterWithClientExtension(
                     new MiniClusterResourceConfiguration.Builder()
                             .setNumberTaskManagers(1)
                             .setNumberSlotsPerTaskManager(DEFAULT_PARALLELISM)
@@ -50,6 +52,20 @@ public class AbstractTestBase extends TestLogger {
 
     @TempDir protected static Path temporaryFolder;
 
+    @AfterEach
+    public final void cleanupRunningJobs() throws Exception {
+        ClusterClient<?> clusterClient = 
MINI_CLUSTER_EXTENSION.createRestClusterClient();
+        for (JobStatusMessage path : clusterClient.listJobs().get()) {
+            if (!path.getJobState().isTerminalState()) {
+                try {
+                    clusterClient.cancel(path.getJobId()).get();
+                } catch (Exception ignored) {
+                    // ignore exceptions when cancelling dangling jobs
+                }
+            }
+        }
+    }
+
     // 
----------------------------------------------------------------------------------------------------------------
     //  Temporary File Utilities
     // 
----------------------------------------------------------------------------------------------------------------
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/MiniClusterWithClientExtension.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/MiniClusterWithClientExtension.java
new file mode 100644
index 00000000..ad226424
--- /dev/null
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/connector/util/MiniClusterWithClientExtension.java
@@ -0,0 +1,225 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.table.store.connector.util;
+
+import org.apache.flink.annotation.Experimental;
+import org.apache.flink.client.program.ClusterClient;
+import org.apache.flink.client.program.MiniClusterClient;
+import org.apache.flink.client.program.rest.RestClusterClient;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.CoreOptions;
+import org.apache.flink.runtime.testutils.InternalMiniClusterExtension;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.streaming.util.TestStreamEnvironment;
+import org.apache.flink.test.junit5.InjectClusterClient;
+import org.apache.flink.test.util.TestEnvironment;
+
+import org.junit.jupiter.api.extension.AfterAllCallback;
+import org.junit.jupiter.api.extension.AfterEachCallback;
+import org.junit.jupiter.api.extension.BeforeAllCallback;
+import org.junit.jupiter.api.extension.BeforeEachCallback;
+import org.junit.jupiter.api.extension.ExtensionContext;
+import org.junit.jupiter.api.extension.ParameterContext;
+import org.junit.jupiter.api.extension.ParameterResolutionException;
+import org.junit.jupiter.api.extension.ParameterResolver;
+
+import java.util.function.Supplier;
+
+/** Mini cluster extension with cluster client which is used to cancel jobs. */
+@Experimental
+public final class MiniClusterWithClientExtension
+        implements BeforeAllCallback,
+                BeforeEachCallback,
+                AfterEachCallback,
+                AfterAllCallback,
+                ParameterResolver {
+
+    private static final ExtensionContext.Namespace NAMESPACE =
+            
ExtensionContext.Namespace.create(MiniClusterWithClientExtension.class);
+
+    private static final String CLUSTER_REST_CLIENT = "clusterRestClient";
+    private static final String MINI_CLUSTER_CLIENT = "miniClusterClient";
+
+    private final Supplier<MiniClusterResourceConfiguration>
+            miniClusterResourceConfigurationSupplier;
+
+    private InternalMiniClusterExtension internalMiniClusterExtension;
+
+    public MiniClusterWithClientExtension(
+            final MiniClusterResourceConfiguration 
miniClusterResourceConfiguration) {
+        this(() -> miniClusterResourceConfiguration);
+    }
+
+    @Experimental
+    public MiniClusterWithClientExtension(
+            Supplier<MiniClusterResourceConfiguration> 
miniClusterResourceConfigurationSupplier) {
+        this.miniClusterResourceConfigurationSupplier = 
miniClusterResourceConfigurationSupplier;
+    }
+
+    // Accessors
+
+    @Override
+    public boolean supportsParameter(
+            ParameterContext parameterContext, ExtensionContext 
extensionContext)
+            throws ParameterResolutionException {
+        Class<?> parameterType = parameterContext.getParameter().getType();
+        if (parameterContext.isAnnotated(InjectClusterClient.class)
+                && ClusterClient.class.isAssignableFrom(parameterType)) {
+            return true;
+        }
+        return 
internalMiniClusterExtension.supportsParameter(parameterContext, 
extensionContext);
+    }
+
+    @Override
+    public Object resolveParameter(
+            ParameterContext parameterContext, ExtensionContext 
extensionContext)
+            throws ParameterResolutionException {
+        Class<?> parameterType = parameterContext.getParameter().getType();
+        if (parameterContext.isAnnotated(InjectClusterClient.class)) {
+            if (parameterType.equals(RestClusterClient.class)) {
+                return extensionContext
+                        .getStore(NAMESPACE)
+                        .getOrComputeIfAbsent(
+                                CLUSTER_REST_CLIENT,
+                                k -> {
+                                    try {
+                                        return new CloseableParameter<>(
+                                                createRestClusterClient(
+                                                        
internalMiniClusterExtension));
+                                    } catch (Exception e) {
+                                        throw new ParameterResolutionException(
+                                                "Cannot create rest cluster 
client", e);
+                                    }
+                                },
+                                CloseableParameter.class)
+                        .get();
+            }
+            // Default to MiniClusterClient
+            return extensionContext
+                    .getStore(NAMESPACE)
+                    .getOrComputeIfAbsent(
+                            MINI_CLUSTER_CLIENT,
+                            k -> {
+                                try {
+                                    return new CloseableParameter<>(
+                                            
createMiniClusterClient(internalMiniClusterExtension));
+                                } catch (Exception e) {
+                                    throw new ParameterResolutionException(
+                                            "Cannot create mini cluster 
client", e);
+                                }
+                            },
+                            CloseableParameter.class)
+                    .get();
+        }
+        return internalMiniClusterExtension.resolveParameter(parameterContext, 
extensionContext);
+    }
+
+    // Lifecycle implementation
+
+    @Override
+    public void beforeAll(ExtensionContext context) throws Exception {
+        internalMiniClusterExtension =
+                new 
InternalMiniClusterExtension(miniClusterResourceConfigurationSupplier.get());
+        internalMiniClusterExtension.beforeAll(context);
+    }
+
+    @Override
+    public void beforeEach(ExtensionContext context) throws Exception {
+        registerEnv(internalMiniClusterExtension);
+    }
+
+    @Override
+    public void afterEach(ExtensionContext context) throws Exception {
+        unregisterEnv(internalMiniClusterExtension);
+    }
+
+    @Override
+    public void afterAll(ExtensionContext context) throws Exception {
+        if (internalMiniClusterExtension != null) {
+            internalMiniClusterExtension.afterAll(context);
+        }
+    }
+
+    // Implementation
+
+    private void registerEnv(InternalMiniClusterExtension 
internalMiniClusterExtension) {
+        final Configuration configuration =
+                
internalMiniClusterExtension.getMiniCluster().getConfiguration();
+
+        final int defaultParallelism =
+                configuration
+                        .getOptional(CoreOptions.DEFAULT_PARALLELISM)
+                        .orElse(internalMiniClusterExtension.getNumberSlots());
+
+        TestEnvironment executionEnvironment =
+                new TestEnvironment(
+                        internalMiniClusterExtension.getMiniCluster(), 
defaultParallelism, false);
+        executionEnvironment.setAsContext();
+        TestStreamEnvironment.setAsContext(
+                internalMiniClusterExtension.getMiniCluster(), 
defaultParallelism);
+    }
+
+    private void unregisterEnv(InternalMiniClusterExtension 
internalMiniClusterExtension) {
+        TestStreamEnvironment.unsetAsContext();
+        TestEnvironment.unsetAsContext();
+    }
+
+    private MiniClusterClient createMiniClusterClient(
+            InternalMiniClusterExtension internalMiniClusterExtension) {
+        return new MiniClusterClient(
+                internalMiniClusterExtension.getClientConfiguration(),
+                internalMiniClusterExtension.getMiniCluster());
+    }
+
+    private RestClusterClient<MiniClusterClient.MiniClusterId> 
createRestClusterClient(
+            InternalMiniClusterExtension internalMiniClusterExtension) throws 
Exception {
+        return new RestClusterClient<>(
+                internalMiniClusterExtension.getClientConfiguration(),
+                MiniClusterClient.MiniClusterId.INSTANCE);
+    }
+
+    public RestClusterClient<MiniClusterClient.MiniClusterId> 
createRestClusterClient()
+            throws Exception {
+        return createRestClusterClient(internalMiniClusterExtension);
+    }
+
+    // Utils
+
+    public Configuration getClientConfiguration() {
+        return internalMiniClusterExtension.getClientConfiguration();
+    }
+
+    private static class CloseableParameter<T extends AutoCloseable>
+            implements ExtensionContext.Store.CloseableResource {
+        private final T autoCloseable;
+
+        CloseableParameter(T autoCloseable) {
+            this.autoCloseable = autoCloseable;
+        }
+
+        public T get() {
+            return autoCloseable;
+        }
+
+        @Override
+        public void close() throws Throwable {
+            this.autoCloseable.close();
+        }
+    }
+}
diff --git 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/kafka/KafkaTableTestBase.java
 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/kafka/KafkaTableTestBase.java
index 5a0e3aa1..667bd456 100644
--- 
a/flink-table-store-connector/src/test/java/org/apache/flink/table/store/kafka/KafkaTableTestBase.java
+++ 
b/flink-table-store-connector/src/test/java/org/apache/flink/table/store/kafka/KafkaTableTestBase.java
@@ -22,7 +22,7 @@ import 
org.apache.flink.api.common.restartstrategy.RestartStrategies;
 import 
org.apache.flink.streaming.api.environment.ExecutionCheckpointingOptions;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
-import org.apache.flink.test.util.AbstractTestBase;
+import org.apache.flink.table.store.connector.util.AbstractTestBase;
 import org.apache.flink.util.DockerImageVersions;
 
 import org.apache.kafka.clients.admin.AdminClient;
@@ -35,9 +35,12 @@ import org.apache.kafka.common.TopicPartition;
 import org.apache.kafka.common.errors.TopicExistsException;
 import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
 import org.apache.kafka.common.serialization.StringDeserializer;
-import org.junit.After;
-import org.junit.Before;
-import org.junit.ClassRule;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.extension.AfterAllCallback;
+import org.junit.jupiter.api.extension.BeforeAllCallback;
+import org.junit.jupiter.api.extension.ExtensionContext;
+import org.junit.jupiter.api.extension.RegisterExtension;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.testcontainers.containers.KafkaContainer;
@@ -66,24 +69,26 @@ public abstract class KafkaTableTestBase extends 
AbstractTestBase {
     private static final Network NETWORK = Network.newNetwork();
     private static final int zkTimeoutMills = 30000;
 
-    @ClassRule
-    public static final KafkaContainer KAFKA_CONTAINER =
-            new 
KafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)) {
-                @Override
-                protected void doStart() {
-                    super.doStart();
-                    if (LOG.isInfoEnabled()) {
-                        this.followOutput(new Slf4jLogConsumer(LOG));
-                    }
-                }
-            }.withEmbeddedZookeeper()
-                    .withNetwork(NETWORK)
-                    .withNetworkAliases(INTER_CONTAINER_KAFKA_ALIAS)
-                    .withEnv(
-                            "KAFKA_TRANSACTION_MAX_TIMEOUT_MS",
-                            String.valueOf(Duration.ofHours(2).toMillis()))
-                    // Disable log deletion to prevent records from being 
deleted during test run
-                    .withEnv("KAFKA_LOG_RETENTION_MS", "-1");
+    @RegisterExtension
+    public static final KafkaContainerExtension KAFKA_CONTAINER =
+            (KafkaContainerExtension)
+                    new 
KafkaContainerExtension(DockerImageName.parse(DockerImageVersions.KAFKA)) {
+                        @Override
+                        protected void doStart() {
+                            super.doStart();
+                            if (LOG.isInfoEnabled()) {
+                                this.followOutput(new Slf4jLogConsumer(LOG));
+                            }
+                        }
+                    }.withEmbeddedZookeeper()
+                            .withNetwork(NETWORK)
+                            .withNetworkAliases(INTER_CONTAINER_KAFKA_ALIAS)
+                            .withEnv(
+                                    "KAFKA_TRANSACTION_MAX_TIMEOUT_MS",
+                                    
String.valueOf(Duration.ofHours(2).toMillis()))
+                            // Disable log deletion to prevent records from 
being deleted during
+                            // test run
+                            .withEnv("KAFKA_LOG_RETENTION_MS", "-1");
 
     protected StreamExecutionEnvironment env;
     protected StreamTableEnvironment tEnv;
@@ -91,7 +96,7 @@ public abstract class KafkaTableTestBase extends 
AbstractTestBase {
     // Timer for scheduling logging task if the test hangs
     private final Timer loggingTimer = new Timer("Debug Logging Timer");
 
-    @Before
+    @BeforeEach
     public void setup() {
         env = StreamExecutionEnvironment.getExecutionEnvironment();
         tEnv = StreamTableEnvironment.create(env);
@@ -114,7 +119,7 @@ public abstract class KafkaTableTestBase extends 
AbstractTestBase {
                 });
     }
 
-    @After
+    @AfterEach
     public void after() throws ExecutionException, InterruptedException {
         // Cancel timer for debug logging
         cancelTimeoutLogger();
@@ -246,4 +251,22 @@ public abstract class KafkaTableTestBase extends 
AbstractTestBase {
                                 beginningOffsets.get(partition),
                                 endOffsets.get(partition)));
     }
+
+    /** Kafka container extension for junit5. */
+    private static class KafkaContainerExtension extends KafkaContainer
+            implements BeforeAllCallback, AfterAllCallback {
+        private KafkaContainerExtension(DockerImageName dockerImageName) {
+            super(dockerImageName);
+        }
+
+        @Override
+        public void beforeAll(ExtensionContext extensionContext) throws 
Exception {
+            this.doStart();
+        }
+
+        @Override
+        public void afterAll(ExtensionContext extensionContext) throws 
Exception {
+            this.close();
+        }
+    }
 }

Reply via email to