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 97fa1144e788f25f9f2c321ea8d18bcbd45f41d6 Author: Joao Boto <[email protected]> AuthorDate: Mon Jan 30 17:30:33 2023 +0100 [FLINK-30790] Change Postgres tests to use new PostgresDatabase --- .../jdbc/catalog/PostgresCatalogTestBase.java | 26 ++------- .../catalog/factory/JdbcCatalogFactoryTest.java | 26 ++------- .../jdbc/databases/postgres/PostgresDatabase.java | 46 ++++++++++++++- .../postgres/PostgresExactlyOnceSinkE2eTest.java | 65 +--------------------- 4 files changed, 57 insertions(+), 106 deletions(-) diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/PostgresCatalogTestBase.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/PostgresCatalogTestBase.java index 2666bbc..bd3b982 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/PostgresCatalogTestBase.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/PostgresCatalogTestBase.java @@ -18,7 +18,7 @@ package org.apache.flink.connector.jdbc.catalog; -import org.apache.flink.connector.jdbc.test.DockerImageVersions; +import org.apache.flink.connector.jdbc.databases.postgres.PostgresDatabase; import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.api.Schema; import org.apache.flink.table.types.logical.DecimalType; @@ -26,11 +26,6 @@ import org.apache.flink.table.types.logical.DecimalType; import org.junit.jupiter.api.BeforeAll; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.testcontainers.containers.PostgreSQLContainer; -import org.testcontainers.containers.output.Slf4jLogConsumer; -import org.testcontainers.junit.jupiter.Container; -import org.testcontainers.junit.jupiter.Testcontainers; -import org.testcontainers.utility.DockerImageName; import java.sql.Connection; import java.sql.DriverManager; @@ -38,17 +33,13 @@ import java.sql.SQLException; import java.sql.Statement; /** Test base for {@link PostgresCatalog}. */ -@Testcontainers -class PostgresCatalogTestBase { +class PostgresCatalogTestBase implements PostgresDatabase { public static final Logger LOG = LoggerFactory.getLogger(PostgresCatalogTestBase.class); - protected static final DockerImageName POSTGRES_IMAGE = - DockerImageName.parse(DockerImageVersions.POSTGRES); - protected static final String TEST_CATALOG_NAME = "mypg"; - protected static final String TEST_USERNAME = "postgres"; - protected static final String TEST_PWD = "postgres"; + protected static final String TEST_USERNAME = CONTAINER.getUsername(); + protected static final String TEST_PWD = CONTAINER.getPassword(); protected static final String TEST_DB = "test"; protected static final String TEST_SCHEMA = "test_schema"; protected static final String TABLE1 = "t1"; @@ -64,17 +55,10 @@ class PostgresCatalogTestBase { protected static String baseUrl; protected static PostgresCatalog catalog; - @Container - static final PostgreSQLContainer<?> POSTGRES_CONTAINER = - new PostgreSQLContainer<>(POSTGRES_IMAGE) - .withUsername(TEST_USERNAME) - .withPassword(TEST_PWD) - .withLogConsumer(new Slf4jLogConsumer(LOG)); - @BeforeAll static void init() throws SQLException { // jdbc:postgresql://localhost:50807/postgres?user=postgres - String jdbcUrl = POSTGRES_CONTAINER.getJdbcUrl(); + String jdbcUrl = CONTAINER.getJdbcUrl(); // jdbc:postgresql://localhost:50807/ baseUrl = jdbcUrl.substring(0, jdbcUrl.lastIndexOf("/")); diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryTest.java index e3aaefc..705e494 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/catalog/factory/JdbcCatalogFactoryTest.java @@ -20,7 +20,7 @@ package org.apache.flink.connector.jdbc.catalog.factory; import org.apache.flink.connector.jdbc.catalog.JdbcCatalog; import org.apache.flink.connector.jdbc.catalog.PostgresCatalog; -import org.apache.flink.connector.jdbc.test.DockerImageVersions; +import org.apache.flink.connector.jdbc.databases.postgres.PostgresDatabase; import org.apache.flink.table.catalog.Catalog; import org.apache.flink.table.catalog.CommonCatalogOptions; import org.apache.flink.table.factories.FactoryUtil; @@ -29,11 +29,6 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.testcontainers.containers.PostgreSQLContainer; -import org.testcontainers.containers.output.Slf4jLogConsumer; -import org.testcontainers.junit.jupiter.Container; -import org.testcontainers.junit.jupiter.Testcontainers; -import org.testcontainers.utility.DockerImageName; import java.sql.SQLException; import java.util.HashMap; @@ -42,8 +37,7 @@ import java.util.Map; import static org.assertj.core.api.Assertions.assertThat; /** Test for {@link JdbcCatalogFactory}. */ -@Testcontainers -class JdbcCatalogFactoryTest { +class JdbcCatalogFactoryTest implements PostgresDatabase { public static final Logger LOG = LoggerFactory.getLogger(JdbcCatalogFactoryTest.class); @@ -51,23 +45,13 @@ class JdbcCatalogFactoryTest { protected static JdbcCatalog catalog; protected static final String TEST_CATALOG_NAME = "mypg"; - protected static final String TEST_USERNAME = "postgres"; - protected static final String TEST_PWD = "postgres"; - - protected static final DockerImageName POSTGRES_IMAGE = - DockerImageName.parse(DockerImageVersions.POSTGRES); - - @Container - static final PostgreSQLContainer<?> POSTGRES_CONTAINER = - new PostgreSQLContainer<>(POSTGRES_IMAGE) - .withUsername(TEST_USERNAME) - .withPassword(TEST_PWD) - .withLogConsumer(new Slf4jLogConsumer(LOG)); + protected static final String TEST_USERNAME = CONTAINER.getUsername(); + protected static final String TEST_PWD = CONTAINER.getPassword(); @BeforeAll static void setup() throws SQLException { // jdbc:postgresql://localhost:50807/postgres?user=postgres - String jdbcUrl = POSTGRES_CONTAINER.getJdbcUrl(); + String jdbcUrl = CONTAINER.getJdbcUrl(); // jdbc:postgresql://localhost:50807/ baseUrl = jdbcUrl.substring(0, jdbcUrl.lastIndexOf("/")); 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 0db1ec7..93e0fe3 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 @@ -19,12 +19,14 @@ package org.apache.flink.connector.jdbc.databases.postgres; import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.apache.flink.connector.jdbc.databases.DatabaseTest; -import org.apache.flink.connector.jdbc.dialect.postgres.PostgresExactlyOnceSinkE2eTest; import org.apache.flink.connector.jdbc.test.DockerImageVersions; import org.testcontainers.containers.PostgreSQLContainer; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import static org.apache.flink.util.Preconditions.checkArgument; /** A Postgres database for testing. * */ @Testcontainers @@ -32,7 +34,7 @@ public interface PostgresDatabase extends DatabaseTest { @Container PostgreSQLContainer<?> CONTAINER = - new PostgresExactlyOnceSinkE2eTest.PostgresXaContainer(DockerImageVersions.POSTGRES) + new PostgresXaContainer(DockerImageVersions.POSTGRES) .withMaxConnections(10) .withMaxTransactions(50); @@ -40,4 +42,44 @@ public interface PostgresDatabase extends DatabaseTest { default DatabaseMetadata getMetadata() { return new PostgresMetadata(CONTAINER); } + + /** {@link PostgreSQLContainer} with XA enabled (by setting max_prepared_transactions). */ + 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/PostgresExactlyOnceSinkE2eTest.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java index daa173c..0a2c39c 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java @@ -1,38 +1,19 @@ package org.apache.flink.connector.jdbc.dialect.postgres; -import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; -import org.apache.flink.connector.jdbc.databases.postgres.PostgresMetadata; -import org.apache.flink.connector.jdbc.test.DockerImageVersions; +import org.apache.flink.connector.jdbc.databases.postgres.PostgresDatabase; 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 - public DatabaseMetadata getMetadata() { - return new PostgresMetadata(CONTAINER); - } +public class PostgresExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest + implements PostgresDatabase { @Override public SerializableSupplier<XADataSource> getDataSourceSupplier() { @@ -44,44 +25,4 @@ public class PostgresExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest { 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(); - } - } }
