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 cbb1c3538ee6b93e086eca7efffeac73279e98c9 Author: Joao Boto <[email protected]> AuthorDate: Tue Feb 7 16:12:25 2023 +0100 [FLINK-30790] Small test fixings --- flink-connector-jdbc/pom.xml | 2 +- .../flink/connector/jdbc/JdbcInputFormatTest.java | 14 ++++--- .../jdbc/catalog/MySqlCatalogTestBase.java | 18 ++++----- .../jdbc/databases/derby/DerbyDatabase.java | 2 +- .../connector/jdbc/databases/h2/H2XaDatabase.java | 2 +- .../jdbc/databases/mysql/MySqlDatabase.java | 2 +- .../jdbc/databases/oracle/OracleDatabase.java | 5 +-- .../jdbc/databases/oracle/OracleMetadata.java | 14 +++++-- .../jdbc/databases/postgres/PostgresDatabase.java | 2 +- .../databases/sqlserver/SqlServerDatabase.java | 2 +- .../jdbc/xa/JdbcExactlyOnceSinkE2eTest.java | 12 +----- .../connector/jdbc/xa/JdbcXaFacadeTestHelper.java | 44 +++++++++------------- .../connector/jdbc/xa/JdbcXaSinkDerbyTest.java | 8 +++- .../connector/jdbc/xa/JdbcXaSinkMigrationTest.java | 13 +------ .../jdbc/xa/JdbcXaSinkNoInsertionTest.java | 3 +- .../connector/jdbc/xa/JdbcXaSinkTestBase.java | 21 ++--------- 16 files changed, 66 insertions(+), 98 deletions(-) diff --git a/flink-connector-jdbc/pom.xml b/flink-connector-jdbc/pom.xml index fa09d49..5fcdb41 100644 --- a/flink-connector-jdbc/pom.xml +++ b/flink-connector-jdbc/pom.xml @@ -39,7 +39,7 @@ under the License. <scala-library.version>2.12.7</scala-library.version> <assertj.version>3.23.1</assertj.version> <postgres.version>42.5.1</postgres.version> - <oracle.version>19.3.0.0</oracle.version> + <oracle.version>21.8.0.0</oracle.version> <byte-buddy.version>1.12.10</byte-buddy.version> </properties> diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcInputFormatTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcInputFormatTest.java index ca6baa2..d766f07 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcInputFormatTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/JdbcInputFormatTest.java @@ -29,8 +29,10 @@ import org.junit.jupiter.api.Test; import java.io.IOException; import java.io.Serializable; +import java.sql.Connection; import java.sql.ResultSet; import java.sql.SQLException; +import java.sql.Statement; import static org.apache.flink.connector.jdbc.JdbcTestFixture.DERBY_EBOOKSHOP_DB; import static org.apache.flink.connector.jdbc.JdbcTestFixture.ROW_TYPE_INFO; @@ -180,8 +182,7 @@ class JdbcInputFormatTest extends JdbcDataTestBase { } @Test - void testDefaultFetchSizeIsUsedIfNotConfiguredOtherwise() - throws SQLException, ClassNotFoundException { + void testDefaultFetchSizeIsUsedIfNotConfiguredOtherwise() throws SQLException { jdbcInputFormat = JdbcInputFormat.buildJdbcInputFormat() .setDrivername(DERBY_EBOOKSHOP_DB.getDriverClass()) @@ -191,10 +192,11 @@ class JdbcInputFormatTest extends JdbcDataTestBase { .finish(); jdbcInputFormat.openInputFormat(); - final int defaultFetchSize = - DERBY_EBOOKSHOP_DB.getConnection().createStatement().getFetchSize(); - - assertThat(jdbcInputFormat.getStatement().getFetchSize()).isEqualTo(defaultFetchSize); + try (Connection dbConn = DERBY_EBOOKSHOP_DB.getConnection(); + Statement dbStatement = dbConn.createStatement(); + Statement inputStatement = jdbcInputFormat.getStatement()) { + assertThat(inputStatement.getFetchSize()).isEqualTo(dbStatement.getFetchSize()); + } } @Test 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 f8b9b07..e46b029 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 @@ -257,7 +257,7 @@ abstract class MySqlCatalogTestBase { } @Test - void testGetDb_DatabaseNotExistException() throws Exception { + void testGetDb_DatabaseNotExistException() { String databaseNotExist = "nonexistent"; assertThatThrownBy(() -> catalog.getDatabase(databaseNotExist)) .satisfies( @@ -275,7 +275,7 @@ abstract class MySqlCatalogTestBase { } @Test - void testDbExists() throws Exception { + void testDbExists() { String databaseNotExist = "nonexistent"; assertThat(catalog.databaseExists(databaseNotExist)).isFalse(); assertThat(catalog.databaseExists(TEST_DB)).isTrue(); @@ -296,7 +296,7 @@ abstract class MySqlCatalogTestBase { } @Test - void testListTables_DatabaseNotExistException() throws DatabaseNotExistException { + void testListTables_DatabaseNotExistException() { String anyDatabase = "anyDatabase"; assertThatThrownBy(() -> catalog.listTables(anyDatabase)) .satisfies( @@ -314,7 +314,7 @@ abstract class MySqlCatalogTestBase { } @Test - void testGetTables_TableNotExistException() throws TableNotExistException { + void testGetTables_TableNotExistException() { String anyTableNotExist = "anyTable"; assertThatThrownBy(() -> catalog.getTable(new ObjectPath(TEST_DB, anyTableNotExist))) .satisfies( @@ -326,7 +326,7 @@ abstract class MySqlCatalogTestBase { } @Test - void testGetTables_TableNotExistException_NoDb() throws TableNotExistException { + void testGetTables_TableNotExistException_NoDb() { String databaseNotExist = "nonexistdb"; String tableNotExist = "anyTable"; assertThatThrownBy(() -> catalog.getTable(new ObjectPath(databaseNotExist, tableNotExist))) @@ -354,8 +354,8 @@ abstract class MySqlCatalogTestBase { .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()); + assertThat(tableSchemaTestPK1.getPrimaryKey()) + .contains(tablePK1.getUnresolvedSchema().getPrimaryKey().get()); // test the PK of test2.t_user Schema tableSchemaTestPK2 = @@ -365,8 +365,8 @@ abstract class MySqlCatalogTestBase { .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()); + assertThat(tableSchemaTestPK2.getPrimaryKey()) + .contains(tablePK2.getUnresolvedSchema().getPrimaryKey().get()); } // ------ test select query. ------ 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 index d2bfae5..2db1ebe 100644 --- 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 @@ -8,7 +8,7 @@ import java.io.OutputStream; import java.sql.DriverManager; import java.sql.SQLException; -/** Derby database for testing. * */ +/** Derby database for testing. */ public interface DerbyDatabase extends DatabaseTest { @SuppressWarnings("unused") // used in string constant in prepareDatabase 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 index 3b49a57..747ec0f 100644 --- 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 @@ -23,7 +23,7 @@ import org.apache.flink.util.FlinkRuntimeException; import java.sql.DriverManager; -/** H2 database for testing. * */ +/** H2 database for testing. */ public interface H2XaDatabase extends DatabaseTest { DatabaseMetadata METADATA = startDatabase(); diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java index 6e8c69e..512b43a 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/mysql/MySqlDatabase.java @@ -36,7 +36,7 @@ import java.sql.Statement; import static org.apache.flink.util.Preconditions.checkArgument; -/** A MySql database for testing. * */ +/** A MySql database for testing. */ @Testcontainers public interface MySqlDatabase extends DatabaseTest { 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 index 0ba1f0e..d22994b 100644 --- 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 @@ -20,12 +20,11 @@ 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.testcontainers.containers.JdbcDatabaseContainer; import org.testcontainers.containers.OracleContainer; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; -/** A Oracle database for testing. * */ +/** A Oracle database for testing. */ @Testcontainers public interface OracleDatabase extends DatabaseTest { @@ -33,7 +32,7 @@ public interface OracleDatabase extends DatabaseTest { String ORACLE_21 = "gvenzl/oracle-xe:21.3.0-slim-faststart"; @Container - JdbcDatabaseContainer<?> CONTAINER = + OracleContainer CONTAINER = new OracleContainer(ORACLE_21) .withStartupTimeoutSeconds(240) .withConnectTimeoutSeconds(120) diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java index 144e165..afee02b 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/oracle/OracleMetadata.java @@ -20,7 +20,7 @@ package org.apache.flink.connector.jdbc.databases.oracle; import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import oracle.jdbc.xa.client.OracleXADataSource; -import org.testcontainers.containers.JdbcDatabaseContainer; +import org.testcontainers.containers.OracleContainer; import javax.sql.XADataSource; @@ -36,14 +36,20 @@ public class OracleMetadata implements DatabaseMetadata { private final String version; private final boolean xaEnabled; - public OracleMetadata(JdbcDatabaseContainer<?> container) { + public OracleMetadata(OracleContainer container) { this(container, false); } - public OracleMetadata(JdbcDatabaseContainer<?> container, boolean hasXaEnabled) { + public OracleMetadata(OracleContainer container, boolean hasXaEnabled) { this.username = container.getUsername(); this.password = container.getPassword(); - this.url = container.getJdbcUrl(); + this.url = + "jdbc:oracle:thin:@" + + container.getHost() + + ":" + + container.getOraclePort() + + ":" + + container.getSid(); this.driver = container.getDriverClassName(); this.version = container.getDockerImageName(); this.xaEnabled = hasXaEnabled; 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 index 65b26f1..daa607b 100644 --- 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 @@ -27,7 +27,7 @@ import org.testcontainers.utility.DockerImageName; import static org.apache.flink.util.Preconditions.checkArgument; -/** A Postgres database for testing. * */ +/** A Postgres database for testing. */ @Testcontainers public interface PostgresDatabase extends DatabaseTest { 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 index 2f0d4a6..59e1a0a 100644 --- 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 @@ -25,7 +25,7 @@ import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; import org.testcontainers.utility.DockerImageName; -/** A SqlServer database for testing. * */ +/** A SqlServer database for testing. */ @Testcontainers public interface SqlServerDatabase extends DatabaseTest { 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 383cfff..6bde49a 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 @@ -152,12 +152,7 @@ public abstract class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { env.execute(); - List<Integer> insertedIds = - getInsertedIds( - getMetadata().getUrl(), - getMetadata().getUser(), - getMetadata().getPassword(), - INPUT_TABLE); + List<Integer> insertedIds = getInsertedIds(getMetadata(), INPUT_TABLE); List<Integer> expectedIds = IntStream.range(0, elementsPerSource * PARALLELISM) .boxed() @@ -287,11 +282,6 @@ public abstract class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { && running && !Thread.currentThread().isInterrupted() && haveActiveSources()) { - if (System.currentTimeMillis() - start > 10_000) { - // debugging FLINK-22889 (TODO: remove after resolved) - LOG.debug("Slept more than 10s", new Exception()); - start = Long.MAX_VALUE; - } try { Thread.sleep(10); } catch (InterruptedException e) { diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaFacadeTestHelper.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaFacadeTestHelper.java index 1372683..9d4dfc3 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaFacadeTestHelper.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaFacadeTestHelper.java @@ -18,11 +18,9 @@ package org.apache.flink.connector.jdbc.xa; import org.apache.flink.connector.jdbc.JdbcTestCheckpoint; - -import javax.sql.XADataSource; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import java.sql.Connection; -import java.sql.DriverManager; import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Statement; @@ -37,20 +35,14 @@ import static org.assertj.core.api.Assertions.assertThat; class JdbcXaFacadeTestHelper implements AutoCloseable { private final String table; - private final String dbUrl; - private final String user; - private final String pass; + private final DatabaseMetadata metadata; private final XaFacade xaFacade; - JdbcXaFacadeTestHelper( - XADataSource xaDataSource, String dbUrl, String table, String user, String pass) - throws Exception { - this.dbUrl = dbUrl; + JdbcXaFacadeTestHelper(DatabaseMetadata metadata, String table) throws Exception { + this.metadata = metadata; this.table = table; - this.xaFacade = XaFacadeImpl.fromXaDataSource(xaDataSource); + this.xaFacade = XaFacadeImpl.fromXaDataSource(metadata.buildXaDataSource()); this.xaFacade.open(); - this.user = user; - this.pass = pass; } void assertPreparedTxCountEquals(int expected) { @@ -72,20 +64,19 @@ class JdbcXaFacadeTestHelper implements AutoCloseable { } private List<Integer> getInsertedIds() throws SQLException { - return getInsertedIds(dbUrl, user, pass, table); + return getInsertedIds(metadata, table); } - static List<Integer> getInsertedIds(String dbUrl, String user, String pass, String table) + static List<Integer> getInsertedIds(DatabaseMetadata metadata, String table) throws SQLException { List<Integer> dbContents = new ArrayList<>(); - try (Connection connection = DriverManager.getConnection(dbUrl, user, pass)) { + try (Connection connection = metadata.getConnection()) { connection.setTransactionIsolation(Connection.TRANSACTION_READ_COMMITTED); connection.setReadOnly(true); - try (Statement st = connection.createStatement()) { - try (ResultSet rs = st.executeQuery("select id from " + table)) { - while (rs.next()) { - dbContents.add(rs.getInt(1)); - } + try (Statement st = connection.createStatement(); + ResultSet rs = st.executeQuery("select id from " + table)) { + while (rs.next()) { + dbContents.add(rs.getInt(1)); } } } @@ -93,14 +84,13 @@ class JdbcXaFacadeTestHelper implements AutoCloseable { } int countInDb() throws SQLException { - try (Connection connection = DriverManager.getConnection(dbUrl)) { + try (Connection connection = metadata.getConnection()) { connection.setTransactionIsolation(Connection.TRANSACTION_READ_COMMITTED); connection.setReadOnly(true); - try (Statement st = connection.createStatement()) { - try (ResultSet rs = st.executeQuery("select count(1) from " + table)) { - rs.next(); - return rs.getInt(1); - } + try (Statement st = connection.createStatement(); + ResultSet rs = st.executeQuery("select count(1) from " + table)) { + rs.next(); + return rs.getInt(1); } } } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkDerbyTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkDerbyTest.java index 3d3e9c5..19b657d 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkDerbyTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkDerbyTest.java @@ -85,7 +85,10 @@ class JdbcXaSinkDerbyTest extends JdbcXaSinkTestBase { void testCommitUponStart() throws Exception { sinkHelper.emitAndSnapshot(JdbcTestFixture.CP0); sinkHelper.close(); - buildAndInit(0, XaFacadeImpl.fromXaDataSource(xaDataSource), sinkHelper.getState()); + buildAndInit( + 0, + XaFacadeImpl.fromXaDataSource(getMetadata().buildXaDataSource()), + sinkHelper.getState()); xaHelper.assertDbContentsEquals(JdbcTestFixture.CP0); } @@ -155,7 +158,8 @@ class JdbcXaSinkDerbyTest extends JdbcXaSinkTestBase { sinkHelper = new JdbcXaSinkTestHelper( buildAndInit( - Integer.MAX_VALUE, XaFacadeImpl.fromXaDataSource(xaDataSource)), + Integer.MAX_VALUE, + XaFacadeImpl.fromXaDataSource(getMetadata().buildXaDataSource())), new TestXaSinkStateHandler()); sinkHelper.emit(TEST_DATA[0]); sinkHelper.emit(TEST_DATA[0]); // duplicate diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkMigrationTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkMigrationTest.java index dbbf1fa..e4199a9 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkMigrationTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkMigrationTest.java @@ -89,12 +89,7 @@ public class JdbcXaSinkMigrationTest extends JdbcTestBase { harness.open(); } try (JdbcXaFacadeTestHelper h = - new JdbcXaFacadeTestHelper( - JdbcXaSinkDerbyTest.derbyXaDs(), - getMetadata().getUrl(), - JdbcTestFixture.INPUT_TABLE, - getMetadata().getUser(), - getMetadata().getPassword())) { + new JdbcXaFacadeTestHelper(getMetadata(), JdbcTestFixture.INPUT_TABLE)) { h.assertDbContentsEquals(CP0); } } @@ -180,11 +175,7 @@ public class JdbcXaSinkMigrationTest extends JdbcTestBase { private static void cancelAllTx() throws Exception { try (JdbcXaFacadeTestHelper xa = new JdbcXaFacadeTestHelper( - derbyXaDs(), - JdbcTestFixture.DERBY_EBOOKSHOP_DB.getUrl(), - JdbcTestFixture.INPUT_TABLE, - JdbcTestFixture.DERBY_EBOOKSHOP_DB.getUser(), - JdbcTestFixture.DERBY_EBOOKSHOP_DB.getPassword())) { + JdbcTestFixture.DERBY_EBOOKSHOP_DB, JdbcTestFixture.INPUT_TABLE)) { xa.cancelAllTx(); } } diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkNoInsertionTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkNoInsertionTest.java index a5f1f6b..9c028d1 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkNoInsertionTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkNoInsertionTest.java @@ -53,7 +53,8 @@ class JdbcXaSinkNoInsertionTest extends JdbcXaSinkTestBase implements H2XaDataba @Test void testNoInsertAfterFacadeClose() throws Exception { - try (XaFacadeImpl xaFacade = XaFacadeImpl.fromXaDataSource(xaDataSource)) { + try (XaFacadeImpl xaFacade = + XaFacadeImpl.fromXaDataSource(getMetadata().buildXaDataSource())) { sinkHelper = new JdbcXaSinkTestHelper( buildAndInit(0, xaFacade), new TestXaSinkStateHandler()); diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkTestBase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkTestBase.java index fb0968d..8b6f4d5 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkTestBase.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/xa/JdbcXaSinkTestBase.java @@ -55,7 +55,6 @@ import org.apache.flink.streaming.api.functions.sink.SinkFunction; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; -import javax.sql.XADataSource; import javax.transaction.xa.Xid; import java.io.Serializable; @@ -79,18 +78,10 @@ abstract class JdbcXaSinkTestBase extends JdbcTestBase { JdbcXaFacadeTestHelper xaHelper; JdbcXaSinkTestHelper sinkHelper; - XADataSource xaDataSource; @BeforeEach void initHelpers() throws Exception { - xaDataSource = getMetadata().buildXaDataSource(); - xaHelper = - new JdbcXaFacadeTestHelper( - getMetadata().buildXaDataSource(), - getMetadata().getUrl(), - INPUT_TABLE, - getMetadata().getUser(), - getMetadata().getPassword()); + xaHelper = new JdbcXaFacadeTestHelper(getMetadata(), INPUT_TABLE); sinkHelper = buildSinkHelper(createStateHandler()); } @@ -106,13 +97,7 @@ abstract class JdbcXaSinkTestBase extends JdbcTestBase { if (xaHelper != null) { xaHelper.close(); } - try (JdbcXaFacadeTestHelper xa = - new JdbcXaFacadeTestHelper( - xaDataSource, - getMetadata().getUrl(), - INPUT_TABLE, - getMetadata().getUser(), - getMetadata().getPassword())) { + try (JdbcXaFacadeTestHelper xa = new JdbcXaFacadeTestHelper(getMetadata(), INPUT_TABLE)) { xa.cancelAllTx(); } } @@ -122,7 +107,7 @@ abstract class JdbcXaSinkTestBase extends JdbcTestBase { } private XaFacadeImpl getXaFacade() { - return XaFacadeImpl.fromXaDataSource(xaDataSource); + return XaFacadeImpl.fromXaDataSource(getMetadata().buildXaDataSource()); } JdbcXaSinkFunction<TestEntry> buildAndInit() throws Exception {
