This is an automated email from the ASF dual-hosted git repository.

dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 703eb79b1 [INLONG-4345][Sort] Sort support sink changelog stream to 
TDSQL Postgres (#4347)
703eb79b1 is described below

commit 703eb79b13cce3c1480414dcc4e4d3d2c16d2259
Author: pacino <[email protected]>
AuthorDate: Thu May 26 16:50:34 2022 +0800

    [INLONG-4345][Sort] Sort support sink changelog stream to TDSQL Postgres 
(#4347)
---
 .../sort/protocol/constant/PostgresConstant.java   |   4 +-
 .../apache/inlong/sort/protocol/node/LoadNode.java |   4 +-
 .../org/apache/inlong/sort/protocol/node/Node.java |   6 +-
 .../sort/protocol/node/load/PostgresLoadNode.java  |   2 +
 ...resLoadNode.java => TDSQLPostgresLoadNode.java} |  19 ++-
 .../node/load/TDSQLPostgresLoadNodeTest.java       |  50 +++++++
 .../sort/jdbc/dialect/TDSQLPostgresDialect.java    | 150 +++++++++++++++++++++
 .../inlong/sort/jdbc/table/JdbcDialects.java       |   4 +-
 .../parser/PostgresLoadNodeFlinkSqlParseTest.java  |   4 +-
 ...=> TDSQLPostgresLoadNodeFlinkSqlParseTest.java} |  31 ++---
 10 files changed, 244 insertions(+), 30 deletions(-)

diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/constant/PostgresConstant.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/constant/PostgresConstant.java
index e43907bf8..e4fa27505 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/constant/PostgresConstant.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/constant/PostgresConstant.java
@@ -49,6 +49,8 @@ public class PostgresConstant {
 
     public static final String URL = "url";
 
-    public static final String JDBC = "jdbc-inlong";
+    public static final String JDBC_INLONG = "jdbc-inlong";
+
+    public static final String JDBC = "jdbc";
 
 }
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
index 106ab22f5..eb9804664 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
@@ -35,6 +35,7 @@ import org.apache.inlong.sort.protocol.node.load.HiveLoadNode;
 import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
 import org.apache.inlong.sort.protocol.node.load.PostgresLoadNode;
 import org.apache.inlong.sort.protocol.node.load.SqlServerLoadNode;
+import org.apache.inlong.sort.protocol.node.load.TDSQLPostgresLoadNode;
 import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
 import org.apache.inlong.sort.protocol.transformation.FilterFunction;
 
@@ -56,7 +57,8 @@ import java.util.Map;
         @JsonSubTypes.Type(value = FileSystemLoadNode.class, name = 
"fileSystemLoad"),
         @JsonSubTypes.Type(value = PostgresLoadNode.class, name = 
"postgresLoad"),
         @JsonSubTypes.Type(value = ClickHouseLoadNode.class, name = 
"clickHouseLoad"),
-        @JsonSubTypes.Type(value = SqlServerLoadNode.class, name = 
"sqlserverLoad")
+        @JsonSubTypes.Type(value = SqlServerLoadNode.class, name = 
"sqlserverLoad"),
+        @JsonSubTypes.Type(value = TDSQLPostgresLoadNode.class, name = 
"tdsqlPostgresLoad"),
 })
 @NoArgsConstructor
 @Data
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/Node.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/Node.java
index 99363bc2a..44cc50098 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/Node.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/Node.java
@@ -28,15 +28,16 @@ import 
org.apache.inlong.sort.protocol.node.extract.MongoExtractNode;
 import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
 import org.apache.inlong.sort.protocol.node.extract.OracleExtractNode;
 import org.apache.inlong.sort.protocol.node.extract.PostgresExtractNode;
-import org.apache.inlong.sort.protocol.node.load.ClickHouseLoadNode;
 import org.apache.inlong.sort.protocol.node.extract.PulsarExtractNode;
 import org.apache.inlong.sort.protocol.node.extract.SqlServerExtractNode;
