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);
-                }
-            }
-        }
-    }
 }

Reply via email to