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 5623cb7cd76e697e5c5c5a271eadb5402622f181 Author: Joao Boto <[email protected]> AuthorDate: Mon Jan 30 17:20:47 2023 +0100 [FLINK-30790] Change SqlServer tests to use new SqlServerDatabase --- .../sqlserver/SqlServerTableSinkITCase.java | 96 +++++++++++----------- .../sqlserver/SqlServerTableSourceITCase.java | 32 ++++---- 2 files changed, 66 insertions(+), 62 deletions(-) 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..42daf51 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 @@ -20,6 +20,7 @@ package org.apache.flink.connector.jdbc.dialect.sqlserver; import org.apache.flink.api.java.tuple.Tuple4; import org.apache.flink.configuration.Configuration; +import org.apache.flink.connector.jdbc.databases.sqlserver.SqlServerDatabase; import org.apache.flink.connector.jdbc.internal.GenericJdbcSinkFunction; import org.apache.flink.runtime.state.StateSnapshotContextSynchronousImpl; import org.apache.flink.streaming.api.datastream.DataStream; @@ -48,9 +49,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.testcontainers.containers.MSSQLServerContainer; -import org.testcontainers.junit.jupiter.Container; -import org.testcontainers.junit.jupiter.Testcontainers; +import org.junit.jupiter.api.condition.DisabledOnOs; +import org.junit.jupiter.api.condition.OS; import java.math.BigDecimal; import java.sql.Connection; @@ -69,15 +69,8 @@ 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}. */ -@Testcontainers -class SqlServerTableSinkITCase extends AbstractTestBase { - - @Container - private static final MSSQLServerContainer<?> container = - new MSSQLServerContainer("mcr.microsoft.com/mssql/server:2019-GA-ubuntu-16.04") - .acceptLicense(); - - private static String containerUrl; +@DisabledOnOs(OS.MAC) +class SqlServerTableSinkITCase extends AbstractTestBase implements SqlServerDatabase { public static final String OUTPUT_TABLE1 = "dynamicSinkForUpsert"; public static final String OUTPUT_TABLE2 = "dynamicSinkForAppend"; @@ -88,14 +81,12 @@ class SqlServerTableSinkITCase extends AbstractTestBase { @BeforeAll static void beforeAll() throws ClassNotFoundException, SQLException { - containerUrl = - String.format( - "%s;username=%s;password=%s", - container.getJdbcUrl(), container.getUsername(), container.getPassword()); - Class.forName(container.getDriverClassName()); + Class.forName(CONTAINER.getDriverClassName()); try (Connection conn = DriverManager.getConnection( - containerUrl, container.getUsername(), container.getPassword()); + CONTAINER.getJdbcUrl(), + CONTAINER.getUsername(), + CONTAINER.getPassword()); Statement stat = conn.createStatement()) { stat.executeUpdate( "CREATE TABLE " @@ -143,10 +134,12 @@ class SqlServerTableSinkITCase extends AbstractTestBase { @AfterAll static void afterAll() throws Exception { TestValuesTableFactory.clearAllData(); - Class.forName(container.getDriverClassName()); + Class.forName(CONTAINER.getDriverClassName()); try (Connection conn = DriverManager.getConnection( - containerUrl, container.getUsername(), container.getPassword()); + CONTAINER.getJdbcUrl(), + CONTAINER.getUsername(), + CONTAINER.getPassword()); Statement stat = conn.createStatement()) { stat.execute("DROP TABLE " + OUTPUT_TABLE1); stat.execute("DROP TABLE " + OUTPUT_TABLE2); @@ -155,7 +148,6 @@ class SqlServerTableSinkITCase extends AbstractTestBase { stat.execute("DROP TABLE " + OUTPUT_TABLE5); stat.execute("DROP TABLE " + USER_TABLE); } - container.stop(); } public static DataStream<Tuple4<Integer, Long, String, Timestamp>> get4TupleDataStream( @@ -207,21 +199,25 @@ class SqlServerTableSinkITCase extends AbstractTestBase { + ") WITH (" + " 'connector'='jdbc'," + " 'url'='" - + containerUrl + + getMetadata().getUrl() + "'," + " 'table-name'='" + OUTPUT_TABLE4 + "'," + " 'username'='" - + container.getUsername() + + getMetadata().getUsername() + "'," + " 'password'='" - + container.getPassword() + + getMetadata().getPassword() + "'" + ")"); tEnv.executeSql("INSERT INTO upsertSink SELECT CAST(1.1 as FLOAT)").await(); - check(new Row[] {Row.of(1.1f)}, containerUrl, "REAL_TABLE", new String[] {"real_data"}); + check( + new Row[] {Row.of(1.1f)}, + getMetadata().getJdbcUrlWithCredentials(), + "REAL_TABLE", + new String[] {"real_data"}); } @Test @@ -259,16 +255,16 @@ class SqlServerTableSinkITCase extends AbstractTestBase { + ") WITH (" + " 'connector'='jdbc'," + " 'url'='" - + containerUrl + + getMetadata().getJdbcUrl() + "'," + " 'table-name'='" + OUTPUT_TABLE1 + "'," + " 'username'='" - + container.getUsername() + + getMetadata().getUsername() + "'," + " 'password'='" - + container.getPassword() + + getMetadata().getPassword() + "'," + " 'sink.buffer-flush.max-rows' = '2'," + " 'sink.buffer-flush.interval' = '0'," @@ -291,7 +287,7 @@ class SqlServerTableSinkITCase extends AbstractTestBase { Row.of(7.0, 1.0, 1, Timestamp.valueOf("1970-01-01 00:00:00.021")), Row.of(9.0, 1.0, 1, Timestamp.valueOf("1970-01-01 00:00:00.015")) }, - containerUrl, + getMetadata().getJdbcUrlWithCredentials(), OUTPUT_TABLE1, new String[] {"cnt", "lencnt", "cTag", "ts"}); } @@ -317,16 +313,16 @@ class SqlServerTableSinkITCase extends AbstractTestBase { + ") WITH (" + " 'connector'='jdbc'," + " 'url'='" - + containerUrl + + getMetadata().getJdbcUrl() + "'," + " 'table-name'='" + OUTPUT_TABLE2 + "'," + " 'username'='" - + container.getUsername() + + getMetadata().getUsername() + "'," + " 'password'='" - + container.getPassword() + + getMetadata().getPassword() + "'" + ")"); @@ -338,7 +334,7 @@ class SqlServerTableSinkITCase extends AbstractTestBase { Row.of(10, 4, Timestamp.valueOf("1970-01-01 00:00:00.01")), Row.of(20, 6, Timestamp.valueOf("1970-01-01 00:00:00.02")) }, - containerUrl, + getMetadata().getJdbcUrlWithCredentials(), OUTPUT_TABLE2, new String[] {"id", "num", "ts"}); } @@ -354,16 +350,16 @@ class SqlServerTableSinkITCase extends AbstractTestBase { + ") WITH ( " + "'connector' = 'jdbc'," + "'url'='" - + containerUrl + + getMetadata().getJdbcUrl() + "'," + "'table-name' = '" + OUTPUT_TABLE3 + "'," + " 'username'='" - + container.getUsername() + + getMetadata().getUsername() + "'," + " 'password'='" - + container.getPassword() + + getMetadata().getPassword() + "'," + "'sink.buffer-flush.max-rows' = '2'," + "'sink.buffer-flush.interval' = '300ms'," @@ -387,7 +383,7 @@ class SqlServerTableSinkITCase extends AbstractTestBase { Row.of("Kim", 42), Row.of("Bob", 1) }, - containerUrl, + getMetadata().getJdbcUrlWithCredentials(), OUTPUT_TABLE3, new String[] {"NAME", "SCORE"}); } @@ -420,16 +416,16 @@ class SqlServerTableSinkITCase extends AbstractTestBase { + ") WITH (\n" + " 'connector' = 'jdbc'," + " 'url'='" - + containerUrl + + getMetadata().getJdbcUrl() + "'," + " 'table-name' = '" + USER_TABLE + "'," + " 'username'='" - + container.getUsername() + + getMetadata().getUsername() + "'," + " 'password'='" - + container.getPassword() + + getMetadata().getPassword() + "'," + " 'sink.buffer-flush.max-rows' = '2'," + " 'sink.buffer-flush.interval' = '0'" @@ -458,7 +454,7 @@ class SqlServerTableSinkITCase extends AbstractTestBase { new BigDecimal("11.30"), new BigDecimal("22.60")) }, - containerUrl, + getMetadata().getJdbcUrlWithCredentials(), USER_TABLE, new String[] {"user_id", "user_name", "email", "balance", "balance2"}); } @@ -467,11 +463,11 @@ class SqlServerTableSinkITCase extends AbstractTestBase { void testFlushBufferWhenCheckpoint() throws Exception { Map<String, String> options = new HashMap<>(); options.put("connector", "jdbc"); - options.put("url", containerUrl); + options.put("url", getMetadata().getJdbcUrl()); options.put("table-name", OUTPUT_TABLE5); options.put("sink.buffer-flush.interval", "0"); - options.put("username", container.getUsername()); - options.put("password", container.getPassword()); + options.put("username", getMetadata().getUsername()); + options.put("password", getMetadata().getPassword()); ResolvedSchema schema = ResolvedSchema.of(Column.physical("id", DataTypes.BIGINT().notNull())); @@ -488,9 +484,17 @@ class SqlServerTableSinkITCase extends AbstractTestBase { sinkFunction.invoke(GenericRowData.of(1L), SinkContextUtil.forTimestamp(1)); sinkFunction.invoke(GenericRowData.of(2L), SinkContextUtil.forTimestamp(1)); - check(new Row[] {}, containerUrl, OUTPUT_TABLE5, new String[] {"id"}); + check( + new Row[] {}, + getMetadata().getJdbcUrlWithCredentials(), + OUTPUT_TABLE5, + new String[] {"id"}); sinkFunction.snapshotState(new StateSnapshotContextSynchronousImpl(1, 1)); - check(new Row[] {Row.of(1L), Row.of(2L)}, containerUrl, OUTPUT_TABLE5, new String[] {"id"}); + check( + new Row[] {Row.of(1L), Row.of(2L)}, + getMetadata().getJdbcUrlWithCredentials(), + OUTPUT_TABLE5, + new String[] {"id"}); sinkFunction.close(); } } 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..a3253ca 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 @@ -18,6 +18,7 @@ package org.apache.flink.connector.jdbc.dialect.sqlserver; +import org.apache.flink.connector.jdbc.databases.sqlserver.SqlServerDatabase; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.TableEnvironment; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; @@ -29,7 +30,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.testcontainers.containers.MSSQLServerContainer; +import org.junit.jupiter.api.condition.DisabledOnOs; +import org.junit.jupiter.api.condition.OS; import java.sql.Connection; import java.sql.DriverManager; @@ -43,12 +45,9 @@ import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; /** The Table Source ITCase for {@link SqlServerDialect}. */ -class SqlServerTableSourceITCase extends AbstractTestBase { +@DisabledOnOs(OS.MAC) +class SqlServerTableSourceITCase extends AbstractTestBase implements SqlServerDatabase { - private static final MSSQLServerContainer container = - new MSSQLServerContainer("mcr.microsoft.com/mssql/server:2019-GA-ubuntu-16.04") - .acceptLicense(); - private static String containerUrl; private static final String INPUT_TABLE = "sql_test_table"; private static StreamExecutionEnvironment env; @@ -56,12 +55,12 @@ class SqlServerTableSourceITCase extends AbstractTestBase { @BeforeAll static void beforeAll() throws ClassNotFoundException, SQLException { - container.start(); - containerUrl = container.getJdbcUrl(); - Class.forName(container.getDriverClassName()); + Class.forName(CONTAINER.getDriverClassName()); try (Connection conn = DriverManager.getConnection( - containerUrl, container.getUsername(), container.getPassword()); + CONTAINER.getJdbcUrl(), + CONTAINER.getUsername(), + CONTAINER.getPassword()); Statement statement = conn.createStatement()) { statement.executeUpdate( "CREATE TABLE " @@ -108,14 +107,15 @@ class SqlServerTableSourceITCase extends AbstractTestBase { @AfterAll static void afterAll() throws Exception { - Class.forName(container.getDriverClassName()); + Class.forName(CONTAINER.getDriverClassName()); try (Connection conn = DriverManager.getConnection( - containerUrl, container.getUsername(), container.getPassword()); + CONTAINER.getJdbcUrl(), + CONTAINER.getUsername(), + CONTAINER.getPassword()); Statement statement = conn.createStatement()) { statement.executeUpdate("DROP TABLE " + INPUT_TABLE); } - container.stop(); } @BeforeEach @@ -213,16 +213,16 @@ class SqlServerTableSourceITCase extends AbstractTestBase { + ") WITH (" + " 'connector'='jdbc'," + " 'url'='" - + containerUrl + + getMetadata().getJdbcUrl() + "'," + " 'table-name'='" + INPUT_TABLE + "'," + " 'username'='" - + container.getUsername() + + getMetadata().getUsername() + "'," + " 'password'='" - + container.getPassword() + + getMetadata().getPassword() + "'" + ")"); }
