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

martijnvisser pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-jdbc.git

commit 7addee18872af6333babbf6314fe073699dbd2f1
Author: Joao Boto <[email protected]>
AuthorDate: Mon Jan 30 11:38:55 2023 +0100

    [FLINK-30790] Refactor MySqlCatalogITCase create tests by database version
---
 .../6b9ab1b0-c14d-4667-bab5-407b81fba98b           |  12 +
 .../jdbc/catalog/MySql56CatalogITCase.java         |  39 ++
 .../jdbc/catalog/MySql57CatalogITCase.java         |  39 ++
 .../connector/jdbc/catalog/MySqlCatalogITCase.java | 366 +------------------
 .../jdbc/catalog/MySqlCatalogTestBase.java         | 396 ++++++++++++++++++---
 .../sqlserver/SqlServerTableSinkITCase.java        |   3 -
 .../sqlserver/SqlServerTableSourceITCase.java      |   3 -
 .../jdbc/xa/JdbcExactlyOnceSinkE2eTest.java        |  12 +-
 8 files changed, 459 insertions(+), 411 deletions(-)

diff --git 
a/flink-connector-jdbc/archunit-violations/6b9ab1b0-c14d-4667-bab5-407b81fba98b 
b/flink-connector-jdbc/archunit-violations/6b9ab1b0-c14d-4667-bab5-407b81fba98b
index 9af409a..e54d3ff 100644
--- 
a/flink-connector-jdbc/archunit-violations/6b9ab1b0-c14d-4667-bab5-407b81fba98b
+++ 
b/flink-connector-jdbc/archunit-violations/6b9ab1b0-c14d-4667-bab5-407b81fba98b
@@ -8,6 +8,18 @@ org.apache.flink.connector.jdbc.catalog.MySqlCatalogITCase 
does not satisfy: onl
 * reside in a package 'org.apache.flink.runtime.*' and contain any fields that 
are static, final, and of type InternalMiniClusterExtension and annotated with 
@RegisterExtension\
 * reside outside of package 'org.apache.flink.runtime.*' and contain any 
fields that are static, final, and of type MiniClusterExtension and annotated 
with @RegisterExtension\
 * reside in a package 'org.apache.flink.runtime.*' and is annotated with 
@ExtendWith with class InternalMiniClusterExtension\
+* reside outside of package 'org.apache.flink.runtime.*' and is annotated with 
@ExtendWith with class MiniClusterExtension\
+ or contain any fields that are public, static, and of type 
MiniClusterWithClientResource and final and annotated with @ClassRule or 
contain any fields that is of type MiniClusterWithClientResource and public and 
final and not static and annotated with @Rule
+org.apache.flink.connector.jdbc.catalog.MySql57CatalogITCase does not satisfy: 
only one of the following predicates match:\
+* reside in a package 'org.apache.flink.runtime.*' and contain any fields that 
are static, final, and of type InternalMiniClusterExtension and annotated with 
@RegisterExtension\
+* reside outside of package 'org.apache.flink.runtime.*' and contain any 
fields that are static, final, and of type MiniClusterExtension and annotated 
with @RegisterExtension\
+* reside in a package 'org.apache.flink.runtime.*' and is annotated with 
@ExtendWith with class InternalMiniClusterExtension\
+* reside outside of package 'org.apache.flink.runtime.*' and is annotated with 
@ExtendWith with class MiniClusterExtension\
+ or contain any fields that are public, static, and of type 
MiniClusterWithClientResource and final and annotated with @ClassRule or 
contain any fields that is of type MiniClusterWithClientResource and public and 
final and not static and annotated with @Rule
+org.apache.flink.connector.jdbc.catalog.MySql56CatalogITCase does not satisfy: 
only one of the following predicates match:\
+* reside in a package 'org.apache.flink.runtime.*' and contain any fields that 
are static, final, and of type InternalMiniClusterExtension and annotated with 
@RegisterExtension\
+* reside outside of package 'org.apache.flink.runtime.*' and contain any 
fields that are static, final, and of type MiniClusterExtension and annotated 
with @RegisterExtension\
+* reside in a package 'org.apache.flink.runtime.*' and is annotated with 
@ExtendWith with class InternalMiniClusterExtension\
 * reside outside of package 'org.apache.flink.runtime.*' and is annotated with 
@ExtendWith with class MiniClusterExtension\
  or contain any fields that are public, static, and of type 
MiniClusterWithClientResource and final and annotated with @ClassRule or 
contain any fields that is of type MiniClusterWithClientResource and public and 
final and not static and annotated with @Rule
 org.apache.flink.connector.jdbc.catalog.PostgresCatalogITCase does not 
