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 867fc0f29a1eb140af0201d65cfc57eaf447dad0 Author: Joao Boto <[email protected]> AuthorDate: Tue Jan 31 17:52:50 2023 +0100 [FLINK-30790] Fix JdbcExactlyOnceSinkE2eTest loop on failure --- .../apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java | 2 +- .../org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java | 2 +- .../flink/connector/jdbc/databases/oracle/OracleDatabase.java | 5 +++-- .../connector/jdbc/dialect/mysql/MySqlExactlyOnceSinkE2eTest.java | 7 +++++++ .../jdbc/dialect/oracle/OracleExactlyOnceSinkE2eTest.java | 7 +++++++ .../jdbc/dialect/postgres/PostgresExactlyOnceSinkE2eTest.java | 7 +++++++ .../apache/flink/connector/jdbc/xa/JdbcExactlyOnceSinkE2eTest.java | 2 +- 7 files changed, 27 insertions(+), 5 deletions(-) diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java index 960a56c..087aa58 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/derby/DerbyMetadata.java @@ -23,7 +23,7 @@ import org.apache.derby.jdbc.EmbeddedXADataSource; import javax.sql.XADataSource; -/** DerbyDbMetadata. */ +/** Derby Metadata. */ public class DerbyMetadata implements DatabaseMetadata { private final String dbName; diff --git a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java index 95db4c6..ee41685 100644 --- a/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java +++ b/flink-connector-jdbc/src/test/java/org/apache/flink/connector/jdbc/databases/h2/H2Metadata.java @@ -22,7 +22,7 @@ import org.apache.flink.connector.jdbc.databases.h2.xa.H2XaDsWrapper; import javax.sql.XADataSource; -/** H2DbMetadata. */ +/** H2 Metadata. */ public class H2Metadata implements DatabaseMetadata { private final String schema; 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 13269d2..0ba1f0e 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 @@ -29,14 +29,15 @@ import org.testcontainers.junit.jupiter.Testcontainers; @Testcontainers public interface OracleDatabase extends DatabaseTest { - String ORACLE_18 = "gvenzl/oracle-xe:18.4.0-slim-faststart"; + String ORACLE_18 = "gvenzl/oracle-xe:18.4.0-slim"; String ORACLE_21 = "gvenzl/oracle-xe:21.3.0-slim-faststart"; @Container JdbcDatabaseContainer<?> CONTAINER = new OracleContainer(ORACLE_21) .withStartupTimeoutSeconds(240) - .withConnectTimeoutSeconds(120); + .withConnectTimeoutSeconds(120) + .usingSid(); @Override default DatabaseMetadata getMetadata() { 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 index 81d4690..24fa8df 100644 --- 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 @@ -1,6 +1,8 @@ package org.apache.flink.connector.jdbc.dialect.mysql; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.apache.flink.connector.jdbc.databases.mysql.MySqlDatabase; +import org.apache.flink.connector.jdbc.databases.mysql.MySqlMetadata; import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.function.SerializableSupplier; @@ -15,6 +17,11 @@ import javax.sql.XADataSource; public class MySqlExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest implements MySqlDatabase { + @Override + public DatabaseMetadata getMetadata() { + return new MySqlMetadata(CONTAINER, true); + } + @Override public SerializableSupplier<XADataSource> getDataSourceSupplier() { return () -> { 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 index 7f9318c..27a41ae 100644 --- 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 @@ -1,6 +1,8 @@ package org.apache.flink.connector.jdbc.dialect.oracle; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.apache.flink.connector.jdbc.databases.oracle.OracleDatabase; +import org.apache.flink.connector.jdbc.databases.oracle.OracleMetadata; import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.function.SerializableSupplier; @@ -17,6 +19,11 @@ import java.sql.SQLException; public class OracleExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest implements OracleDatabase { + @Override + public DatabaseMetadata getMetadata() { + return new OracleMetadata(CONTAINER, true); + } + @Override public SerializableSupplier<XADataSource> getDataSourceSupplier() { return () -> { 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 0a2c39c..609523a 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,6 +1,8 @@ package org.apache.flink.connector.jdbc.dialect.postgres; +import org.apache.flink.connector.jdbc.databases.DatabaseMetadata; import org.apache.flink.connector.jdbc.databases.postgres.PostgresDatabase; +import org.apache.flink.connector.jdbc.databases.postgres.PostgresMetadata; import org.apache.flink.connector.jdbc.xa.JdbcExactlyOnceSinkE2eTest; import org.apache.flink.util.function.SerializableSupplier; @@ -15,6 +17,11 @@ import javax.sql.XADataSource; public class PostgresExactlyOnceSinkE2eTest extends JdbcExactlyOnceSinkE2eTest implements PostgresDatabase { + @Override + public DatabaseMetadata getMetadata() { + return new PostgresMetadata(CONTAINER, true); + } + @Override public SerializableSupplier<XADataSource> getDataSourceSupplier() { return () -> { 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 d229bec..383cfff 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 @@ -130,7 +130,7 @@ public abstract class JdbcExactlyOnceSinkE2eTest extends JdbcTestBase { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(PARALLELISM); - env.setRestartStrategy(fixedDelayRestart(Integer.MAX_VALUE, Time.milliseconds(100))); + env.setRestartStrategy(fixedDelayRestart(elementsPerSource * 2, Time.milliseconds(100))); env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime); env.enableCheckpointing(50, CheckpointingMode.EXACTLY_ONCE); // timeout checkpoints as some tasks may fail while triggering
