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

Reply via email to