satisfy: only one of the following predicates match:\
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySql56CatalogITCase.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySql56CatalogITCase.java
new file mode 100644
index 0000000..2ff3ed0
--- /dev/null
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySql56CatalogITCase.java
@@ -0,0 +1,39 @@
+/*
+ * 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.connector.jdbc.catalog;
+
+import org.apache.flink.connector.jdbc.test.DockerImageVersions;
+
+import org.testcontainers.containers.MySQLContainer;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+
+/** E2E test for {@link MySqlCatalog}. */
+@Testcontainers
+public class MySql56CatalogITCase extends MySqlCatalogTestBase {
+
+    @Container
+    private static final MySQLContainer<?> CONTAINER =
+            createContainer(DockerImageVersions.MYSQL_5_6);
+
+    @Override
+    protected String getDatabaseUrl() {
+        return CONTAINER.getJdbcUrl().substring(0, 
CONTAINER.getJdbcUrl().lastIndexOf("/"));
+    }
+}
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySql57CatalogITCase.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySql57CatalogITCase.java
new file mode 100644
index 0000000..0a1dc8b
--- /dev/null
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySql57CatalogITCase.java
@@ -0,0 +1,39 @@
+/*
+ * 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.connector.jdbc.catalog;
+
+import org.apache.flink.connector.jdbc.test.DockerImageVersions;
+
+import org.testcontainers.containers.MySQLContainer;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+
+/** E2E test for {@link MySqlCatalog}. */
+@Testcontainers
+public class MySql57CatalogITCase extends MySqlCatalogTestBase {
+
+    @Container
+    private static final MySQLContainer<?> CONTAINER =
+            createContainer(DockerImageVersions.MYSQL_5_7);
+
+    @Override
+    protected String getDatabaseUrl() {
+        return CONTAINER.getJdbcUrl().substring(0, 
CONTAINER.getJdbcUrl().lastIndexOf("/"));
+    }
+}
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogITCase.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogITCase.java
index 9d57400..73b5acf 100644
--- 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogITCase.java
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogITCase.java
@@ -18,367 +18,21 @@
 
 package org.apache.flink.connector.jdbc.catalog;
 
-import org.apache.flink.table.api.DataTypes;
-import org.apache.flink.table.api.EnvironmentSettings;
-import org.apache.flink.table.api.Schema;
-import org.apache.flink.table.api.TableEnvironment;
-import org.apache.flink.table.catalog.CatalogBaseTable;
-import org.apache.flink.table.catalog.ObjectPath;
-import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException;
-import org.apache.flink.table.catalog.exceptions.TableNotExistException;
-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.CollectionUtil;
+import org.apache.flink.connector.jdbc.test.DockerImageVersions;
 
-import org.apache.flink.shaded.guava30.com.google.common.collect.Lists;
-
-import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.TestTemplate;
-import org.junit.jupiter.api.extension.ExtendWith;
-
-import java.math.BigDecimal;
-import java.sql.Date;
-import java.sql.Time;
-import java.sql.Timestamp;
-import java.util.Arrays;
-import java.util.Collection;
-import java.util.Collections;
-import java.util.List;
-
-import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
-import static 
org.apache.flink.table.api.config.ExecutionConfigOptions.TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM;
-import static org.assertj.core.api.Assertions.assertThat;
-import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import org.testcontainers.containers.MySQLContainer;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
 
 /** E2E test for {@link MySqlCatalog}. */
