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:
