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 b0a25a688679595b0f2d128fb4b88758ce8629da Author: Joao Boto <[email protected]> AuthorDate: Thu Feb 2 09:12:37 2023 +0100 [FLINK-30790] Create unified databases for testing --- .../flink/connector/jdbc/JdbcTestFixture.java | 5 +- .../jdbc/catalog/MySql56CatalogITCase.java | 2 +- .../jdbc/catalog/MySql57CatalogITCase.java | 2 +- .../connector/jdbc/databases/DatabaseMetadata.java | 40 ++++++++++++++ .../connector/jdbc/databases/DatabaseTest.java | 24 +++++++++ .../jdbc/databases/derby/DerbyDatabase.java | 50 +++++++++++++++++ .../derby/DerbyMetadata.java} | 31 +++++++---- .../h2/H2Metadata.java} | 47 ++++++++-------- .../connector/jdbc/databases/h2/H2XaDatabase.java | 50 +++++++++++++++++ .../mysql/MySqlDatabase.java} | 62 ++++++++++------------ .../mysql/MySqlMetadata.java | 37 +++++++------ .../jdbc/databases/oracle/OracleDatabase.java | 38 +++++++++++++ .../oracle/OracleMetadata.java | 36 +++++++------ .../jdbc/databases/postgres/PostgresDatabase.java | 43 +++++++++++++++ .../postgres/PostgresMetadata.java | 34 +++++++----- .../databases/sqlserver/SqlServerDatabase.java | 40 ++++++++++++++ .../sqlserver/SqlServerMetadata.java} | 38 +++++++------ .../dialect/mysql/MySqlExactlyOnceSinkE2eTest.java | 1 + .../oracle/OracleExactlyOnceSinkE2eTest.java | 1 + .../postgres/PostgresExactlyOnceSinkE2eTest.java | 1 + .../flink/connector/jdbc/xa/h2/H2XaDsWrapper.java | 2 +- 21 files changed, 452 insertions(+), 132 deletions(-) diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcTestFixture.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcTestFixture.java index e9b2953..d00fbad 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcTestFixture.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcTestFixture.java @@ -20,6 +20,7 @@ package org.apache.flink.connector.jdbc; import org.apache.flink.api.common.typeinfo.BasicTypeInfo; import org.apache.flink.api.java.typeutils.RowTypeInfo; +import org.apache.flink.connector.jdbc.databases.derby.DerbyMetadata; import org.apache.flink.connector.jdbc.xa.h2.H2DbMetadata; import org.apache.flink.table.types.logical.RowType; @@ -74,8 +75,8 @@ public class JdbcTestFixture { }; private static final String EBOOKSHOP_SCHEMA_NAME = "ebookshop"; - public static final DerbyDbMetadata DERBY_EBOOKSHOP_DB = - new DerbyDbMetadata(EBOOKSHOP_SCHEMA_NAME); + public static final DerbyMetadata DERBY_EBOOKSHOP_DB = + new DerbyMetadata(EBOOKSHOP_SCHEMA_NAME); public static final H2DbMetadata H2_EBOOKSHOP_DB = new H2DbMetadata(EBOOKSHOP_SCHEMA_NAME); /** TestEntry. */ 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 index 2ff3ed0..3a1c554 100644 --- 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 @@ -24,7 +24,7 @@ import org.testcontainers.containers.MySQLContainer; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; -/** E2E test for {@link MySqlCatalog}. */ +/** E2E test for {@link MySqlCatalog} with MySql version 5.6. */ @Testcontainers public class MySql56CatalogITCase extends MySqlCatalogTestBase { 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 index 0a1dc8b..350bea8 100644 --- 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 @@ -24,7 +24,7 @@ import org.testcontainers.containers.MySQLContainer; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; -/** E2E test for {@link MySqlCatalog}. */ +/** E2E test for {@link MySqlCatalog} with MySql version 5.7. */ @Testcontainers public class MySql57CatalogITCase extends MySqlCatalogTestBase { diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/DatabaseMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/DatabaseMetadata.java new file mode 100644 index 0000000..fcfc656 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/DatabaseMetadata.java @@ -0,0 +1,40 @@ +package org.apache.flink.connector.jdbc.databases; + +import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.JdbcConnectionOptions; + +import javax.sql.XADataSource; + +import java.io.Serializable; + +public interface DatabaseMetadata extends Serializable, DbMetadata { + + default String getUrl(){ + return getJdbcUrl(); + } + + default String getUser() { + return getUsername(); + } + + String getJdbcUrl(); + + String getUsername(); + + String getPassword(); + + XADataSource buildXaDataSource(); + + String getDriverClass(); + + String getVersion(); + + default JdbcConnectionOptions toConnectionOptions() { + return new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() + .withDriverName(getDriverClass()) + .withUrl(getJdbcUrl()) + .withUsername(getUsername()) + .withPassword(getPassword()) + .build(); + } +} \ No newline at end of file diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/DatabaseTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/DatabaseTest.java new file mode 100644 index 0000000..802a468 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/DatabaseTest.java @@ -0,0 +1,24 @@ +/* + * 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.databases; + +/** Base interface for tests that have dependency in a database. */ +public interface DatabaseTest { + + DatabaseMetadata getMetadata(); +} diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyDatabase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyDatabase.java new file mode 100644 index 0000000..f0f6e3a --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyDatabase.java @@ -0,0 +1,50 @@ +package org.apache.flink.connector.jdbc.databases.derby; + +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseTest; +import org.apache.flink.util.FlinkRuntimeException; + +import java.io.OutputStream; +import java.sql.DriverManager; +import java.sql.SQLException; + +/** Derby database for testing. * */ +public interface DerbyDatabase extends DatabaseTest { + + @SuppressWarnings("unused") // used in string constant in prepareDatabase + OutputStream DEV_NULL = + new OutputStream() { + @Override + public void write(int b) {} + }; + + DatabaseMetadata METADATA = startDatabase(); + + @Override + default DatabaseMetadata getMetadata() { + return METADATA; + } + + static DatabaseMetadata startDatabase() { + DatabaseMetadata metadata = new DerbyMetadata("test"); + try { + System.setProperty( + "derby.stream.error.field", + DerbyDatabase.class.getCanonicalName() + ".DEV_NULL"); + Class.forName(metadata.getDriverClass()); + DriverManager.getConnection(String.format("%s;create=true", metadata.getJdbcUrl())).close(); + } catch (Exception e) { + throw new FlinkRuntimeException(e); + } + return metadata; + } + + default void stopDatabase() throws Exception { + DatabaseMetadata metadata = getMetadata(); + try { + DriverManager.getConnection(String.format("%s;shutdown=true", metadata.getJdbcUrl())) + .close(); + } catch (SQLException ignored) { + } + } +} diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DerbyDbMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java similarity index 72% copy from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DerbyDbMetadata.java copy to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java index 85e4123..f7cc8fb 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DerbyDbMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java @@ -15,22 +15,20 @@ * limitations under the License. */ -package org.apache.flink.connector.jdbc; +package org.apache.flink.connector.jdbc.databases.derby; + +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.apache.derby.jdbc.EmbeddedXADataSource; import javax.sql.XADataSource; /** DerbyDbMetadata. */ -public class DerbyDbMetadata implements DbMetadata { +public class DerbyMetadata implements DatabaseMetadata { private final String dbName; - private final String dbInitUrl; - private final String url; - public DerbyDbMetadata(String schemaName) { + public DerbyMetadata(String schemaName) { dbName = "memory:" + schemaName; - url = "jdbc:derby:" + dbName; - dbInitUrl = url + ";create=true"; } public String getDbName() { @@ -38,8 +36,18 @@ public class DerbyDbMetadata implements DbMetadata { } @Override - public String getInitUrl() { - return dbInitUrl; + public String getJdbcUrl() { + return String.format("jdbc:derby:%s", dbName); + } + + @Override + public String getUsername() { + return ""; + } + + @Override + public String getPassword() { + return ""; } @Override @@ -55,7 +63,8 @@ public class DerbyDbMetadata implements DbMetadata { } @Override - public String getUrl() { - return url; + public String getVersion() { + return "derby:memory"; } + } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DerbyDbMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java similarity index 55% rename from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DerbyDbMetadata.java rename to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java index 85e4123..027fb53 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DerbyDbMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java @@ -15,47 +15,52 @@ * limitations under the License. */ -package org.apache.flink.connector.jdbc; +package org.apache.flink.connector.jdbc.databases.h2; -import org.apache.derby.jdbc.EmbeddedXADataSource; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.xa.h2.H2XaDsWrapper; import javax.sql.XADataSource; -/** DerbyDbMetadata. */ -public class DerbyDbMetadata implements DbMetadata { - private final String dbName; - private final String dbInitUrl; - private final String url; +/** H2DbMetadata. */ +public class H2Metadata implements DatabaseMetadata { - public DerbyDbMetadata(String schemaName) { - dbName = "memory:" + schemaName; - url = "jdbc:derby:" + dbName; - dbInitUrl = url + ";create=true"; + private final String schema; + + public H2Metadata(String schema) { + this.schema = schema; + } + + @Override + public String getJdbcUrl() { + return String.format("jdbc:h2:mem:%s", schema); } - public String getDbName() { - return dbName; + @Override + public String getUsername() { + return ""; } @Override - public String getInitUrl() { - return dbInitUrl; + public String getPassword() { + return ""; } @Override public XADataSource buildXaDataSource() { - EmbeddedXADataSource ds = new EmbeddedXADataSource(); - ds.setDatabaseName(dbName); - return ds; + final org.h2.jdbcx.JdbcDataSource ds = new org.h2.jdbcx.JdbcDataSource(); + ds.setUrl(getJdbcUrl()); + return new H2XaDsWrapper(ds); } @Override public String getDriverClass() { - return "org.apache.derby.jdbc.EmbeddedDriver"; + return "org.h2.Driver"; } @Override - public String getUrl() { - return url; + public String getVersion() { + return "h2:mem"; } + } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2XaDatabase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2XaDatabase.java new file mode 100644 index 0000000..3b49a57 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2XaDatabase.java @@ -0,0 +1,50 @@ +/* + * 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.databases.h2; + +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseTest; +import org.apache.flink.util.FlinkRuntimeException; + +import java.sql.DriverManager; + +/** H2 database for testing. * */ +public interface H2XaDatabase extends DatabaseTest { + + DatabaseMetadata METADATA = startDatabase(); + + @Override + default DatabaseMetadata getMetadata() { + return METADATA; + } + + static DatabaseMetadata startDatabase() { + DatabaseMetadata metadata = new H2Metadata("test"); + try { + Class.forName(metadata.getDriverClass()); + DriverManager.getConnection( + String.format( + "%s;DB_CLOSE_DELAY=-1;INIT=CREATE SCHEMA IF NOT EXISTS %s\\;SET SCHEMA %s", + metadata.getJdbcUrl(), "test", "test")) + .close(); + } catch (Exception e) { + throw new FlinkRuntimeException(e); + } + return metadata; + } +} diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java similarity index 83% copy from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java copy to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java index 1da2f7c..f8e70a9 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java @@ -1,12 +1,27 @@ -package org.apache.flink.connector.jdbc.dialect.mysql; +/* + * 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.databases.mysql; -import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseTest; import org.apache.flink.connector.jdbc.test.DockerImageVersions; -import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.ExceptionUtils; -import org.apache.flink.util.function.SerializableSupplier; -import com.mysql.cj.jdbc.MysqlXADataSource; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testcontainers.containers.MySQLContainer; @@ -14,8 +29,6 @@ import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; import org.testcontainers.utility.DockerImageName; -import javax.sql.XADataSource; - import java.sql.Connection; import java.sql.DriverManager; import java.sql.ResultSet; @@ -24,42 +37,21 @@ import java.sql.Statement; import static org.apache.flink.util.Preconditions.checkArgument; -/** - * A simple end-to-end test for {@link JdbcExactlyOnceSinkE2eTest}. Check for issues with errors on - * closing connections. - */ +/** A MySql database for testing. * */ @Testcontainers -public class MySqlExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest { +public interface MySqlDatabase extends DatabaseTest { @Container - private static final MySqlXaContainer CONTAINER = - new MySqlXaContainer(DockerImageVersions.MYSQL) - .withLockWaitTimeout( - (CHECKPOINT_TIMEOUT_MS + TASK_CANCELLATION_TIMEOUT_MS) * 2); - - @Override - protected String getDockerVersion() { - return CONTAINER.getDockerImageName(); - } + MySqlXaContainer CONTAINER = + new MySqlXaContainer(DockerImageVersions.MYSQL).withLockWaitTimeout(50_000L); @Override - protected DbMetadata getDbMetadata() { + default DatabaseMetadata getMetadata() { return new MySqlMetadata(CONTAINER); } - @Override - public SerializableSupplier<XADataSource> getDataSourceSupplier() { - return () -> { - MysqlXADataSource xaDataSource = new MysqlXADataSource(); - xaDataSource.setUrl(CONTAINER.getJdbcUrl()); - xaDataSource.setUser(CONTAINER.getUsername()); - xaDataSource.setPassword(CONTAINER.getPassword()); - return xaDataSource; - }; - } - /** {@link MySQLContainer} with XA enabled. */ - static class MySqlXaContainer extends MySQLContainer<MySqlXaContainer> { + class MySqlXaContainer extends MySQLContainer<MySqlXaContainer> { private long lockWaitTimeout = 0; private volatile InnoDbStatusLogger innoDbStatusLogger; @@ -116,7 +108,7 @@ public class MySqlExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest { } /** InnoDB status logger. */ - static class InnoDbStatusLogger { + class InnoDbStatusLogger { private static final Logger LOG = LoggerFactory.getLogger(InnoDbStatusLogger.class); private final Thread thread; private volatile boolean running; diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlMetadata.java similarity index 77% rename from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlMetadata.java rename to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlMetadata.java index 0d626a3..4bb3d5a 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlMetadata.java @@ -15,17 +15,17 @@ * limitations under the License. */ -package org.apache.flink.connector.jdbc.dialect.mysql; +package org.apache.flink.connector.jdbc.databases.mysql; -import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import com.mysql.cj.jdbc.MysqlXADataSource; import org.testcontainers.containers.MySQLContainer; import javax.sql.XADataSource; -/** Postgres Metadata. */ -public class MySqlMetadata implements DbMetadata { +/** MySql Metadata. */ +public class MySqlMetadata implements DatabaseMetadata { private final String username; private final String password; @@ -34,11 +34,11 @@ public class MySqlMetadata implements DbMetadata { private final String version; private final boolean xaEnabled; - protected MySqlMetadata(MySQLContainer<?> container) { + public MySqlMetadata(MySQLContainer<?> container) { this(container, false); } - protected MySqlMetadata(MySQLContainer<?> container, boolean hasXaEnabled) { + public MySqlMetadata(MySQLContainer<?> container, boolean hasXaEnabled) { this.username = container.getUsername(); this.password = container.getPassword(); this.url = container.getJdbcUrl(); @@ -48,10 +48,20 @@ public class MySqlMetadata implements DbMetadata { } @Override - public String getUrl() { + public String getJdbcUrl() { return this.url; } + @Override + public String getUsername() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } + @Override public XADataSource buildXaDataSource() { if (!xaEnabled) { @@ -59,8 +69,8 @@ public class MySqlMetadata implements DbMetadata { } MysqlXADataSource xaDataSource = new MysqlXADataSource(); - xaDataSource.setUrl(getUrl()); - xaDataSource.setUser(getUser()); + xaDataSource.setUrl(getJdbcUrl()); + xaDataSource.setUser(getUsername()); xaDataSource.setPassword(getPassword()); return xaDataSource; } @@ -71,12 +81,9 @@ public class MySqlMetadata implements DbMetadata { } @Override - public String getUser() { - return this.username; + public String getVersion() { + return this.version; } - @Override - public String getPassword() { - return this.password; - } + } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleDatabase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleDatabase.java new file mode 100644 index 0000000..d800e38 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleDatabase.java @@ -0,0 +1,38 @@ +/* + * 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.databases.oracle; + +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseTest; +import org.apache.flink.connector.jdbc.dialect.oracle.OracleContainer; + +import org.testcontainers.containers.JdbcDatabaseContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; + +/** A Oracle database for testing. * */ +@Testcontainers +public interface OracleDatabase extends DatabaseTest { + + @Container JdbcDatabaseContainer<?> CONTAINER = new OracleContainer(); + + @Override + default DatabaseMetadata getMetadata() { + return new OracleMetadata(CONTAINER); + } +} diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java similarity index 78% rename from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleMetadata.java rename to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java index 8487c89..be2c70f 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java @@ -15,9 +15,9 @@ * limitations under the License. */ -package org.apache.flink.connector.jdbc.dialect.oracle; +package org.apache.flink.connector.jdbc.databases.oracle; -import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import oracle.jdbc.xa.client.OracleXADataSource; import org.testcontainers.containers.JdbcDatabaseContainer; @@ -26,8 +26,8 @@ import javax.sql.XADataSource; import java.sql.SQLException; -/** Postgres Metadata. */ -public class OracleMetadata implements DbMetadata { +/** Oracle Metadata. */ +public class OracleMetadata implements DatabaseMetadata { private final String username; private final String password; @@ -36,11 +36,11 @@ public class OracleMetadata implements DbMetadata { private final String version; private final boolean xaEnabled; - protected OracleMetadata(JdbcDatabaseContainer<?> container) { + public OracleMetadata(JdbcDatabaseContainer<?> container) { this(container, false); } - protected OracleMetadata(JdbcDatabaseContainer<?> container, boolean hasXaEnabled) { + public OracleMetadata(JdbcDatabaseContainer<?> container, boolean hasXaEnabled) { this.username = container.getUsername(); this.password = container.getPassword(); this.url = container.getJdbcUrl(); @@ -50,10 +50,20 @@ public class OracleMetadata implements DbMetadata { } @Override - public String getUrl() { + public String getJdbcUrl() { return this.url; } + @Override + public String getUsername() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } + @Override public XADataSource buildXaDataSource() { if (!xaEnabled) { @@ -61,8 +71,8 @@ public class OracleMetadata implements DbMetadata { } try { OracleXADataSource xaDataSource = new OracleXADataSource(); - xaDataSource.setURL(getUrl()); - xaDataSource.setUser(getUser()); + xaDataSource.setURL(getJdbcUrl()); + xaDataSource.setUser(getUsername()); xaDataSource.setPassword(getPassword()); return xaDataSource; } catch (SQLException e) { @@ -76,12 +86,8 @@ public class OracleMetadata implements DbMetadata { } @Override - public String getUser() { - return this.username; + public String getVersion() { + return version; } - @Override - public String getPassword() { - return this.password; - } } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/postgres/PostgresDatabase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/postgres/PostgresDatabase.java new file mode 100644 index 0000000..0db1ec7 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/postgres/PostgresDatabase.java @@ -0,0 +1,43 @@ +/* + * 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.databases.postgres; + +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseTest; +import org.apache.flink.connector.jdbc.dialect.postgres.PostgresExactlyOnceSinkE2eTest; +import org.apache.flink.connector.jdbc.test.DockerImageVersions; + +import org.testcontainers.containers.PostgreSQLContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; + +/** A Postgres database for testing. * */ +@Testcontainers +public interface PostgresDatabase extends DatabaseTest { + + @Container + PostgreSQLContainer<?> CONTAINER = + new PostgresExactlyOnceSinkE2eTest.PostgresXaContainer(DockerImageVersions.POSTGRES) + .withMaxConnections(10) + .withMaxTransactions(50); + + @Override + default DatabaseMetadata getMetadata() { + return new PostgresMetadata(CONTAINER); + } +} diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/postgres/PostgresMetadata.java similarity index 78% copy from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java copy to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/postgres/PostgresMetadata.java index 8e171da..e3d6c04 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/postgres/PostgresMetadata.java @@ -15,9 +15,9 @@ * limitations under the License. */ -package org.apache.flink.connector.jdbc.dialect.postgres; +package org.apache.flink.connector.jdbc.databases.postgres; -import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.postgresql.xa.PGXADataSource; import org.testcontainers.containers.PostgreSQLContainer; @@ -25,7 +25,7 @@ import org.testcontainers.containers.PostgreSQLContainer; import javax.sql.XADataSource; /** Postgres Metadata. */ -public class PostgresMetadata implements DbMetadata { +public class PostgresMetadata implements DatabaseMetadata { private final String username; private final String password; @@ -34,11 +34,11 @@ public class PostgresMetadata implements DbMetadata { private final String version; private final boolean xaEnabled; - protected PostgresMetadata(PostgreSQLContainer<?> container) { + public PostgresMetadata(PostgreSQLContainer<?> container) { this(container, false); } - protected PostgresMetadata(PostgreSQLContainer<?> container, boolean hasXaEnabled) { + public PostgresMetadata(PostgreSQLContainer<?> container, boolean hasXaEnabled) { this.username = container.getUsername(); this.password = container.getPassword(); this.url = container.getJdbcUrl(); @@ -48,10 +48,20 @@ public class PostgresMetadata implements DbMetadata { } @Override - public String getUrl() { + public String getJdbcUrl() { return this.url; } + @Override + public String getUsername() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } + @Override public XADataSource buildXaDataSource() { if (!xaEnabled) { @@ -59,8 +69,8 @@ public class PostgresMetadata implements DbMetadata { } PGXADataSource xaDataSource = new PGXADataSource(); - xaDataSource.setUrl(getUrl()); - xaDataSource.setUser(getUser()); + xaDataSource.setUrl(getJdbcUrl()); + xaDataSource.setUser(getUsername()); xaDataSource.setPassword(getPassword()); return xaDataSource; } @@ -71,12 +81,8 @@ public class PostgresMetadata implements DbMetadata { } @Override - public String getUser() { - return this.username; + public String getVersion() { + return version; } - @Override - public String getPassword() { - return this.password; - } } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/sqlserver/SqlServerDatabase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/sqlserver/SqlServerDatabase.java new file mode 100644 index 0000000..9f01f5c --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/sqlserver/SqlServerDatabase.java @@ -0,0 +1,40 @@ +/* + * 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.databases.sqlserver; + +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseTest; +import org.apache.flink.connector.jdbc.test.DockerImageVersions; + +import org.testcontainers.containers.MSSQLServerContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; + +/** A SqlServer database for testing. * */ +@Testcontainers +public interface SqlServerDatabase extends DatabaseTest { + + @Container + MSSQLServerContainer<?> CONTAINER = + new MSSQLServerContainer<>(DockerImageVersions.MSSQL_SERVER).acceptLicense(); + + @Override + default DatabaseMetadata getMetadata() { + return new SqlServerMetadata(CONTAINER); + } +} diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/sqlserver/SqlServerMetadata.java similarity index 74% rename from flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java rename to flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/sqlserver/SqlServerMetadata.java index 8e171da..5435746 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/sqlserver/SqlServerMetadata.java @@ -15,17 +15,17 @@ * limitations under the License. */ -package org.apache.flink.connector.jdbc.dialect.postgres; +package org.apache.flink.connector.jdbc.databases.sqlserver; -import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.postgresql.xa.PGXADataSource; -import org.testcontainers.containers.PostgreSQLContainer; +import org.testcontainers.containers.MSSQLServerContainer; import javax.sql.XADataSource; -/** Postgres Metadata. */ -public class PostgresMetadata implements DbMetadata { +/** SqlServer Metadata. */ +public class SqlServerMetadata implements DatabaseMetadata { private final String username; private final String password; @@ -34,11 +34,11 @@ public class PostgresMetadata implements DbMetadata { private final String version; private final boolean xaEnabled; - protected PostgresMetadata(PostgreSQLContainer<?> container) { + protected SqlServerMetadata(MSSQLServerContainer<?> container) { this(container, false); } - protected PostgresMetadata(PostgreSQLContainer<?> container, boolean hasXaEnabled) { + protected SqlServerMetadata(MSSQLServerContainer<?> container, boolean hasXaEnabled) { this.username = container.getUsername(); this.password = container.getPassword(); this.url = container.getJdbcUrl(); @@ -48,10 +48,20 @@ public class PostgresMetadata implements DbMetadata { } @Override - public String getUrl() { + public String getJdbcUrl() { return this.url; } + @Override + public String getUsername() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } + @Override public XADataSource buildXaDataSource() { if (!xaEnabled) { @@ -59,8 +69,8 @@ public class PostgresMetadata implements DbMetadata { } PGXADataSource xaDataSource = new PGXADataSource(); - xaDataSource.setUrl(getUrl()); - xaDataSource.setUser(getUser()); + xaDataSource.setUrl(getJdbcUrl()); + xaDataSource.setUser(getUsername()); xaDataSource.setPassword(getPassword()); return xaDataSource; } @@ -71,12 +81,8 @@ public class PostgresMetadata implements DbMetadata { } @Override - public String getUser() { - return this.username; + public String getVersion() { + return version; } - @Override - public String getPassword() { - return this.password; - } } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java index 1da2f7c..d03c928 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java @@ -1,6 +1,7 @@ package org.apache.flink.connector.jdbc.dialect.mysql; import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.mysql.MySqlMetadata; import org.apache.flink.connector.jdbc.test.DockerImageVersions; import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.ExceptionUtils; diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleExactlyOnceSinkE2eTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleExactlyOnceSinkE2eTest.java index 857568e..b3dce62 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleExactlyOnceSinkE2eTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleExactlyOnceSinkE2eTest.java @@ -1,6 +1,7 @@ package org.apache.flink.connector.jdbc.dialect.oracle; import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.oracle.OracleMetadata; import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.function.SerializableSupplier; diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java index f8634a3..5001f14 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java @@ -1,6 +1,7 @@ package org.apache.flink.connector.jdbc.dialect.postgres; import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.databases.postgres.PostgresMetadata; import org.apache.flink.connector.jdbc.test.DockerImageVersions; import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.function.SerializableSupplier; diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/h2/H2XaDsWrapper.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/h2/H2XaDsWrapper.java index 0d70727..9d348d9 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/h2/H2XaDsWrapper.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/h2/H2XaDsWrapper.java @@ -33,7 +33,7 @@ public class H2XaDsWrapper implements XADataSource { private final XADataSource wrapped; - H2XaDsWrapper(XADataSource wrapped) { + public H2XaDsWrapper(XADataSource wrapped) { this.wrapped = wrapped; }
