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 ba9dbc91a6dc19a49683dc6b558e46be61fd1c88 Author: Joao Boto <[email protected]> AuthorDate: Mon Jan 30 15:01:22 2023 +0100 [FLINK-30790] Refactor JdbcExactlyOnceSinkE2eTest create tests by dialect --- .../apache/flink/connector/jdbc/DbMetadata.java | 4 +- .../dialect/mysql/MySqlExactlyOnceSinkE2eTest.java | 220 +++++++++ .../jdbc/dialect/mysql/MySqlMetadata.java | 82 ++++ .../oracle/OracleExactlyOnceSinkE2eTest.java | 46 ++ .../jdbc/dialect/oracle/OracleMetadata.java | 87 ++++ .../postgres/PostgresExactlyOnceSinkE2eTest.java | 91 ++++ .../jdbc/dialect/postgres/PostgresMetadata.java | 82 ++++ .../sqlserver/SqlServerTableSinkITCase.java | 3 + .../sqlserver/SqlServerTableSourceITCase.java | 3 + .../connector/jdbc/test/DockerImageVersions.java | 19 +- .../jdbc/xa/JdbcExactlyOnceSinkE2eTest.java | 541 ++------------------- 11 files changed, 671 insertions(+), 507 deletions(-) diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DbMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DbMetadata.java index 21f93d1..55c5317 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DbMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/DbMetadata.java @@ -24,7 +24,9 @@ import java.io.Serializable; /** Describes a database: driver, schema and urls. */ public interface DbMetadata extends Serializable { - String getInitUrl(); + default String getInitUrl() { + return getUrl(); + } String getUrl(); 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 new file mode 100644 index 0000000..1da2f7c --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java @@ -0,0 +1,220 @@ +package org.apache.flink.connector.jdbc.dialect.mysql; + +import org.apache.flink.connector.jdbc.DbMetadata; +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; +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; +import java.sql.SQLException; +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. + */ +@Testcontainers +public class MySqlExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest { + + @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(); + } + + @Override + protected DbMetadata getDbMetadata() { + 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> { + private long lockWaitTimeout = 0; + private volatile InnoDbStatusLogger innoDbStatusLogger; + + public MySqlXaContainer(String dockerImageName) { + super(DockerImageName.parse(dockerImageName)); + } + + public MySqlXaContainer withLockWaitTimeout(long lockWaitTimeout) { + checkArgument(lockWaitTimeout >= 0, "lockWaitTimeout should be greater than 0"); + this.lockWaitTimeout = lockWaitTimeout; + return this.self(); + } + + @Override + public void start() { + super.start(); + // prevent XAER_RMERR: Fatal error occurred in the transaction branch - check your + // data for consistency works for mysql v8+ + try (Connection connection = + DriverManager.getConnection(getJdbcUrl(), "root", getPassword())) { + prepareDb(connection, lockWaitTimeout); + } catch (SQLException e) { + ExceptionUtils.rethrow(e); + } + + this.innoDbStatusLogger = + new InnoDbStatusLogger( + getJdbcUrl(), "root", getPassword(), lockWaitTimeout / 2); + innoDbStatusLogger.start(); + } + + @Override + public void stop() { + try { + innoDbStatusLogger.stop(); + } catch (Exception e) { + ExceptionUtils.rethrow(e); + } finally { + super.stop(); + } + } + + private void prepareDb(Connection connection, long lockWaitTimeout) throws SQLException { + try (Statement st = connection.createStatement()) { + st.execute("GRANT XA_RECOVER_ADMIN ON *.* TO '" + getUsername() + "'@'%'"); + st.execute("FLUSH PRIVILEGES"); + // if the reason of task cancellation failure is waiting for a lock + // then failing transactions with a relevant message would ease debugging + st.execute("SET GLOBAL innodb_lock_wait_timeout = " + lockWaitTimeout); + // st.execute("SET GLOBAL innodb_status_output = ON"); + // st.execute("SET GLOBAL innodb_status_output_locks = ON"); + } + } + } + + /** InnoDB status logger. */ + static class InnoDbStatusLogger { + private static final Logger LOG = LoggerFactory.getLogger(InnoDbStatusLogger.class); + private final Thread thread; + private volatile boolean running; + + private InnoDbStatusLogger(String url, String user, String password, long intervalMs) { + running = true; + thread = + new Thread( + () -> { + LOG.info("Logging InnoDB status every {}ms", intervalMs); + try (Connection connection = + DriverManager.getConnection(url, user, password)) { + while (running) { + Thread.sleep(intervalMs); + queryAndLog(connection); + } + } catch (Exception e) { + LOG.warn("failed", e); + } finally { + LOG.info("Logging InnoDB status stopped"); + } + }); + } + + public void start() { + thread.start(); + } + + public void stop() throws InterruptedException { + running = false; + thread.join(); + } + + private void queryAndLog(Connection connection) throws SQLException { + try (Statement st = connection.createStatement()) { + showBlockedTrx(st); + showAllTrx(st); + showEngineStatus(st); + showRecoveredTrx(st); + // additional query: show full processlist \G; -- only shows live + } + } + + private void showRecoveredTrx(Statement st) throws SQLException { + try (ResultSet rs = st.executeQuery("xa recover convert xid ")) { + while (rs.next()) { + LOG.debug( + "recovered trx: {} {} {} {}", + rs.getString(1), + rs.getString(2), + rs.getString(3), + rs.getString(4)); + } + } + } + + private void showEngineStatus(Statement st) throws SQLException { + LOG.debug("Engine status"); + try (ResultSet rs = st.executeQuery("show engine innodb status")) { + while (rs.next()) { + LOG.debug(rs.getString(3)); + } + } + } + + private void showAllTrx(Statement st) throws SQLException { + LOG.debug("All TRX"); + try (ResultSet rs = st.executeQuery("select * from information_schema.innodb_trx")) { + while (rs.next()) { + LOG.debug( + "trx_id: {}, trx_state: {}, trx_started: {}, trx_requested_lock_id: {}, trx_wait_started: {}, trx_mysql_thread_id: {},", + rs.getString("trx_id"), + rs.getString("trx_state"), + rs.getString("trx_started"), + rs.getString("trx_requested_lock_id"), + rs.getString("trx_wait_started"), + rs.getString("trx_mysql_thread_id") /* 0 for recovered*/); + } + } + } + + private void showBlockedTrx(Statement st) throws SQLException { + LOG.debug("Blocked TRX"); + try (ResultSet rs = + st.executeQuery( + " SELECT waiting_trx_id, waiting_pid, waiting_query, blocking_trx_id, blocking_pid, blocking_query " + + "FROM sys.innodb_lock_waits; ")) { + while (rs.next()) { + LOG.debug( + "waiting_trx_id: {}, waiting_pid: {}, waiting_query: {}, blocking_trx_id: {}, blocking_pid: {}, blocking_query: {}", + rs.getString(1), + rs.getString(2), + rs.getString(3), + rs.getString(4), + rs.getString(5), + rs.getString(6)); + } + } + } + } +} 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/dialect/mysql/MySqlMetadata.java new file mode 100644 index 0000000..0d626a3 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/mysql/MySqlMetadata.java @@ -0,0 +1,82 @@ +/* + * 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.dialect.mysql; + +import org.apache.flink.connector.jdbc.DbMetadata; + +import com.mysql.cj.jdbc.MysqlXADataSource; +import org.testcontainers.containers.MySQLContainer; + +import javax.sql.XADataSource; + +/** Postgres Metadata. */ +public class MySqlMetadata implements DbMetadata { + + private final String username; + private final String password; + private final String url; + private final String driver; + private final String version; + private final boolean xaEnabled; + + protected MySqlMetadata(MySQLContainer<?> container) { + this(container, false); + } + + protected MySqlMetadata(MySQLContainer<?> container, boolean hasXaEnabled) { + this.username = container.getUsername(); + this.password = container.getPassword(); + this.url = container.getJdbcUrl(); + this.driver = container.getDriverClassName(); + this.version = container.getDockerImageName(); + this.xaEnabled = hasXaEnabled; + } + + @Override + public String getUrl() { + return this.url; + } + + @Override + public XADataSource buildXaDataSource() { + if (!xaEnabled) { + throw new UnsupportedOperationException(); + } + + MysqlXADataSource xaDataSource = new MysqlXADataSource(); + xaDataSource.setUrl(getUrl()); + xaDataSource.setUser(getUser()); + xaDataSource.setPassword(getPassword()); + return xaDataSource; + } + + @Override + public String getDriverClass() { + return this.driver; + } + + @Override + public String getUser() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } +} 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 new file mode 100644 index 0000000..857568e --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleExactlyOnceSinkE2eTest.java @@ -0,0 +1,46 @@ +package org.apache.flink.connector.jdbc.dialect.oracle; + +import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; +import org.apache.flink.util.function.SerializableSupplier; + +import oracle.jdbc.xa.client.OracleXADataSource; +import org.testcontainers.containers.JdbcDatabaseContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; + +import javax.sql.XADataSource; + +import java.sql.SQLException; + +/** A simple end-to-end test for {@link JdbcExactlyOnceSinkE2eTest}. */ +@Testcontainers +public class OracleExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest { + + @Container private static final JdbcDatabaseContainer<?> CONTAINER = new OracleContainer(); + + @Override + protected String getDockerVersion() { + return CONTAINER.getDockerImageName(); + } + + @Override + protected DbMetadata getDbMetadata() { + return new OracleMetadata(CONTAINER); + } + + @Override + public SerializableSupplier<XADataSource> getDataSourceSupplier() { + return () -> { + try { + OracleXADataSource xaDataSource = new OracleXADataSource(); + xaDataSource.setURL(CONTAINER.getJdbcUrl()); + xaDataSource.setUser(CONTAINER.getUsername()); + xaDataSource.setPassword(CONTAINER.getPassword()); + return xaDataSource; + } catch (SQLException ex) { + throw new RuntimeException(ex); + } + }; + } +} 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/dialect/oracle/OracleMetadata.java new file mode 100644 index 0000000..8487c89 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/oracle/OracleMetadata.java @@ -0,0 +1,87 @@ +/* + * 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.dialect.oracle; + +import org.apache.flink.connector.jdbc.DbMetadata; + +import oracle.jdbc.xa.client.OracleXADataSource; +import org.testcontainers.containers.JdbcDatabaseContainer; + +import javax.sql.XADataSource; + +import java.sql.SQLException; + +/** Postgres Metadata. */ +public class OracleMetadata implements DbMetadata { + + private final String username; + private final String password; + private final String url; + private final String driver; + private final String version; + private final boolean xaEnabled; + + protected OracleMetadata(JdbcDatabaseContainer<?> container) { + this(container, false); + } + + protected OracleMetadata(JdbcDatabaseContainer<?> container, boolean hasXaEnabled) { + this.username = container.getUsername(); + this.password = container.getPassword(); + this.url = container.getJdbcUrl(); + this.driver = container.getDriverClassName(); + this.version = container.getDockerImageName(); + this.xaEnabled = hasXaEnabled; + } + + @Override + public String getUrl() { + return this.url; + } + + @Override + public XADataSource buildXaDataSource() { + if (!xaEnabled) { + throw new UnsupportedOperationException(); + } + try { + OracleXADataSource xaDataSource = new OracleXADataSource(); + xaDataSource.setURL(getUrl()); + xaDataSource.setUser(getUser()); + xaDataSource.setPassword(getPassword()); + return xaDataSource; + } catch (SQLException e) { + throw new RuntimeException(e); + } + } + + @Override + public String getDriverClass() { + return this.driver; + } + + @Override + public String getUser() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } +} 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 new file mode 100644 index 0000000..f8634a3 --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java @@ -0,0 +1,91 @@ +package org.apache.flink.connector.jdbc.dialect.postgres; + +import org.apache.flink.connector.jdbc.DbMetadata; +import org.apache.flink.connector.jdbc.test.DockerImageVersions; +import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; +import org.apache.flink.util.function.SerializableSupplier; + +import org.postgresql.xa.PGXADataSource; +import org.testcontainers.containers.PostgreSQLContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import javax.sql.XADataSource; + +import static org.apache.flink.util.Preconditions.checkArgument; + +/** + * A simple end-to-end test for {@link JdbcExactlyOnceSinkE2eTest}. Check for issues with suspending + * connections (requires pooling) and honoring limits (properly closing connections). + */ +@Testcontainers +public class PostgresExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest { + + @Container + private static final PostgreSQLContainer<?> CONTAINER = + new PostgresXaContainer(DockerImageVersions.POSTGRES) + .withMaxConnections(PARALLELISM * 2) + .withMaxTransactions(50); + + @Override + protected String getDockerVersion() { + return CONTAINER.getDockerImageName(); + } + + @Override + protected DbMetadata getDbMetadata() { + return new PostgresMetadata(CONTAINER); + } + + @Override + public SerializableSupplier<XADataSource> getDataSourceSupplier() { + return () -> { + PGXADataSource xaDataSource = new PGXADataSource(); + xaDataSource.setUrl(CONTAINER.getJdbcUrl()); + xaDataSource.setUser(CONTAINER.getUsername()); + xaDataSource.setPassword(CONTAINER.getPassword()); + return xaDataSource; + }; + } + + /** {@link PostgreSQLContainer} with XA enabled (by setting max_prepared_transactions). */ + public static class PostgresXaContainer extends PostgreSQLContainer<PostgresXaContainer> { + private static final int SUPERUSER_RESERVED_CONNECTIONS = 1; + private int maxConnections = SUPERUSER_RESERVED_CONNECTIONS + 1; + private int maxTransactions = 1; + + public PostgresXaContainer(String dockerImageName) { + super(DockerImageName.parse(dockerImageName)); + } + + public PostgresXaContainer withMaxConnections(int maxConnections) { + checkArgument( + maxConnections > SUPERUSER_RESERVED_CONNECTIONS, + "maxConnections should be greater than superuser_reserved_connections"); + this.maxConnections = maxConnections; + return this.self(); + } + + public PostgresXaContainer withMaxTransactions(int maxTransactions) { + checkArgument(maxTransactions > 1, "maxTransactions should be greater 1"); + this.maxTransactions = maxTransactions; + return this.self(); + } + + @Override + public void start() { + setCommand( + "postgres", + "-c", + "superuser_reserved_connections=" + SUPERUSER_RESERVED_CONNECTIONS, + "-c", + "max_connections=" + maxConnections, + "-c", + "max_prepared_transactions=" + maxTransactions, + "-c", + "fsync=off"); + super.start(); + } + } +} 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/dialect/postgres/PostgresMetadata.java new file mode 100644 index 0000000..8e171da --- /dev/null +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresMetadata.java @@ -0,0 +1,82 @@ +/* + * 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.dialect.postgres; + +import org.apache.flink.connector.jdbc.DbMetadata; + +import org.postgresql.xa.PGXADataSource; +import org.testcontainers.containers.PostgreSQLContainer; + +import javax.sql.XADataSource; + +/** Postgres Metadata. */ +public class PostgresMetadata implements DbMetadata { + + private final String username; + private final String password; + private final String url; + private final String driver; + private final String version; + private final boolean xaEnabled; + + protected PostgresMetadata(PostgreSQLContainer<?> container) { + this(container, false); + } + + protected PostgresMetadata(PostgreSQLContainer<?> container, boolean hasXaEnabled) { + this.username = container.getUsername(); + this.password = container.getPassword(); + this.url = container.getJdbcUrl(); + this.driver = container.getDriverClassName(); + this.version = container.getDockerImageName(); + this.xaEnabled = hasXaEnabled; + } + + @Override + public String getUrl() { + return this.url; + } + + @Override + public XADataSource buildXaDataSource() { + if (!xaEnabled) { + throw new UnsupportedOperationException(); + } + + PGXADataSource xaDataSource = new PGXADataSource(); + xaDataSource.setUrl(getUrl()); + xaDataSource.setUser(getUser()); + xaDataSource.setPassword(getPassword()); + return xaDataSource; + } + + @Override + public String getDriverClass() { + return this.driver; + } + + @Override + public String getUser() { + return this.username; + } + + @Override + public String getPassword() { + return this.password; + } +} 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 b40466b..4c6dbb2 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,6 +48,8 @@ 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; @@ -69,6 +71,7 @@ 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 9abfe54..bdb9a08 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,6 +29,8 @@ 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; @@ -43,6 +45,7 @@ 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/test/DockerImageVersions.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/test/DockerImageVersions.java index 4551219..22451f7 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/test/DockerImageVersions.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/test/DockerImageVersions.java @@ -24,5 +24,22 @@ package org.apache.flink.connector.jdbc.test; */ public class DockerImageVersions { - public static final String POSTGRES = "postgres:9.6.12"; + public static final String MSSQL_SERVER_2017 = "mcr.microsoft.com/mssql/server:2017-CU12"; + public static final String MSSQL_SERVER_2019 = + "mcr.microsoft.com/mssql/server:2019-GA-ubuntu-16.04"; + + public static final String MSSQL_SERVER = MSSQL_SERVER_2019; + + public static final String MYSQL_5_6 = "mysql:5.6.51"; + public static final String MYSQL_5_7 = "mysql:5.7.41"; + public static final String MYSQL_8_0 = "mysql:8.0.32"; + public static final String MYSQL = MYSQL_8_0; + + public static final String ORACLE_18 = "gvenzl/oracle-xe:18.4.0-slim-faststart"; + public static final String ORACLE_21 = "gvenzl/oracle-xe:21.3.0-slim-faststart"; + public static final String ORACLE = ORACLE_21; + + public static final String POSTGRES_9 = "postgres:9.6.24"; + public static final String POSTGRES_15 = "postgres:15.1"; + public static final String POSTGRES = POSTGRES_15; } 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 5c5867b..f22f8a9 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 @@ -23,14 +23,13 @@ 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; import org.apache.flink.connector.jdbc.JdbcSink; import org.apache.flink.connector.jdbc.JdbcTestBase; import org.apache.flink.connector.jdbc.JdbcTestFixture.TestEntry; -import org.apache.flink.connector.jdbc.dialect.oracle.OracleContainer; -import org.apache.flink.connector.jdbc.test.DockerImageVersions; import org.apache.flink.runtime.state.CheckpointListener; import org.apache.flink.runtime.state.FunctionInitializationContext; import org.apache.flink.runtime.state.FunctionSnapshotContext; @@ -41,36 +40,19 @@ import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction; import org.apache.flink.streaming.api.functions.source.SourceFunction; -import org.apache.flink.test.util.MiniClusterWithClientResource; -import org.apache.flink.testutils.junit.extensions.parameterized.Parameter; -import org.apache.flink.testutils.junit.extensions.parameterized.ParameterizedTestExtension; -import org.apache.flink.testutils.junit.extensions.parameterized.Parameters; +import org.apache.flink.test.junit5.MiniClusterExtension; import org.apache.flink.util.ExceptionUtils; import org.apache.flink.util.function.SerializableSupplier; -import com.mysql.cj.jdbc.MysqlXADataSource; -import oracle.jdbc.xa.client.OracleXADataSource; import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.TestTemplate; -import org.junit.jupiter.api.extension.ExtendWith; -import org.postgresql.xa.PGXADataSource; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.RegisterExtension; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.testcontainers.containers.JdbcDatabaseContainer; -import org.testcontainers.containers.MySQLContainer; -import org.testcontainers.containers.PostgreSQLContainer; import javax.sql.XADataSource; -import java.sql.Connection; -import java.sql.DriverManager; -import java.sql.ResultSet; -import java.sql.SQLException; -import java.sql.Statement; import java.time.Duration; -import java.util.Arrays; -import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Random; @@ -89,36 +71,43 @@ import static org.apache.flink.connector.jdbc.JdbcTestFixture.INPUT_TABLE; import static org.apache.flink.connector.jdbc.JdbcTestFixture.INSERT_TEMPLATE; import static org.apache.flink.connector.jdbc.xa.JdbcXaFacadeTestHelper.getInsertedIds; import static org.apache.flink.streaming.api.environment.ExecutionCheckpointingOptions.CHECKPOINTING_TIMEOUT; -import static org.apache.flink.util.Preconditions.checkArgument; import static org.apache.flink.util.Preconditions.checkState; import static org.assertj.core.api.Assertions.assertThat; /** A simple end-to-end test for {@link JdbcXaSinkFunction}. */ -@ExtendWith(ParameterizedTestExtension.class) -public class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { +public abstract class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { private static final Random RANDOM = new Random(System.currentTimeMillis()); private static final Logger LOG = LoggerFactory.getLogger(JdbcExactlyOnceSinkE2eTest.class); - private static final long CHECKPOINT_TIMEOUT_MS = 20_000L; - private static final long TASK_CANCELLATION_TIMEOUT_MS = 20_000L; + protected static final int PARALLELISM = 4; + protected static final long CHECKPOINT_TIMEOUT_MS = 20_000L; + protected static final long TASK_CANCELLATION_TIMEOUT_MS = 20_000L; - private interface JdbcExactlyOnceSinkTestEnv { - void start(); + abstract protected SerializableSupplier<XADataSource> getDataSourceSupplier(); - void stop(); + abstract protected String getDockerVersion(); - JdbcDatabaseContainer<?> getContainer(); + @RegisterExtension + static final MiniClusterExtension MINI_CLUSTER = createCluster(); - SerializableSupplier<XADataSource> getDataSourceSupplier(); + private static MiniClusterExtension createCluster() { + Configuration configuration = new Configuration(); + // single failover region to allow checkpointing even after some sources have finished and + // restart all tasks if at least one fails + configuration.set(EXECUTION_FAILOVER_STRATEGY, "full"); + // cancel tasks eagerly to reduce the risk of running out of memory with many restarts + configuration.set(TASK_CANCELLATION_TIMEOUT, TASK_CANCELLATION_TIMEOUT_MS); + configuration.set(CHECKPOINTING_TIMEOUT, Duration.ofMillis(CHECKPOINT_TIMEOUT_MS)); - int getParallelism(); + return new MiniClusterExtension( + new MiniClusterResourceConfiguration.Builder() + .setNumberTaskManagers(PARALLELISM) + .setConfiguration(configuration) + .build() + ); } - @Parameter public JdbcExactlyOnceSinkTestEnv dbEnv; - - private MiniClusterWithClientResource cluster; - // track active sources for: // 1. if any cancels, cancel others ASAP // 2. wait for others (to participate in checkpointing) @@ -129,69 +118,24 @@ public class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { // not using SharedObjects because we want to explicitly control which tag (attempt) to use private static final Map<Integer, CountDownLatch> inactiveMappers = new ConcurrentHashMap<>(); - @Parameters(name = "{0}") - public static Collection<JdbcExactlyOnceSinkTestEnv> parameters() { - return Arrays.asList( - // PGSQL: check for issues with suspending connections (requires pooling) and - // honoring limits (properly closing connections). - new PgSqlJdbcExactlyOnceSinkTestEnv(4), - // MYSQL: check for issues with errors on closing connections. - new MySqlJdbcExactlyOnceSinkTestEnv(4), - // ORACLE - default tests. - new OracleJdbcExactlyOnceSinkTestEnv(4) - // MSSQL - not testing: XA transactions need to be enabled via GUI (plus EULA). - // DB2 - not testing: requires auth configuration (plus EULA). - // MARIADB - not testing: XA rollback doesn't recognize recovered transactions. - ); - } - - @BeforeEach - public void before() throws Exception { - Configuration configuration = new Configuration(); - // single failover region to allow checkpointing even after some sources have finished and - // restart all tasks if at least one fails - configuration.set(EXECUTION_FAILOVER_STRATEGY, "full"); - // cancel tasks eagerly to reduce the risk of running out of memory with many restarts - configuration.set(TASK_CANCELLATION_TIMEOUT, TASK_CANCELLATION_TIMEOUT_MS); - configuration.set(CHECKPOINTING_TIMEOUT, Duration.ofMillis(CHECKPOINT_TIMEOUT_MS)); - cluster = - new MiniClusterWithClientResource( - new MiniClusterResourceConfiguration.Builder() - .setConfiguration(configuration) - // Get enough TMs to run the job. Parallelize using TMs (rather than - // slots) for better isolation - this test tends to exhaust memory - // by restarts and fast sources - .setNumberTaskManagers(dbEnv.getParallelism()) - .build()); - cluster.before(); - dbEnv.start(); - super.before(); - } - @AfterEach @Override public void after() { - // no need for cleanup - done by test container tear down - if (cluster != null) { - cluster.after(); - cluster = null; - } - dbEnv.stop(); activeSources.clear(); inactiveMappers.clear(); } - @TestTemplate + @Test void testInsert() throws Exception { long started = System.currentTimeMillis(); - LOG.info("Test insert for {}", dbEnv); + LOG.info("Test insert for {}", getDockerVersion()); int elementsPerSource = 50; int numElementsPerCheckpoint = 7; int minElementsPerFailure = numElementsPerCheckpoint / 3; int maxElementsPerFailure = numElementsPerCheckpoint * 3; StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); - env.setParallelism(dbEnv.getParallelism()); + env.setParallelism(PARALLELISM); env.setRestartStrategy(fixedDelayRestart(Integer.MAX_VALUE, Time.milliseconds(100))); env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime); env.enableCheckpointing(50, CheckpointingMode.EXACTLY_ONCE); @@ -200,7 +144,7 @@ public class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { // NOTE: keep operator chaining enabled to prevent memory exhaustion by sources while maps // are still initializing env.addSource(new TestEntrySource(elementsPerSource, numElementsPerCheckpoint)) - .setParallelism(dbEnv.getParallelism()) + .setParallelism(PARALLELISM) .map(new FailingMapper(minElementsPerFailure, maxElementsPerFailure)) .addSink( JdbcSink.exactlyOnceSink( @@ -210,18 +154,18 @@ public class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { JdbcExactlyOnceOptions.builder() .withTransactionPerConnection(true) .build(), - this.dbEnv.getDataSourceSupplier())); + this.getDataSourceSupplier())); env.execute(); List<Integer> insertedIds = getInsertedIds( - dbEnv.getContainer().getJdbcUrl(), - dbEnv.getContainer().getUsername(), - dbEnv.getContainer().getPassword(), + getDbMetadata().getUrl(), + getDbMetadata().getUser(), + getDbMetadata().getPassword(), INPUT_TABLE); List<Integer> expectedIds = - IntStream.range(0, elementsPerSource * dbEnv.getParallelism()) + IntStream.range(0, elementsPerSource * PARALLELISM) .boxed() .collect(Collectors.toList()); assertThat(insertedIds) @@ -229,45 +173,10 @@ public class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { .containsExactlyInAnyOrderElementsOf(expectedIds); LOG.info( "Test insert for {} finished in {} ms.", - dbEnv, + getDockerVersion(), System.currentTimeMillis() - started); } - @Override - protected DbMetadata getDbMetadata() { - return new DbMetadata() { - @Override - public String getInitUrl() { - return dbEnv.getContainer().getJdbcUrl(); - } - - @Override - public String getUrl() { - return dbEnv.getContainer().getJdbcUrl(); - } - - @Override - public XADataSource buildXaDataSource() { - throw new UnsupportedOperationException(); - } - - @Override - public String getDriverClass() { - return dbEnv.getContainer().getDriverClassName(); - } - - @Override - public String getUser() { - return dbEnv.getContainer().getUsername(); - } - - @Override - public String getPassword() { - return dbEnv.getContainer().getPassword(); - } - }; - } - /** {@link SourceFunction} emits {@link TestEntry test entries} and waits for the checkpoint. */ private static class TestEntrySource extends RichParallelSourceFunction<TestEntry> implements CheckpointListener, CheckpointedFunction { @@ -482,382 +391,4 @@ public class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { super("java.lang.Exception: Artificial failure", null, true, false); } } - - private static class MySqlJdbcExactlyOnceSinkTestEnv implements JdbcExactlyOnceSinkTestEnv { - private final int parallelism; - private final JdbcDatabaseContainer<?> db; - - public MySqlJdbcExactlyOnceSinkTestEnv(int parallelism) { - this.parallelism = parallelism; - this.db = new MySqlXaDb(); - } - - @Override - public void start() { - db.start(); - } - - @Override - public void stop() { - db.close(); - } - - @Override - public JdbcDatabaseContainer<?> getContainer() { - return db; - } - - @Override - public SerializableSupplier<XADataSource> getDataSourceSupplier() { - return new MySqlXaDataSourceFactory( - db.getJdbcUrl(), db.getUsername(), db.getPassword()); - } - - @Override - public int getParallelism() { - return parallelism; - } - - private static final class MySqlXaDb extends MySQLContainer<MySqlXaDb> { - private static final String IMAGE_NAME = "mysql:8.0.23"; // version 5 had issues with XA - private volatile InnoDbStatusLogger innoDbStatusLogger; - - @Override - public String toString() { - return IMAGE_NAME; - } - - public MySqlXaDb() { - super(IMAGE_NAME); - } - - @Override - public void start() { - super.start(); - long lockWaitTimeout = (CHECKPOINT_TIMEOUT_MS + TASK_CANCELLATION_TIMEOUT_MS) * 2; - // prevent XAER_RMERR: Fatal error occurred in the transaction branch - check your - // data for consistency works for mysql v8+ - try (Connection connection = - DriverManager.getConnection(getJdbcUrl(), "root", getPassword())) { - prepareDb(connection, lockWaitTimeout); - } catch (SQLException e) { - ExceptionUtils.rethrow(e); - } - this.innoDbStatusLogger = - new InnoDbStatusLogger( - getJdbcUrl(), "root", getPassword(), lockWaitTimeout / 2); - innoDbStatusLogger.start(); - } - - @Override - public void stop() { - try { - innoDbStatusLogger.stop(); - } catch (Exception e) { - ExceptionUtils.rethrow(e); - } finally { - super.stop(); - } - } - - private void prepareDb(Connection connection, long lockWaitTimeout) - throws SQLException { - try (Statement st = connection.createStatement()) { - st.execute("GRANT XA_RECOVER_ADMIN ON *.* TO '" + getUsername() + "'@'%'"); - st.execute("FLUSH PRIVILEGES"); - // if the reason of task cancellation failure is waiting for a lock - // then failing transactions with a relevant message would ease debugging - st.execute("SET GLOBAL innodb_lock_wait_timeout = " + lockWaitTimeout); - // st.execute("SET GLOBAL innodb_status_output = ON"); - // st.execute("SET GLOBAL innodb_status_output_locks = ON"); - } - } - } - - private static class MySqlXaDataSourceFactory - implements SerializableSupplier<XADataSource> { - private final String jdbcUrl; - private final String username; - private final String password; - - public MySqlXaDataSourceFactory(String jdbcUrl, String username, String password) { - this.jdbcUrl = jdbcUrl; - this.username = username; - this.password = password; - } - - @Override - public XADataSource get() { - MysqlXADataSource xaDataSource = new MysqlXADataSource(); - xaDataSource.setUrl(jdbcUrl); - xaDataSource.setUser(username); - xaDataSource.setPassword(password); - return xaDataSource; - } - } - - @Override - public String toString() { - return db + ", parallelism=" + parallelism; - } - - private static class InnoDbStatusLogger { - private static final Logger LOG = LoggerFactory.getLogger(InnoDbStatusLogger.class); - private final Thread thread; - private volatile boolean running; - - private InnoDbStatusLogger(String url, String user, String password, long intervalMs) { - running = true; - thread = - new Thread( - () -> { - LOG.info("Logging InnoDB status every {}ms", intervalMs); - try (Connection connection = - DriverManager.getConnection(url, user, password)) { - while (running) { - Thread.sleep(intervalMs); - queryAndLog(connection); - } - } catch (Exception e) { - LOG.warn("failed", e); - } finally { - LOG.info("Logging InnoDB status stopped"); - } - }); - } - - public void start() { - thread.start(); - } - - public void stop() throws InterruptedException { - running = false; - thread.join(); - } - - private void queryAndLog(Connection connection) throws SQLException { - try (Statement st = connection.createStatement()) { - showBlockedTrx(st); - showAllTrx(st); - showEngineStatus(st); - showRecoveredTrx(st); - // additional query: show full processlist \G; -- only shows live - } - } - - private void showRecoveredTrx(Statement st) throws SQLException { - try (ResultSet rs = st.executeQuery("xa recover convert xid ")) { - while (rs.next()) { - LOG.debug( - "recovered trx: {} {} {} {}", - rs.getString(1), - rs.getString(2), - rs.getString(3), - rs.getString(4)); - } - } - } - - private void showEngineStatus(Statement st) throws SQLException { - LOG.debug("Engine status"); - try (ResultSet rs = st.executeQuery("show engine innodb status")) { - while (rs.next()) { - LOG.debug(rs.getString(3)); - } - } - } - - private void showAllTrx(Statement st) throws SQLException { - LOG.debug("All TRX"); - try (ResultSet rs = - st.executeQuery("select * from information_schema.innodb_trx")) { - while (rs.next()) { - LOG.debug( - "trx_id: {}, trx_state: {}, trx_started: {}, trx_requested_lock_id: {}, trx_wait_started: {}, trx_mysql_thread_id: {},", - rs.getString("trx_id"), - rs.getString("trx_state"), - rs.getString("trx_started"), - rs.getString("trx_requested_lock_id"), - rs.getString("trx_wait_started"), - rs.getString("trx_mysql_thread_id") /* 0 for recovered*/); - } - } - } - - private void showBlockedTrx(Statement st) throws SQLException { - LOG.debug("Blocked TRX"); - try (ResultSet rs = - st.executeQuery( - " SELECT waiting_trx_id, waiting_pid, waiting_query, blocking_trx_id, blocking_pid, blocking_query " - + "FROM sys.innodb_lock_waits; ")) { - while (rs.next()) { - LOG.debug( - "waiting_trx_id: {}, waiting_pid: {}, waiting_query: {}, blocking_trx_id: {}, blocking_pid: {}, blocking_query: {}", - rs.getString(1), - rs.getString(2), - rs.getString(3), - rs.getString(4), - rs.getString(5), - rs.getString(6)); - } - } - } - } - } - - private static class PgSqlJdbcExactlyOnceSinkTestEnv implements JdbcExactlyOnceSinkTestEnv { - private final int parallelism; - private final PgXaDb db; - - private PgSqlJdbcExactlyOnceSinkTestEnv(int parallelism) { - this.parallelism = parallelism; - this.db = new PgXaDb(parallelism * 2, 50); - } - - @Override - public void start() { - db.start(); - } - - @Override - public void stop() { - db.close(); - } - - @Override - public JdbcDatabaseContainer<?> getContainer() { - return db; - } - - @Override - public SerializableSupplier<XADataSource> getDataSourceSupplier() { - return new PgXaDataSourceFactory(db.getJdbcUrl(), db.getUsername(), db.getPassword()); - } - - @Override - public int getParallelism() { - return parallelism; - } - - @Override - public String toString() { - return db + ", parallelism=" + parallelism; - } - - /** {@link PostgreSQLContainer} with XA enabled (by setting max_prepared_transactions). */ - private static final class PgXaDb extends PostgreSQLContainer<PgXaDb> { - private static final String IMAGE_NAME = DockerImageVersions.POSTGRES; - private static final int SUPERUSER_RESERVED_CONNECTIONS = 1; - - @Override - public String toString() { - return IMAGE_NAME; - } - - public PgXaDb(int maxConnections, int maxTransactions) { - super(IMAGE_NAME); - checkArgument( - maxConnections > SUPERUSER_RESERVED_CONNECTIONS, - "maxConnections should be greater than superuser_reserved_connections"); - setCommand( - "postgres", - "-c", - "superuser_reserved_connections=" + SUPERUSER_RESERVED_CONNECTIONS, - "-c", - "max_connections=" + maxConnections, - "-c", - "max_prepared_transactions=" + maxTransactions, - "-c", - "fsync=off"); - } - } - - private static class PgXaDataSourceFactory implements SerializableSupplier<XADataSource> { - private final String jdbcUrl; - private final String username; - private final String password; - - public PgXaDataSourceFactory(String jdbcUrl, String username, String password) { - this.jdbcUrl = jdbcUrl; - this.username = username; - this.password = password; - } - - @Override - public XADataSource get() { - PGXADataSource xaDataSource = new PGXADataSource(); - xaDataSource.setUrl(jdbcUrl); - xaDataSource.setUser(username); - xaDataSource.setPassword(password); - return xaDataSource; - } - } - } - - private static class OracleJdbcExactlyOnceSinkTestEnv implements JdbcExactlyOnceSinkTestEnv { - private final int parallelism; - private final OracleContainer db; - - private OracleJdbcExactlyOnceSinkTestEnv(int parallelism) { - this.parallelism = parallelism; - this.db = new OracleContainer(); - } - - @Override - public void start() { - db.start(); - } - - @Override - public void stop() { - db.close(); - } - - @Override - public JdbcDatabaseContainer<?> getContainer() { - return db; - } - - @Override - public SerializableSupplier<XADataSource> getDataSourceSupplier() { - return new OracleXaDataSourceFactory( - db.getJdbcUrl(), db.getUsername(), db.getPassword()); - } - - @Override - public int getParallelism() { - return parallelism; - } - - @Override - public String toString() { - return db + ", parallelism=" + parallelism; - } - - private static class OracleXaDataSourceFactory - implements SerializableSupplier<XADataSource> { - private final String jdbcUrl; - private final String username; - private final String password; - - public OracleXaDataSourceFactory(String jdbcUrl, String username, String password) { - this.jdbcUrl = jdbcUrl; - this.username = username; - this.password = password; - } - - @Override - public XADataSource get() { - try { - OracleXADataSource xaDataSource = new OracleXADataSource(); - xaDataSource.setURL(jdbcUrl); - xaDataSource.setUser(username); - xaDataSource.setPassword(password); - return xaDataSource; - } catch (SQLException ex) { - throw new RuntimeException(ex); - } - } - } - } }