+import org.apache.inlong.sort.protocol.node.load.ClickHouseLoadNode;
 import org.apache.inlong.sort.protocol.node.load.FileSystemLoadNode;
 import org.apache.inlong.sort.protocol.node.load.HbaseLoadNode;
 import org.apache.inlong.sort.protocol.node.load.HiveLoadNode;
 import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
 import org.apache.inlong.sort.protocol.node.load.PostgresLoadNode;
 import org.apache.inlong.sort.protocol.node.load.SqlServerLoadNode;
+import org.apache.inlong.sort.protocol.node.load.TDSQLPostgresLoadNode;
 import org.apache.inlong.sort.protocol.node.transform.DistinctNode;
 import org.apache.inlong.sort.protocol.node.transform.TransformNode;
 
@@ -69,7 +70,8 @@ import java.util.TreeMap;
         @JsonSubTypes.Type(value = PostgresLoadNode.class, name = 
"postgresLoad"),
         @JsonSubTypes.Type(value = FileSystemLoadNode.class, name = 
"fileSystemLoad"),
         @JsonSubTypes.Type(value = ClickHouseLoadNode.class, name = 
"clickHouseLoad"),
-        @JsonSubTypes.Type(value = SqlServerLoadNode.class, name = 
"sqlserverLoad")
+        @JsonSubTypes.Type(value = SqlServerLoadNode.class, name = 
"sqlserverLoad"),
+        @JsonSubTypes.Type(value = TDSQLPostgresLoadNode.class, name = 
"tdsqlPostgresLoad"),
 })
 public interface Node {
 
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
index 2e6549aa7..214479e0b 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
@@ -22,6 +22,7 @@ import com.google.common.base.Preconditions;
 import lombok.Data;
 import lombok.EqualsAndHashCode;
 import lombok.NoArgsConstructor;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTypeName;
 import org.apache.inlong.sort.protocol.FieldInfo;
@@ -67,6 +68,7 @@ public class PostgresLoadNode extends LoadNode implements 
Serializable {
     @JsonProperty("primaryKey")
     private String primaryKey;
 
+    @JsonCreator
     public PostgresLoadNode(
             @JsonProperty("id") String id,
             @JsonProperty("name") String name,
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/TDSQLPostgresLoadNode.java
similarity index 85%
copy from 
inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
copy to 
inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/TDSQLPostgresLoadNode.java
index 2e6549aa7..bd2e3960e 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/PostgresLoadNode.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/TDSQLPostgresLoadNode.java
@@ -22,6 +22,7 @@ import com.google.common.base.Preconditions;
 import lombok.Data;
 import lombok.EqualsAndHashCode;
 import lombok.NoArgsConstructor;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTypeName;
 import org.apache.inlong.sort.protocol.FieldInfo;
@@ -36,13 +37,17 @@ import java.util.List;
 import java.util.Map;
 
 /**
- * Postgres load node can load data into Postgres
+ * TDSQLPostgres load node can load data into TDSQL Postgres
+ * @see <a herf="https://cloud.tencent.com/product/tbase";>TDSQL Postgres</a>
+ * TDSQL Postgres is an enterprise-level distributed HTAP database. Through a 
single database cluster to provide
+ * users with highly consistent distributed database services and 
high-performance data warehouse services,
+ * a set of integrated enterprise-level solutions is formed.
  */
 @EqualsAndHashCode(callSuper = true)
-@JsonTypeName("postgresLoad")
+@JsonTypeName("tdsqlPostgresLoad")
 @Data
 @NoArgsConstructor
-public class PostgresLoadNode extends LoadNode implements Serializable {
+public class TDSQLPostgresLoadNode extends LoadNode implements Serializable {
 
     private static final long serialVersionUID = 1L;
 
@@ -61,13 +66,13 @@ public class PostgresLoadNode extends LoadNode implements 
Serializable {
     @JsonProperty("tableName")
     private String tableName;
     /**
-     * Please declare primary key for sink table when query contains 
update/delete record if your version support
-     * upsert. You can change source stream mode is "append" if your version 
can't support upsert.
+     * Please declare primary key for sink table when query contains 
update/delete record.
      */
     @JsonProperty("primaryKey")
     private String primaryKey;
 
-    public PostgresLoadNode(
+    @JsonCreator
+    public TDSQLPostgresLoadNode(
             @JsonProperty("id") String id,
             @JsonProperty("name") String name,
             @JsonProperty("fields") List<FieldInfo> fields,
@@ -92,7 +97,7 @@ public class PostgresLoadNode extends LoadNode implements 
Serializable {
     @Override
     public Map<String, String> tableOptions() {
         Map<String, String> options = super.tableOptions();
-        options.put(PostgresConstant.CONNECTOR, PostgresConstant.JDBC);
+        options.put(PostgresConstant.CONNECTOR, PostgresConstant.JDBC_INLONG);
         options.put(PostgresConstant.URL, url);
         options.put(PostgresConstant.USERNAME, username);
         options.put(PostgresConstant.PASSWORD, password);
diff --git 
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/TDSQLPostgresLoadNodeTest.java
 
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/TDSQLPostgresLoadNodeTest.java
new file mode 100644
index 000000000..b6a26319b
--- /dev/null
+++ 
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/TDSQLPostgresLoadNodeTest.java
@@ -0,0 +1,50 @@
+/*
+ *   Licensed to the Apache Software Foundation (ASF) under one
+ *   or more contributor license agreements.  See the NOTICE file
+ *   distributed with this work for additional information
+ *   regarding copyright ownership.  The ASF licenses this file
+ *   to you under the Apache License, Version 2.0 (the
+ *   "License"); you may not use this file except in compliance
+ *   with the License.  You may obtain a copy of the License at
+ *
+ *       http://www.apache.org/licenses/LICENSE-2.0
+ *
+ *   Unless required by applicable law or agreed to in writing, software
+ *   distributed under the License is distributed on an "AS IS" BASIS,
+ *   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ *   See the License for the specific language governing permissions and
+ *   limitations under the License.
+ */
+
+package org.apache.inlong.sort.protocol.node.load;
+
+import org.apache.inlong.sort.SerializeBaseTest;
+import org.apache.inlong.sort.formats.common.StringFormatInfo;
+import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
+
+import java.util.Arrays;
+
+/**
+ * Test for {@link TDSQLPostgresLoadNode}
+ */
+public class TDSQLPostgresLoadNodeTest extends 
SerializeBaseTest<TDSQLPostgresLoadNode> {
+
+    /**
+     * Get test object
+     *
+     * @return The test object
+     */
+    @Override
+    public TDSQLPostgresLoadNode getTestObject() {
+        return new TDSQLPostgresLoadNode("1", "tdsqlPostgres_output", 
Arrays.asList(new FieldInfo("name",
+                new StringFormatInfo())),
+                Arrays.asList(new FieldRelationShip(new FieldInfo("name", new 
StringFormatInfo()),
+                        new FieldInfo("name", new StringFormatInfo()))), null, 
null, 1, null,
+                "jdbc:postgresql://localhost:5432/postgres",
+                "postgres",
+                "inlong",
+                "user",
+                "name");
+    }
+}
diff --git 
a/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/dialect/TDSQLPostgresDialect.java
 
b/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/dialect/TDSQLPostgresDialect.java
new file mode 100644
index 000000000..3984426a0
--- /dev/null
+++ 
b/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/dialect/TDSQLPostgresDialect.java
@@ -0,0 +1,150 @@
+/*
+ *   Licensed to the Apache Software Foundation (ASF) under one
+ *   or more contributor license agreements.  See the NOTICE file
+ *   distributed with this work for additional information
+ *   regarding copyright ownership.  The ASF licenses this file
+ *   to you under the Apache License, Version 2.0 (the
+ *   "License"); you may not use this file except in compliance
+ *   with the License.  You may obtain a copy of the License at
+ *
+ *       http://www.apache.org/licenses/LICENSE-2.0
+ *
+ *   Unless required by applicable law or agreed to in writing, software
+ *   distributed under the License is distributed on an "AS IS" BASIS,
+ *   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ *   See the License for the specific language governing permissions and
+ *   limitations under the License.
+ */
+
+package org.apache.inlong.sort.jdbc.dialect;
+
+import org.apache.commons.lang3.StringUtils;
+import org.apache.flink.connector.jdbc.internal.converter.JdbcRowConverter;
+import org.apache.flink.connector.jdbc.internal.converter.PostgresRowConverter;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.RowType;
+import org.apache.inlong.sort.jdbc.table.AbstractJdbcDialect;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.Optional;
+import java.util.stream.Collectors;
+
+/**
+ * copy from flink-jdbc
+ * modify upsert statement implement
+ */
+public class TDSQLPostgresDialect extends AbstractJdbcDialect {
+
+    private static final long serialVersionUID = 1L;
+
+    // Define MAX/MIN precision of TIMESTAMP type according to PostgreSQL docs:
+    // https://www.postgresql.org/docs/12/datatype-datetime.html
+    private static final int MAX_TIMESTAMP_PRECISION = 6;
+    private static final int MIN_TIMESTAMP_PRECISION = 1;
+
+    // Define MAX/MIN precision of DECIMAL type according to PostgreSQL docs:
+    // 
https://www.postgresql.org/docs/12/datatype-numeric.html#DATATYPE-NUMERIC-DECIMAL
+    private static final int MAX_DECIMAL_PRECISION = 1000;
+    private static final int MIN_DECIMAL_PRECISION = 1;
+
+    @Override
+    public boolean canHandle(String url) {
+        return url.startsWith("jdbc:postgresql:");
+    }
+
+    @Override
+    public JdbcRowConverter getRowConverter(RowType rowType) {
+        return new PostgresRowConverter(rowType);
+    }
+
+    @Override
+    public String getLimitClause(long limit) {
+        return "LIMIT " + limit;
+    }
+
+    @Override
+    public Optional<String> defaultDriverName() {
+        return Optional.of("org.postgresql.Driver");
+    }
+
+    /** Postgres upsert query. It use ON CONFLICT ... DO UPDATE SET.. to 
replace into Postgres. */
+    @Override
+    public Optional<String> getUpsertStatement(
+            String tableName, String[] fieldNames, String[] uniqueKeyFields) {
+        String uniqueColumns =
+                Arrays.stream(uniqueKeyFields)
+                        .map(this::quoteIdentifier)
+                        .collect(Collectors.joining(", "));
+        List<String> uniqueKeyFieldList = Arrays.asList(uniqueKeyFields);
+        String updateClause =
+                Arrays.stream(fieldNames)
+                        .filter(f -> !uniqueKeyFieldList.contains(f))
+                        .map(f -> quoteIdentifier(f) + "=EXCLUDED." + 
quoteIdentifier(f))
+                        .collect(Collectors.joining(", "));
+        String str = getInsertIntoStatement(tableName, fieldNames);
+        if (StringUtils.isNotEmpty(updateClause)) {
+            str = str + " ON CONFLICT ("
+                    + uniqueColumns
+                    + ")"
+                    + " DO UPDATE SET "
+                    + updateClause;
+        }
+        return Optional.of(str);
+    }
+
+    @Override
+    public String quoteIdentifier(String identifier) {
+        return identifier;
+    }
+
+    @Override
+    public String dialectName() {
+        return "PostgreSQL";
+    }
+
+    @Override
+    public int maxDecimalPrecision() {
+        return MAX_DECIMAL_PRECISION;
+    }
+
+    @Override
+    public int minDecimalPrecision() {
+        return MIN_DECIMAL_PRECISION;
+    }
+
+    @Override
+    public int maxTimestampPrecision() {
+        return MAX_TIMESTAMP_PRECISION;
+    }
+
+    @Override
+    public int minTimestampPrecision() {
+        return MIN_TIMESTAMP_PRECISION;
+    }
+
+    @Override
+    public List<LogicalTypeRoot> unsupportedTypes() {
+        // The data types used in PostgreSQL are list at:
+        // https://www.postgresql.org/docs/12/datatype.html
+
+        // TODO: We can't convert BINARY data type to
+        //  PrimitiveArrayTypeInfo.BYTE_PRIMITIVE_ARRAY_TYPE_INFO in
+        // LegacyTypeInfoDataTypeConverter.
+        return Arrays.asList(
+                LogicalTypeRoot.BINARY,
+                LogicalTypeRoot.TIMESTAMP_WITH_TIME_ZONE,
+                LogicalTypeRoot.INTERVAL_YEAR_MONTH,
+                LogicalTypeRoot.INTERVAL_DAY_TIME,
+                LogicalTypeRoot.MULTISET,
+                LogicalTypeRoot.MAP,
+                LogicalTypeRoot.ROW,
+                LogicalTypeRoot.DISTINCT_TYPE,
+                LogicalTypeRoot.STRUCTURED_TYPE,
+                LogicalTypeRoot.NULL,
+                LogicalTypeRoot.RAW,
+                LogicalTypeRoot.SYMBOL,
+                LogicalTypeRoot.UNRESOLVED);
+    }
+
+}
diff --git 
a/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/table/JdbcDialects.java
 
b/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/table/JdbcDialects.java
index 3d7aab912..a9d473661 100644
--- 
a/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/table/JdbcDialects.java
+++ 
b/inlong-sort/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/table/JdbcDialects.java
@@ -20,8 +20,8 @@ package org.apache.inlong.sort.jdbc.table;
 
 import org.apache.flink.connector.jdbc.dialect.JdbcDialect;
 import org.apache.flink.connector.jdbc.dialect.MySQLDialect;
-import org.apache.flink.connector.jdbc.dialect.PostgresDialect;
 import org.apache.inlong.sort.jdbc.dialect.SqlServerDialect;
+import org.apache.inlong.sort.jdbc.dialect.TDSQLPostgresDialect;
 
 import java.util.ArrayList;
 import java.util.List;
@@ -36,7 +36,7 @@ public final class JdbcDialects {
 
     static {
         DIALECTS.add(new MySQLDialect());
-        DIALECTS.add(new PostgresDialect());
+        DIALECTS.add(new TDSQLPostgresDialect());
         DIALECTS.add(new SqlServerDialect());
     }
 
diff --git 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
index 5fccfb00a..49317f3f1 100644
--- 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
+++ 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
@@ -22,10 +22,10 @@ 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.IntFormatInfo;
 import org.apache.inlong.sort.formats.common.StringFormatInfo;
+import org.apache.inlong.sort.parser.impl.FlinkSqlParser;
+import org.apache.inlong.sort.parser.result.ParseResult;
 import org.apache.inlong.sort.protocol.FieldInfo;
 import org.apache.inlong.sort.protocol.GroupInfo;
 import org.apache.inlong.sort.protocol.StreamInfo;
diff --git 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/TDSQLPostgresLoadNodeFlinkSqlParseTest.java
similarity index 87%
copy from 
inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
copy to 
inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/TDSQLPostgresLoadNodeFlinkSqlParseTest.java
index 5fccfb00a..22033f8f5 100644
--- 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/PostgresLoadNodeFlinkSqlParseTest.java
+++ 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/TDSQLPostgresLoadNodeFlinkSqlParseTest.java
@@ -22,16 +22,16 @@ 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.IntFormatInfo;
 import org.apache.inlong.sort.formats.common.StringFormatInfo;
+import org.apache.inlong.sort.parser.impl.FlinkSqlParser;
+import org.apache.inlong.sort.parser.result.ParseResult;
 import org.apache.inlong.sort.protocol.FieldInfo;
 import org.apache.inlong.sort.protocol.GroupInfo;
 import org.apache.inlong.sort.protocol.StreamInfo;
 import org.apache.inlong.sort.protocol.node.Node;
 import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
-import org.apache.inlong.sort.protocol.node.load.PostgresLoadNode;
+import org.apache.inlong.sort.protocol.node.load.TDSQLPostgresLoadNode;
 import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
 import 
org.apache.inlong.sort.protocol.transformation.relation.NodeRelationShip;
 import org.junit.Assert;
@@ -45,9 +45,9 @@ import java.util.Map;
 import java.util.stream.Collectors;
 
 /**
- * Test for {@link PostgresLoadNode}
+ * Test for {@link TDSQLPostgresLoadNode}
  */
-public class PostgresLoadNodeFlinkSqlParseTest extends AbstractTestBase {
+public class TDSQLPostgresLoadNodeFlinkSqlParseTest extends AbstractTestBase {
 
     /**
      * build mysql extract node
@@ -66,22 +66,22 @@ public class PostgresLoadNodeFlinkSqlParseTest extends 
AbstractTestBase {
     }
 
     /**
-     * build postgres load node
+     * build  load node
      *
-     * @return postgres load node
+     * @return TDSQL Postgres load node
      */
-    private PostgresLoadNode buildPostgresLoadNode() {
-        return new PostgresLoadNode("2", "postgres_output", Arrays.asList(new 
FieldInfo("name",
+    private TDSQLPostgresLoadNode buildTDSQLPostgresLoadNode() {
+        return new TDSQLPostgresLoadNode("2", "tdsqlPostgres_output", 
Arrays.asList(new FieldInfo("name",
                 new StringFormatInfo()), new FieldInfo("age", new 
IntFormatInfo())),
                 Arrays.asList(new FieldRelationShip(new FieldInfo("name", new 
StringFormatInfo()),
                                 new FieldInfo("name", new StringFormatInfo())),
                         new FieldRelationShip(new FieldInfo("age", new 
IntFormatInfo()),
                                 new FieldInfo("age", new IntFormatInfo()))), 
null, null, 1, null,
-                "jdbc:postgresql://localhost:5432/postgres",
-                "postgres",
+                "jdbc:postgresql://localhost:5432/tdsql",
+                "tdsqlpostgres",
                 "inlong",
-                "public.user",
-                "name,age");
+                "public.test",
+                "name");
     }
 
     /**
@@ -98,7 +98,8 @@ public class PostgresLoadNodeFlinkSqlParseTest extends 
AbstractTestBase {
     }
 
     /**
-     * Test flink sql task for extract is mysql {@link MySqlExtractNode} and 
load is postgres {@link PostgresLoadNode}
+     * Test flink sql task for extract is mysql {@link MySqlExtractNode} and 
load is tdsql postgres
+     * {@link TDSQLPostgresLoadNode}
      *
      * @throws Exception The exception may be thrown when executing
      */
@@ -115,7 +116,7 @@ public class PostgresLoadNodeFlinkSqlParseTest extends 
AbstractTestBase {
                 .build();
         StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, 
settings);
         Node inputNode = buildMySQLExtractNode();
-        Node outputNode = buildPostgresLoadNode();
+        Node outputNode = buildTDSQLPostgresLoadNode();
         StreamInfo streamInfo = new StreamInfo("1", Arrays.asList(inputNode, 
outputNode),
                 
Collections.singletonList(buildNodeRelation(Collections.singletonList(inputNode),
                         Collections.singletonList(outputNode))));

Reply via email to