-@ExtendWith(ParameterizedTestExtension.class)
+@Testcontainers
 public class MySqlCatalogITCase extends MySqlCatalogTestBase {
 
-    private static final List<Row> ALL_TYPES_ROWS =
-            Lists.newArrayList(
-                    Row.ofKind(
-                            RowKind.INSERT,
-                            1L,
-                            -1L,
-                            new BigDecimal(1),
-                            null,
-                            true,
-                            null,
-                            "hello",
-                            Date.valueOf("2021-08-04").toLocalDate(),
-                            Timestamp.valueOf("2021-08-04 
01:54:16").toLocalDateTime(),
-                            new BigDecimal(-1),
-                            new BigDecimal(1),
-                            -1.0d,
-                            1.0d,
-                            "enum2",
-                            -9.1f,
-                            9.1f,
-                            -1,
-                            1L,
-                            -1,
-                            1L,
-                            null,
-                            "col_longtext",
-                            null,
-                            -1,
-                            1,
-                            "col_mediumtext",
-                            new BigDecimal(-99),
-                            new BigDecimal(99),
-                            -1.0d,
-                            1.0d,
-                            "set_ele1",
-                            Short.parseShort("-1"),
-                            1,
-                            "col_text",
-                            Time.valueOf("10:32:34").toLocalTime(),
-                            Timestamp.valueOf("2021-08-04 
01:54:16").toLocalDateTime(),
-                            "col_tinytext",
-                            Byte.parseByte("-1"),
-                            Short.parseShort("1"),
-                            null,
-                            "col_varchar",
-                            Timestamp.valueOf("2021-08-04 
01:54:16.463").toLocalDateTime(),
-                            Time.valueOf("09:33:43").toLocalTime(),
-                            Timestamp.valueOf("2021-08-04 
01:54:16.463").toLocalDateTime(),
-                            null),
-                    Row.ofKind(
-                            RowKind.INSERT,
-                            2L,
-                            -1L,
-                            new BigDecimal(1),
-                            null,
-                            true,
-                            null,
-                            "hello",
-                            Date.valueOf("2021-08-04").toLocalDate(),
-                            Timestamp.valueOf("2021-08-04 
01:53:19").toLocalDateTime(),
-                            new BigDecimal(-1),
-                            new BigDecimal(1),
-                            -1.0d,
-                            1.0d,
-                            "enum2",
-                            -9.1f,
-                            9.1f,
-                            -1,
-                            1L,
-                            -1,
-                            1L,
-                            null,
-                            "col_longtext",
-                            null,
-                            -1,
-                            1,
-                            "col_mediumtext",
-                            new BigDecimal(-99),
-                            new BigDecimal(99),
-                            -1.0d,
-                            1.0d,
-                            "set_ele1,set_ele12",
-                            Short.parseShort("-1"),
-                            1,
-                            "col_text",
-                            Time.valueOf("10:32:34").toLocalTime(),
-                            Timestamp.valueOf("2021-08-04 
01:53:19").toLocalDateTime(),
-                            "col_tinytext",
-                            Byte.parseByte("-1"),
-                            Short.parseShort("1"),
-                            null,
-                            "col_varchar",
-                            Timestamp.valueOf("2021-08-04 
01:53:19.098").toLocalDateTime(),
-                            Time.valueOf("09:33:43").toLocalTime(),
-                            Timestamp.valueOf("2021-08-04 
01:53:19.098").toLocalDateTime(),
-                            null));
-
-    private final MySqlCatalog catalog;
-    private TableEnvironment tEnv;
-
-    @BeforeEach
-    void setup() {
-        this.tEnv = 
TableEnvironment.create(EnvironmentSettings.inStreamingMode());
-        tEnv.getConfig().set(TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM, 1);
-
-        // Use mysql catalog.
-        tEnv.registerCatalog(TEST_CATALOG_NAME, catalog);
-        tEnv.useCatalog(TEST_CATALOG_NAME);
-    }
-
-    public MySqlCatalogITCase(String version) {
-        catalog = CATALOGS.get(version);
-    }
-
-    @Parameters(name = "version = {0}")
-    public static Collection<String> params() {
-        return DOCKER_IMAGE_NAMES;
-    }
-
-    // ------ databases ------
-
-    @TestTemplate
-    void testGetDb_DatabaseNotExistException() throws Exception {
-        String databaseNotExist = "nonexistent";
-        assertThatThrownBy(() -> catalog.getDatabase(databaseNotExist))
-                .satisfies(
-                        anyCauseMatches(
-                                DatabaseNotExistException.class,
-                                String.format(
-                                        "Database %s does not exist in 
Catalog",
-                                        databaseNotExist)));
-    }
-
-    @TestTemplate
-    void testListDatabases() {
-        List<String> actual = catalog.listDatabases();
-        assertThat(actual).containsExactly(TEST_DB, TEST_DB2);
-    }
-
-    @TestTemplate
-    void testDbExists() throws Exception {
-        String databaseNotExist = "nonexistent";
-        assertThat(catalog.databaseExists(databaseNotExist)).isFalse();
-        assertThat(catalog.databaseExists(TEST_DB)).isTrue();
-    }
-
-    // ------ tables ------
-
-    @TestTemplate
-    void testListTables() throws DatabaseNotExistException {
-        List<String> actual = catalog.listTables(TEST_DB);
-        assertThat(actual)
-                .isEqualTo(
-                        Arrays.asList(
-                                TEST_TABLE_ALL_TYPES,
-                                TEST_SINK_TABLE_ALL_TYPES,
-                                TEST_TABLE_SINK_FROM_GROUPED_BY,
-                                TEST_TABLE_PK));
-    }
-
-    @TestTemplate
-    void testListTables_DatabaseNotExistException() throws 
DatabaseNotExistException {
-        String anyDatabase = "anyDatabase";
-        assertThatThrownBy(() -> catalog.listTables(anyDatabase))
-                .satisfies(
-                        anyCauseMatches(
-                                DatabaseNotExistException.class,
-                                String.format(
-                                        "Database %s does not exist in 
Catalog", anyDatabase)));
-    }
-
-    @TestTemplate
-    void testTableExists() {
-        String tableNotExist = "nonexist";
-        assertThat(catalog.tableExists(new ObjectPath(TEST_DB, 
tableNotExist))).isFalse();
-        assertThat(catalog.tableExists(new ObjectPath(TEST_DB, 
TEST_TABLE_ALL_TYPES))).isTrue();
-    }
-
-    @TestTemplate
-    void testGetTables_TableNotExistException() throws TableNotExistException {
-        String anyTableNotExist = "anyTable";
-        assertThatThrownBy(() -> catalog.getTable(new ObjectPath(TEST_DB, 
anyTableNotExist)))
-                .satisfies(
-                        anyCauseMatches(
-                                TableNotExistException.class,
-                                String.format(
-                                        "Table (or view) %s.%s does not exist 
in Catalog",
-                                        TEST_DB, anyTableNotExist)));
-    }
-
-    @TestTemplate
-    void testGetTables_TableNotExistException_NoDb() throws 
TableNotExistException {
-        String databaseNotExist = "nonexistdb";
-        String tableNotExist = "anyTable";
-        assertThatThrownBy(() -> catalog.getTable(new 
ObjectPath(databaseNotExist, tableNotExist)))
-                .satisfies(
-                        anyCauseMatches(
-                                TableNotExistException.class,
-                                String.format(
-                                        "Table (or view) %s.%s does not exist 
in Catalog",
-                                        databaseNotExist, tableNotExist)));
-    }
-
-    @TestTemplate
-    void testGetTable() throws TableNotExistException {
-        CatalogBaseTable table = catalog.getTable(new ObjectPath(TEST_DB, 
TEST_TABLE_ALL_TYPES));
-        assertThat(table.getUnresolvedSchema()).isEqualTo(TABLE_SCHEMA);
-    }
-
-    @TestTemplate
-    void testGetTablePrimaryKey() throws TableNotExistException {
-        // test the PK of test.t_user
-        Schema tableSchemaTestPK1 =
-                Schema.newBuilder()
-                        .column("uid", DataTypes.BIGINT().notNull())
-                        .column("col_bigint", DataTypes.BIGINT())
-                        .primaryKeyNamed("PRIMARY", 
Collections.singletonList("uid"))
-                        .build();
-        CatalogBaseTable tablePK1 = catalog.getTable(new ObjectPath(TEST_DB, 
TEST_TABLE_PK));
-        assertThat(tableSchemaTestPK1.getPrimaryKey().get())
-                
.isEqualTo(tablePK1.getUnresolvedSchema().getPrimaryKey().get());
-
-        // test the PK of test2.t_user
-        Schema tableSchemaTestPK2 =
-                Schema.newBuilder()
-                        .column("pid", DataTypes.INT().notNull())
-                        .column("col_varchar", DataTypes.VARCHAR(255))
-                        .primaryKeyNamed("PRIMARY", 
Collections.singletonList("pid"))
-                        .build();
-        CatalogBaseTable tablePK2 = catalog.getTable(new ObjectPath(TEST_DB2, 
TEST_TABLE_PK));
-        assertThat(tableSchemaTestPK2.getPrimaryKey().get())
-                
.isEqualTo(tablePK2.getUnresolvedSchema().getPrimaryKey().get());
-    }
-
-    // ------ test select query. ------
-
-    @TestTemplate
-    void testSelectField() {
-        List<Row> results =
-                CollectionUtil.iteratorToList(
-                        tEnv.sqlQuery(String.format("select pid from %s", 
TEST_TABLE_ALL_TYPES))
-                                .execute()
-                                .collect());
-        assertThat(results)
-                .isEqualTo(
-                        Lists.newArrayList(
-                                Row.ofKind(RowKind.INSERT, 1L), 
Row.ofKind(RowKind.INSERT, 2L)));
-    }
-
-    @TestTemplate
-    void testWithoutCatalogDB() {
-        List<Row> results =
-                CollectionUtil.iteratorToList(
-                        tEnv.sqlQuery(String.format("select * from %s", 
TEST_TABLE_ALL_TYPES))
-                                .execute()
-                                .collect());
-
-        assertThat(results).isEqualTo(ALL_TYPES_ROWS);
-    }
-
-    @TestTemplate
-    void testWithoutCatalog() {
-        List<Row> results =
-                CollectionUtil.iteratorToList(
-                        tEnv.sqlQuery(
-                                        String.format(
-                                                "select * from `%s`.`%s`",
-                                                TEST_DB, TEST_TABLE_ALL_TYPES))
-                                .execute()
-                                .collect());
-        assertThat(results).isEqualTo(ALL_TYPES_ROWS);
-    }
-
-    @TestTemplate
-    void testFullPath() {
-        List<Row> results =
-                CollectionUtil.iteratorToList(
-                        tEnv.sqlQuery(
-                                        String.format(
-                                                "select * from %s.%s.`%s`",
-                                                TEST_CATALOG_NAME,
-                                                catalog.getDefaultDatabase(),
-                                                TEST_TABLE_ALL_TYPES))
-                                .execute()
-                                .collect());
-        assertThat(results).isEqualTo(ALL_TYPES_ROWS);
-    }
-
-    @TestTemplate
-    void testSelectToInsert() throws Exception {
-
-        String sql =
-                String.format(
-                        "insert into `%s` select * from `%s`",
-                        TEST_SINK_TABLE_ALL_TYPES, TEST_TABLE_ALL_TYPES);
-        tEnv.executeSql(sql).await();
-
-        List<Row> results =
-                CollectionUtil.iteratorToList(
-                        tEnv.sqlQuery(String.format("select * from %s", 
TEST_SINK_TABLE_ALL_TYPES))
-                                .execute()
-                                .collect());
-        assertThat(results).isEqualTo(ALL_TYPES_ROWS);
-    }
-
-    @TestTemplate
-    void testGroupByInsert() throws Exception {
-        // Changes primary key for the next record.
-        tEnv.executeSql(
-                        String.format(
-                                "insert into `%s` select max(`pid`) `pid`, 
`col_bigint` from `%s` "
-                                        + "group by `col_bigint` ",
-                                TEST_TABLE_SINK_FROM_GROUPED_BY, 
TEST_TABLE_ALL_TYPES))
-                .await();
+    @Container
+    private static final MySQLContainer<?> CONTAINER = 
createContainer(DockerImageVersions.MYSQL);
 
-        List<Row> results =
-                CollectionUtil.iteratorToList(
-                        tEnv.sqlQuery(
-                                        String.format(
-                                                "select * from `%s`",
-                                                
TEST_TABLE_SINK_FROM_GROUPED_BY))
-                                .execute()
-                                .collect());
-        
assertThat(results).isEqualTo(Lists.newArrayList(Row.ofKind(RowKind.INSERT, 2L, 
-1L)));
+    @Override
+    protected String getDatabaseUrl() {
+        return CONTAINER.getJdbcUrl().substring(0, 
CONTAINER.getJdbcUrl().lastIndexOf("/"));
     }
 }
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogTestBase.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogTestBase.java
index c58dcf4..f8b9b07 100644
--- 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogTestBase.java
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/MySqlCatalogTestBase.java
@@ -19,31 +19,46 @@
 package org.apache.flink.connector.jdbc.catalog;
 
 import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.api.EnvironmentSettings;
 import org.apache.flink.table.api.Schema;
+import org.apache.flink.table.api.TableEnvironment;
+import org.apache.flink.table.catalog.CatalogBaseTable;
+import org.apache.flink.table.catalog.ObjectPath;
+import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException;
+import org.apache.flink.table.catalog.exceptions.TableNotExistException;
+import org.apache.flink.types.Row;
+import org.apache.flink.types.RowKind;
+import org.apache.flink.util.CollectionUtil;
 
 import org.apache.flink.shaded.guava30.com.google.common.collect.Lists;
 
-import org.junit.jupiter.api.AfterAll;
-import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.testcontainers.containers.MySQLContainer;
 import org.testcontainers.containers.output.Slf4jLogConsumer;
 import org.testcontainers.utility.DockerImageName;
 
-import java.sql.SQLException;
+import java.math.BigDecimal;
+import java.sql.Date;
+import java.sql.Time;
+import java.sql.Timestamp;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 
+import static org.apache.flink.core.testutils.FlinkAssertions.anyCauseMatches;
+import static 
org.apache.flink.table.api.config.ExecutionConfigOptions.TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
 /** Test base for {@link MySqlCatalog}. */
-class MySqlCatalogTestBase {
+abstract class MySqlCatalogTestBase {
 
     public static final Logger LOG = 
LoggerFactory.getLogger(MySqlCatalogTestBase.class);
-
-    protected static final List<String> DOCKER_IMAGE_NAMES =
-            Arrays.asList("mysql:5.6.51", "mysql:5.7.40", "mysql:8.0.31");
     protected static final String TEST_CATALOG_NAME = "mysql_catalog";
     protected static final String TEST_USERNAME = "mysql";
     protected static final String TEST_PWD = "mysql";
@@ -111,39 +126,338 @@ class MySqlCatalogTestBase {
                     .primaryKeyNamed("PRIMARY", Lists.newArrayList("pid"))
                     .build();
 
-    public static final Map<String, MySQLContainer<?>> MYSQL_CONTAINERS = new 
HashMap<>();
-    public static final Map<String, MySqlCatalog> CATALOGS = new HashMap<>();
-
-    @BeforeAll
-    static void beforeAll() throws SQLException {
-        for (String dockerImageName : DOCKER_IMAGE_NAMES) {
-            MySQLContainer<?> container =
-                    new 
MySQLContainer<>(DockerImageName.parse(dockerImageName))
-                            .withUsername("root")
-                            .withPassword("")
-                            .withEnv(DEFAULT_CONTAINER_ENV_MAP)
-                            .withInitScript(MYSQL_INIT_SCRIPT)
-                            .withLogConsumer(new Slf4jLogConsumer(LOG));
-            container.start();
-            MYSQL_CONTAINERS.put(dockerImageName, container);
-            CATALOGS.put(
-                    dockerImageName,
-                    new MySqlCatalog(
-                            Thread.currentThread().getContextClassLoader(),
-                            TEST_CATALOG_NAME,
-                            TEST_DB,
-                            TEST_USERNAME,
-                            TEST_PWD,
-                            container
-                                    .getJdbcUrl()
-                                    .substring(0, 
container.getJdbcUrl().lastIndexOf("/"))));
-        }
-    }
-
-    @AfterAll
-    static void cleanup() {
-        for (MySQLContainer<?> container : MYSQL_CONTAINERS.values()) {
-            container.stop();
-        }
+    protected static final List<Row> TABLE_ROWS =
+            Lists.newArrayList(
+                    Row.ofKind(
+                            RowKind.INSERT,
+                            1L,
+                            -1L,
+                            new BigDecimal(1),
+                            null,
+                            true,
+                            null,
+                            "hello",
+                            Date.valueOf("2021-08-04").toLocalDate(),
+                            Timestamp.valueOf("2021-08-04 
01:54:16").toLocalDateTime(),
+                            new BigDecimal(-1),
+                            new BigDecimal(1),
+                            -1.0d,
+                            1.0d,
+                            "enum2",
+                            -9.1f,
+                            9.1f,
+                            -1,
+                            1L,
+                            -1,
+                            1L,
+                            null,
+                            "col_longtext",
+                            null,
+                            -1,
+                            1,
+                            "col_mediumtext",
+                            new BigDecimal(-99),
+                            new BigDecimal(99),
+                            -1.0d,
+                            1.0d,
+                            "set_ele1",
+                            Short.parseShort("-1"),
+                            1,
+                            "col_text",
+                            Time.valueOf("10:32:34").toLocalTime(),
+                            Timestamp.valueOf("2021-08-04 
01:54:16").toLocalDateTime(),
+                            "col_tinytext",
+                            Byte.parseByte("-1"),
+                            Short.parseShort("1"),
+                            null,
+                            "col_varchar",
+                            Timestamp.valueOf("2021-08-04 
01:54:16.463").toLocalDateTime(),
+                            Time.valueOf("09:33:43").toLocalTime(),
+                            Timestamp.valueOf("2021-08-04 
01:54:16.463").toLocalDateTime(),
+                            null),
+                    Row.ofKind(
+                            RowKind.INSERT,
+                            2L,
+                            -1L,
+                            new BigDecimal(1),
+                            null,
+                            true,
+                            null,
+                            "hello",
+                            Date.valueOf("2021-08-04").toLocalDate(),
+                            Timestamp.valueOf("2021-08-04 
01:53:19").toLocalDateTime(),
+                            new BigDecimal(-1),
+                            new BigDecimal(1),
+                            -1.0d,
+                            1.0d,
+                            "enum2",
+                            -9.1f,
+                            9.1f,
+                            -1,
+                            1L,
+                            -1,
+                            1L,
+                            null,
+                            "col_longtext",
+                            null,
+                            -1,
+                            1,
+                            "col_mediumtext",
+                            new BigDecimal(-99),
+                            new BigDecimal(99),
+                            -1.0d,
+                            1.0d,
+                            "set_ele1,set_ele12",
+                            Short.parseShort("-1"),
+                            1,
+                            "col_text",
+                            Time.valueOf("10:32:34").toLocalTime(),
+                            Timestamp.valueOf("2021-08-04 
01:53:19").toLocalDateTime(),
+                            "col_tinytext",
+                            Byte.parseByte("-1"),
+                            Short.parseShort("1"),
+                            null,
+                            "col_varchar",
+                            Timestamp.valueOf("2021-08-04 
01:53:19.098").toLocalDateTime(),
+                            Time.valueOf("09:33:43").toLocalTime(),
+                            Timestamp.valueOf("2021-08-04 
01:53:19.098").toLocalDateTime(),
+                            null));
+
+    private MySqlCatalog catalog;
+    private TableEnvironment tEnv;
+
+    protected static MySQLContainer<?> createContainer(String dockerImage) {
+        return new MySQLContainer<>(DockerImageName.parse(dockerImage))
+                .withUsername("root")
+                .withPassword("")
+                .withEnv(DEFAULT_CONTAINER_ENV_MAP)
+                .withInitScript(MYSQL_INIT_SCRIPT)
+                .withLogConsumer(new Slf4jLogConsumer(LOG));
+    }
+
+    protected abstract String getDatabaseUrl();
+
+    @BeforeEach
+    void setup() {
+        catalog =
+                new MySqlCatalog(
+                        Thread.currentThread().getContextClassLoader(),
+                        TEST_CATALOG_NAME,
+                        TEST_DB,
+                        TEST_USERNAME,
+                        TEST_PWD,
+                        getDatabaseUrl());
+
+        this.tEnv = 
TableEnvironment.create(EnvironmentSettings.inStreamingMode());
+        tEnv.getConfig().set(TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM, 1);
+
+        // Use mysql catalog.
+        tEnv.registerCatalog(TEST_CATALOG_NAME, catalog);
+        tEnv.useCatalog(TEST_CATALOG_NAME);
+    }
+
+    @Test
+    void testGetDb_DatabaseNotExistException() throws Exception {
+        String databaseNotExist = "nonexistent";
+        assertThatThrownBy(() -> catalog.getDatabase(databaseNotExist))
+                .satisfies(
+                        anyCauseMatches(
+                                DatabaseNotExistException.class,
+                                String.format(
+                                        "Database %s does not exist in 
Catalog",
+                                        databaseNotExist)));
+    }
+
+    @Test
+    void testListDatabases() {
+        List<String> actual = catalog.listDatabases();
+        assertThat(actual).containsExactly(TEST_DB, TEST_DB2);
+    }
+
+    @Test
+    void testDbExists() throws Exception {
+        String databaseNotExist = "nonexistent";
+        assertThat(catalog.databaseExists(databaseNotExist)).isFalse();
+        assertThat(catalog.databaseExists(TEST_DB)).isTrue();
+    }
+
+    // ------ tables ------
+
+    @Test
+    void testListTables() throws DatabaseNotExistException {
+        List<String> actual = catalog.listTables(TEST_DB);
+        assertThat(actual)
+                .isEqualTo(
+                        Arrays.asList(
+                                TEST_TABLE_ALL_TYPES,
+                                TEST_SINK_TABLE_ALL_TYPES,
+                                TEST_TABLE_SINK_FROM_GROUPED_BY,
+                                TEST_TABLE_PK));
+    }
+
+    @Test
+    void testListTables_DatabaseNotExistException() throws 
DatabaseNotExistException {
+        String anyDatabase = "anyDatabase";
+        assertThatThrownBy(() -> catalog.listTables(anyDatabase))
+                .satisfies(
+                        anyCauseMatches(
+                                DatabaseNotExistException.class,
+                                String.format(
+                                        "Database %s does not exist in 
Catalog", anyDatabase)));
+    }
+
+    @Test
+    void testTableExists() {
+        String tableNotExist = "nonexist";
+        assertThat(catalog.tableExists(new ObjectPath(TEST_DB, 
tableNotExist))).isFalse();
+        assertThat(catalog.tableExists(new ObjectPath(TEST_DB, 
TEST_TABLE_ALL_TYPES))).isTrue();
+    }
+
+    @Test
+    void testGetTables_TableNotExistException() throws TableNotExistException {
+        String anyTableNotExist = "anyTable";
+        assertThatThrownBy(() -> catalog.getTable(new ObjectPath(TEST_DB, 
anyTableNotExist)))
+                .satisfies(
+                        anyCauseMatches(
+                                TableNotExistException.class,
+                                String.format(
+                                        "Table (or view) %s.%s does not exist 
in Catalog",
+                                        TEST_DB, anyTableNotExist)));
+    }
+
+    @Test
+    void testGetTables_TableNotExistException_NoDb() throws 
TableNotExistException {
+        String databaseNotExist = "nonexistdb";
+        String tableNotExist = "anyTable";
+        assertThatThrownBy(() -> catalog.getTable(new 
ObjectPath(databaseNotExist, tableNotExist)))
+                .satisfies(
+                        anyCauseMatches(
+                                TableNotExistException.class,
+                                String.format(
+                                        "Table (or view) %s.%s does not exist 
in Catalog",
+                                        databaseNotExist, tableNotExist)));
+    }
+
+    @Test
+    void testGetTable() throws TableNotExistException {
+        CatalogBaseTable table = catalog.getTable(new ObjectPath(TEST_DB, 
TEST_TABLE_ALL_TYPES));
+        assertThat(table.getUnresolvedSchema()).isEqualTo(TABLE_SCHEMA);
+    }
+
+    @Test
+    void testGetTablePrimaryKey() throws TableNotExistException {
+        // test the PK of test.t_user
+        Schema tableSchemaTestPK1 =
+                Schema.newBuilder()
+                        .column("uid", DataTypes.BIGINT().notNull())
+                        .column("col_bigint", DataTypes.BIGINT())
+                        .primaryKeyNamed("PRIMARY", 
Collections.singletonList("uid"))
+                        .build();
+        CatalogBaseTable tablePK1 = catalog.getTable(new ObjectPath(TEST_DB, 
TEST_TABLE_PK));
+        assertThat(tableSchemaTestPK1.getPrimaryKey().get())
+                
.isEqualTo(tablePK1.getUnresolvedSchema().getPrimaryKey().get());
+
+        // test the PK of test2.t_user
+        Schema tableSchemaTestPK2 =
+                Schema.newBuilder()
+                        .column("pid", DataTypes.INT().notNull())
+                        .column("col_varchar", DataTypes.VARCHAR(255))
+                        .primaryKeyNamed("PRIMARY", 
Collections.singletonList("pid"))
+                        .build();
+        CatalogBaseTable tablePK2 = catalog.getTable(new ObjectPath(TEST_DB2, 
TEST_TABLE_PK));
+        assertThat(tableSchemaTestPK2.getPrimaryKey().get())
+                
.isEqualTo(tablePK2.getUnresolvedSchema().getPrimaryKey().get());
+    }
+
+    // ------ test select query. ------
+
+    @Test
+    void testSelectField() {
+        List<Row> results =
+                CollectionUtil.iteratorToList(
+                        tEnv.sqlQuery(String.format("select pid from %s", 
TEST_TABLE_ALL_TYPES))
+                                .execute()
+                                .collect());
+        assertThat(results)
+                .isEqualTo(
+                        Lists.newArrayList(
+                                Row.ofKind(RowKind.INSERT, 1L), 
Row.ofKind(RowKind.INSERT, 2L)));
+    }
+
+    @Test
+    void testWithoutCatalogDB() {
+        List<Row> results =
+                CollectionUtil.iteratorToList(
+                        tEnv.sqlQuery(String.format("select * from %s", 
TEST_TABLE_ALL_TYPES))
+                                .execute()
+                                .collect());
+
+        assertThat(results).isEqualTo(TABLE_ROWS);
+    }
+
+    @Test
+    void testWithoutCatalog() {
+        List<Row> results =
+                CollectionUtil.iteratorToList(
+                        tEnv.sqlQuery(
+                                        String.format(
+                                                "select * from `%s`.`%s`",
+                                                TEST_DB, TEST_TABLE_ALL_TYPES))
+                                .execute()
+                                .collect());
+        assertThat(results).isEqualTo(TABLE_ROWS);
+    }
+
+    @Test
+    void testFullPath() {
+        List<Row> results =
+                CollectionUtil.iteratorToList(
+                        tEnv.sqlQuery(
+                                        String.format(
+                                                "select * from %s.%s.`%s`",
+                                                TEST_CATALOG_NAME,
+                                                catalog.getDefaultDatabase(),
+                                                TEST_TABLE_ALL_TYPES))
+                                .execute()
+                                .collect());
+        assertThat(results).isEqualTo(TABLE_ROWS);
+    }
+
+    @Test
+    void testSelectToInsert() throws Exception {
+
+        String sql =
+                String.format(
+                        "insert into `%s` select * from `%s`",
+                        TEST_SINK_TABLE_ALL_TYPES, TEST_TABLE_ALL_TYPES);
+        tEnv.executeSql(sql).await();
+
+        List<Row> results =
+                CollectionUtil.iteratorToList(
+                        tEnv.sqlQuery(String.format("select * from %s", 
TEST_SINK_TABLE_ALL_TYPES))
+                                .execute()
+                                .collect());
+        assertThat(results).isEqualTo(TABLE_ROWS);
+    }
+
+    @Test
+    void testGroupByInsert() throws Exception {
+        // Changes primary key for the next record.
+        tEnv.executeSql(
+                        String.format(
+                                "insert into `%s` select max(`pid`) `pid`, 
`col_bigint` from `%s` "
+                                        + "group by `col_bigint` ",
+                                TEST_TABLE_SINK_FROM_GROUPED_BY, 
TEST_TABLE_ALL_TYPES))
+                .await();
+
+        List<Row> results =
+                CollectionUtil.iteratorToList(
+                        tEnv.sqlQuery(
+                                        String.format(
+                                                "select * from `%s`",
+                                                
TEST_TABLE_SINK_FROM_GROUPED_BY))
+                                .execute()
+                                .collect());
+        
assertThat(results).isEqualTo(Lists.newArrayList(Row.ofKind(RowKind.INSERT, 2L, 
-1L)));
     }
 }
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSinkITCase.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSinkITCase.java
index 4c6dbb2..b40466b 100644
--- 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSinkITCase.java
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSinkITCase.java
@@ -48,8 +48,6 @@ import org.apache.flink.types.Row;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.condition.DisabledOnOs;
-import org.junit.jupiter.api.condition.OS;
 import org.testcontainers.containers.MSSQLServerContainer;
 import org.testcontainers.junit.jupiter.Container;
 import org.testcontainers.junit.jupiter.Testcontainers;
@@ -71,7 +69,6 @@ import static org.apache.flink.table.api.Expressions.$;
 import static 
org.apache.flink.table.factories.utils.FactoryMocks.createTableSink;
 
 /** The Table Sink ITCase for {@link SqlServerDialect}. */
-@DisabledOnOs(OS.MAC)
 @Testcontainers
 class SqlServerTableSinkITCase extends AbstractTestBase {
 
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSourceITCase.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSourceITCase.java
index bdb9a08..9abfe54 100644
--- 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSourceITCase.java
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/sqlserver/SqlServerTableSourceITCase.java
@@ -29,8 +29,6 @@ import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.condition.DisabledOnOs;
-import org.junit.jupiter.api.condition.OS;
 import org.testcontainers.containers.MSSQLServerContainer;
 
 import java.sql.Connection;
@@ -45,7 +43,6 @@ import java.util.stream.Stream;
 import static org.assertj.core.api.Assertions.assertThat;
 
 /** The Table Source ITCase for {@link SqlServerDialect}. */
-@DisabledOnOs(OS.MAC)
 class SqlServerTableSourceITCase extends AbstractTestBase {
 
     private static final MSSQLServerContainer container =
diff --git 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcExactlyOnceSinkE2eTest.java
 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcExactlyOnceSinkE2eTest.java
index f22f8a9..13e7875 100644
--- 
a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcExactlyOnceSinkE2eTest.java
+++ 
b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcExactlyOnceSinkE2eTest.java
@@ -22,8 +22,6 @@ import org.apache.flink.api.common.state.ListState;
 import org.apache.flink.api.common.state.ListStateDescriptor;
 import org.apache.flink.api.common.time.Time;
 import org.apache.flink.configuration.Configuration;
-import org.apache.flink.connector.jdbc.DbMetadata;
-import org.apache.flink.connector.jdbc.DerbyDbMetadata;
 import org.apache.flink.connector.jdbc.JdbcExactlyOnceOptions;
 import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
 import org.apache.flink.connector.jdbc.JdbcITCase;
@@ -84,12 +82,11 @@ public abstract class JdbcExactlyOnceSinkE2eTest extends 
JdbcTestBase {
     protected static final long CHECKPOINT_TIMEOUT_MS = 20_000L;
     protected static final long TASK_CANCELLATION_TIMEOUT_MS = 20_000L;
 
-    abstract protected SerializableSupplier<XADataSource> 
getDataSourceSupplier();
+    protected abstract SerializableSupplier<XADataSource> 
getDataSourceSupplier();
 
-    abstract protected String getDockerVersion();
+    protected abstract String getDockerVersion();
 
-    @RegisterExtension
-    static final MiniClusterExtension MINI_CLUSTER = createCluster();
+    @RegisterExtension static final MiniClusterExtension MINI_CLUSTER = 
createCluster();
 
     private static MiniClusterExtension createCluster() {
         Configuration configuration = new Configuration();
@@ -104,8 +101,7 @@ public abstract class JdbcExactlyOnceSinkE2eTest extends 
JdbcTestBase {
                 new MiniClusterResourceConfiguration.Builder()
                         .setNumberTaskManagers(PARALLELISM)
                         .setConfiguration(configuration)
-                        .build()
-        );
+                        .build());
     }
 
     // track active sources for:


Reply via email to