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();
+ }
+ }
}