This is an automated email from the ASF dual-hosted git repository. dockerzhang pushed a commit to branch branch-1.7 in repository https://gitbox.apache.org/repos/asf/inlong.git
commit 009d0088d5a01170d1efffa98a3e6a42bb3ad95e Author: emhui <[email protected]> AuthorDate: Fri May 12 16:06:55 2023 +0800 [INLONG-8013][Sort] Fix UT failed caused by OOM (#8015) --- ...tgreSQLTest.java => AllMigrateMongoDBTest.java} | 39 ++++++---------------- .../inlong/sort/parser/AllMigrateOracleTest.java | 3 +- .../sort/parser/AllMigratePostgreSQLTest.java | 3 +- .../apache/inlong/sort/parser/AllMigrateTest.java | 3 +- .../parser/AllMigrateWithSpecifyingFieldTest.java | 4 ++- .../sort/parser/ClickHouseSqlParserTest.java | 3 +- .../inlong/sort/parser/DorisMultipleSinkTest.java | 3 +- .../inlong/sort/parser/ESMultipleSinkTest.java | 3 +- .../sort/parser/FilesystemSqlParserTest.java | 3 +- .../sort/parser/NativeFlinkSqlParserTest.java | 3 +- .../inlong/sort/parser/PulsarSqlParserTest.java | 3 +- 11 files changed, 31 insertions(+), 39 deletions(-) diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateMongoDBTest.java similarity index 72% copy from inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java copy to inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateMongoDBTest.java index df0f805b7..2bd69280c 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateMongoDBTest.java @@ -26,6 +26,7 @@ import java.util.stream.Collectors; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.common.enums.MetaField; import org.apache.inlong.sort.formats.common.StringFormatInfo; import org.apache.inlong.sort.parser.impl.FlinkSqlParser; @@ -35,33 +36,24 @@ import org.apache.inlong.sort.protocol.GroupInfo; import org.apache.inlong.sort.protocol.MetaFieldInfo; import org.apache.inlong.sort.protocol.StreamInfo; import org.apache.inlong.sort.protocol.node.Node; -import org.apache.inlong.sort.protocol.node.extract.PostgresExtractNode; -import org.apache.inlong.sort.protocol.node.format.CanalJsonFormat; +import org.apache.inlong.sort.protocol.node.extract.MongoExtractNode; import org.apache.inlong.sort.protocol.node.format.CsvFormat; -import org.apache.inlong.sort.protocol.node.load.DorisLoadNode; import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode; import org.apache.inlong.sort.protocol.transformation.FieldRelation; import org.apache.inlong.sort.protocol.transformation.relation.NodeRelation; import org.junit.Assert; import org.junit.Test; -/** - * A demo of transferring data between PostgreSQL Server and other database, using postgres-cdc-inlong connector. - */ -public class AllMigratePostgreSQLTest { +public class AllMigrateMongoDBTest extends AbstractTestBase { - private PostgresExtractNode buildAllMigrateExtractNode() { + private MongoExtractNode buildAllMigrateExtractNode() { List<FieldInfo> fields = Arrays.asList( new MetaFieldInfo("data", MetaField.DATA_CANAL)); - Map<String, String> option = new HashMap<>(); - option.put("source.multiple.enable", "true"); - List<String> tableNames = Arrays.asList("table"); - // List<String> tableNames = Arrays.asList("*"); - PostgresExtractNode node = new PostgresExtractNode("1", "pg_input", fields, - null, option, null, tableNames, "localhost", "user", - "password", "db", "schema", 1234, - null, "Asia/Shanghai", "initial"); - return node; + Map<String, String> options = new HashMap<>(); + options.put("source.multiple.enable", "true"); + return new MongoExtractNode("1", "mongodb_input", fields, + null, options, "db.test", "localhost:27017", + "root", "inlong", "db"); } private KafkaLoadNode buildAllMigrateKafkaNode() { @@ -72,21 +64,11 @@ public class AllMigratePostgreSQLTest { CsvFormat csvFormat = new CsvFormat(); csvFormat.setDisableQuoteCharacter(true); return new KafkaLoadNode("2", "kafka_output", fields, relations, null, null, - "pg_db", "localhost:9092", + "test", "localhost:9092", csvFormat, null, null, null); } - private DorisLoadNode buildAllMigrateDorisNode() { - List<FieldInfo> fields = Arrays.asList(new FieldInfo("data", new StringFormatInfo())); - List<FieldRelation> relations = Arrays.asList(new FieldRelation(new FieldInfo("data", new StringFormatInfo()), - new FieldInfo("data", new StringFormatInfo()))); - return new DorisLoadNode("2", "doris_output", fields, relations, null, null, - null, null, "localhost:1234", "user", "password", - null, null, true, new CanalJsonFormat(), - "${database}", "${table}"); - } - private NodeRelation buildNodeRelation(List<Node> inputs, List<Node> outputs) { List<String> inputIds = inputs.stream().map(Node::getId).collect(Collectors.toList()); List<String> outputIds = outputs.stream().map(Node::getId).collect(Collectors.toList()); @@ -111,7 +93,6 @@ public class AllMigratePostgreSQLTest { StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, settings); Node inputNode = buildAllMigrateExtractNode(); Node outputNode = buildAllMigrateKafkaNode(); - // Node outputNode = buildAllMigrateDorisNode(); StreamInfo streamInfo = new StreamInfo("1", Arrays.asList(inputNode, outputNode), Collections.singletonList(buildNodeRelation(Collections.singletonList(inputNode), Collections.singletonList(outputNode)))); diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateOracleTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateOracleTest.java index 36029a0b2..9091bf55d 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateOracleTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateOracleTest.java @@ -26,6 +26,7 @@ import java.util.stream.Collectors; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.common.enums.MetaField; import org.apache.inlong.sort.formats.common.StringFormatInfo; import org.apache.inlong.sort.formats.common.VarBinaryFormatInfo; @@ -45,7 +46,7 @@ import org.apache.inlong.sort.protocol.transformation.relation.NodeRelation; import org.junit.Assert; import org.junit.Test; -public class AllMigrateOracleTest { +public class AllMigrateOracleTest extends AbstractTestBase { private OracleExtractNode buildAllMigrateExtractNode() { List<FieldInfo> fields = Arrays.asList( diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java index df0f805b7..4bb96462c 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigratePostgreSQLTest.java @@ -26,6 +26,7 @@ import java.util.stream.Collectors; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.common.enums.MetaField; import org.apache.inlong.sort.formats.common.StringFormatInfo; import org.apache.inlong.sort.parser.impl.FlinkSqlParser; @@ -48,7 +49,7 @@ import org.junit.Test; /** * A demo of transferring data between PostgreSQL Server and other database, using postgres-cdc-inlong connector. */ -public class AllMigratePostgreSQLTest { +public class AllMigratePostgreSQLTest extends AbstractTestBase { private PostgresExtractNode buildAllMigrateExtractNode() { List<FieldInfo> fields = Arrays.asList( diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java index fddceb58e..e024f7c01 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateTest.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.common.enums.MetaField; import org.apache.inlong.sort.formats.common.StringFormatInfo; import org.apache.inlong.sort.formats.common.VarBinaryFormatInfo; @@ -47,7 +48,7 @@ import java.util.List; import java.util.Map; import java.util.stream.Collectors; -public class AllMigrateTest { +public class AllMigrateTest extends AbstractTestBase { private MySqlExtractNode buildAllMigrateExtractNode() { diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateWithSpecifyingFieldTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateWithSpecifyingFieldTest.java index 352a359c0..0194336d3 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateWithSpecifyingFieldTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/AllMigrateWithSpecifyingFieldTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.common.enums.MetaField; import org.apache.inlong.sort.formats.common.VarBinaryFormatInfo; import org.apache.inlong.sort.parser.impl.FlinkSqlParser; @@ -46,7 +47,7 @@ import java.util.LinkedHashMap; import java.util.Map; import java.util.stream.Collectors; -public class AllMigrateWithSpecifyingFieldTest { +public class AllMigrateWithSpecifyingFieldTest extends AbstractTestBase { private MySqlExtractNode buildAllMigrateMySQLExtractNode() { @@ -123,4 +124,5 @@ public class AllMigrateWithSpecifyingFieldTest { ParseResult result = parser.parse(); Assert.assertTrue(result.tryExecute()); } + } diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ClickHouseSqlParserTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ClickHouseSqlParserTest.java index 19e9a5740..176aee8c1 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ClickHouseSqlParserTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ClickHouseSqlParserTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.sort.formats.common.LongFormatInfo; import org.apache.inlong.sort.formats.common.StringFormatInfo; import org.apache.inlong.sort.jdbc.dialect.ClickHouseDialect; @@ -46,7 +47,7 @@ import java.util.stream.Collectors; /** * Test for {@link ClickHouseLoadNode} and {@link ClickHouseDialect} */ -public class ClickHouseSqlParserTest { +public class ClickHouseSqlParserTest extends AbstractTestBase { public MySqlExtractNode buildMySQLExtractNode(String id) { List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new LongFormatInfo()), diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/DorisMultipleSinkTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/DorisMultipleSinkTest.java index 0d598bd5d..31db6249c 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/DorisMultipleSinkTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/DorisMultipleSinkTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.sort.formats.common.VarBinaryFormatInfo; import org.apache.inlong.sort.parser.impl.FlinkSqlParser; import org.apache.inlong.sort.parser.result.ParseResult; @@ -47,7 +48,7 @@ import java.util.stream.Collectors; /** * Test for {@link org.apache.inlong.sort.protocol.node.load.DorisLoadNode} */ -public class DorisMultipleSinkTest { +public class DorisMultipleSinkTest extends AbstractTestBase { private KafkaExtractNode buildKafkaExtractNode() { List<FieldInfo> fields = Collections.singletonList(new FieldInfo("raw", new VarBinaryFormatInfo())); diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ESMultipleSinkTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ESMultipleSinkTest.java index 842a574ac..0a23c4b48 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ESMultipleSinkTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/ESMultipleSinkTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.sort.formats.common.VarBinaryFormatInfo; import org.apache.inlong.sort.parser.impl.FlinkSqlParser; import org.apache.inlong.sort.parser.result.ParseResult; @@ -47,7 +48,7 @@ import java.util.stream.Collectors; /** * Test for {@link ElasticsearchLoadNode} */ -public class ESMultipleSinkTest { +public class ESMultipleSinkTest extends AbstractTestBase { private KafkaExtractNode buildKafkaExtractNode() { List<FieldInfo> fields = Collections.singletonList(new FieldInfo("raw", new VarBinaryFormatInfo())); diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/FilesystemSqlParserTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/FilesystemSqlParserTest.java index 68a62e9f3..4b959f79d 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/FilesystemSqlParserTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/FilesystemSqlParserTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.sort.formats.common.IntFormatInfo; import org.apache.inlong.sort.formats.common.LongFormatInfo; import org.apache.inlong.sort.formats.common.StringFormatInfo; @@ -46,7 +47,7 @@ import java.util.stream.Collectors; /** * Test for {@link FileSystemLoadNode} */ -public class FilesystemSqlParserTest { +public class FilesystemSqlParserTest extends AbstractTestBase { /** * build mysql extract node diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/NativeFlinkSqlParserTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/NativeFlinkSqlParserTest.java index a8f60f0be..2d661d85c 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/NativeFlinkSqlParserTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/NativeFlinkSqlParserTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.sort.parser.impl.NativeFlinkSqlParser; import org.apache.inlong.sort.parser.result.ParseResult; import org.junit.Assert; @@ -28,7 +29,7 @@ import org.junit.Test; /** * Test for {@link NativeFlinkSqlParser} */ -public class NativeFlinkSqlParserTest { +public class NativeFlinkSqlParserTest extends AbstractTestBase { @Test public void testNativeFlinkSqlParser() throws Exception { diff --git a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PulsarSqlParserTest.java b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PulsarSqlParserTest.java index 8c72bfbfb..f46778d82 100644 --- a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PulsarSqlParserTest.java +++ b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PulsarSqlParserTest.java @@ -20,6 +20,7 @@ package org.apache.inlong.sort.parser; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; +import org.apache.flink.test.util.AbstractTestBase; import org.apache.inlong.sort.parser.impl.FlinkSqlParser; import org.apache.inlong.sort.parser.result.ParseResult; import org.apache.inlong.sort.formats.common.LongFormatInfo; @@ -43,7 +44,7 @@ import java.util.Collections; import java.util.List; import java.util.stream.Collectors; -public class PulsarSqlParserTest { +public class PulsarSqlParserTest extends AbstractTestBase { private KafkaLoadNode buildKafkaLoadNode() { List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new LongFormatInfo()),
