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

healchow 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 b91fa37a4 [INLONG-4665][Sort] The primary key of the MongoDB CDC 
connector must be _id (#4666)
b91fa37a4 is described below

commit b91fa37a49c2293c725e0a9e3911e4a1967ebbf2
Author: Oneal65 <[email protected]>
AuthorDate: Thu Jun 16 12:15:57 2022 +0800

    [INLONG-4665][Sort] The primary key of the MongoDB CDC connector must be 
_id (#4666)
---
 .../service/sort/util/ExtractNodeUtils.java        |  1 -
 .../protocol/node/extract/MongoExtractNode.java    | 35 +++++++++++++++-------
 .../node/extract/MongoExtractNodeTest.java         |  2 +-
 .../sort/parser/MongoExtractFlinkSqlParseTest.java |  9 ++----
 4 files changed, 28 insertions(+), 19 deletions(-)

diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/ExtractNodeUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/ExtractNodeUtils.java
index 406376798..f23c31f2d 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/ExtractNodeUtils.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/ExtractNodeUtils.java
@@ -385,7 +385,6 @@ public class ExtractNodeUtils {
                 fieldInfos,
                 null,
                 properties,
-                source.getPrimaryKey(),
                 source.getCollection(),
                 source.getHosts(),
                 source.getUsername(),
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNode.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNode.java
index a559d7883..ed3f18064 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNode.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNode.java
@@ -20,6 +20,7 @@ package org.apache.inlong.sort.protocol.node.extract;
 
 import com.google.common.base.Preconditions;
 import java.io.Serializable;
+import java.util.ArrayList;
 import java.util.List;
 import java.util.Map;
 import javax.annotation.Nonnull;
@@ -31,6 +32,7 @@ import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonInc
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonInclude.Include;
 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.formats.common.StringFormatInfo;
 import org.apache.inlong.sort.protocol.FieldInfo;
 import org.apache.inlong.sort.protocol.node.ExtractNode;
 import org.apache.inlong.sort.protocol.transformation.WatermarkField;
@@ -45,6 +47,13 @@ public class MongoExtractNode extends ExtractNode implements 
Serializable {
 
     private static final long serialVersionUID = 1L;
 
+    /**
+     * the primary key must be "_id"
+     *
+     * @see <a 
href="https://ververica.github.io/flink-cdc-connectors/release-2.1/content/connectors/mongodb-cdc.html#";>MongoDB
 CDC Connector</a>
+     */
+    private static final String ID = "_id";
+
     @JsonInclude(Include.NON_NULL)
     @JsonProperty("primaryKey")
     private String primaryKey;
@@ -61,23 +70,27 @@ public class MongoExtractNode extends ExtractNode 
implements Serializable {
 
     @JsonCreator
     public MongoExtractNode(@JsonProperty("id") String id,
-        @JsonProperty("name") String name,
-        @JsonProperty("fields") List<FieldInfo> fields,
-        @Nullable @JsonProperty("watermarkField") WatermarkField 
waterMarkField,
-        @JsonProperty("properties") Map<String, String> properties,
-        @JsonProperty("primaryKey") String primaryKey,
-        @JsonProperty("collection") @Nonnull String collection,
-        @JsonProperty("hostname") String hostname,
-        @JsonProperty("username") String username,
-        @JsonProperty("password") String password,
-        @JsonProperty("database") String database) {
+            @JsonProperty("name") String name,
+            @JsonProperty("fields") List<FieldInfo> fields,
+            @Nullable @JsonProperty("watermarkField") WatermarkField 
waterMarkField,
+            @JsonProperty("properties") Map<String, String> properties,
+            @JsonProperty("collection") @Nonnull String collection,
+            @JsonProperty("hostname") String hostname,
+            @JsonProperty("username") String username,
+            @JsonProperty("password") String password,
+            @JsonProperty("database") String database) {
         super(id, name, fields, waterMarkField, properties);
+        if (fields.stream().noneMatch(m -> m.getName().equals(ID))) {
+            List<FieldInfo> allFields = new ArrayList<>(fields);
+            allFields.add(new FieldInfo(ID, new StringFormatInfo()));
+            this.setFields(allFields);
+        }
         this.collection = Preconditions.checkNotNull(collection, "collection 
is null");
         this.hosts = Preconditions.checkNotNull(hostname, "hostname is null");
         this.username = Preconditions.checkNotNull(username, "username is 
null");
         this.password = Preconditions.checkNotNull(password, "password is 
null");
         this.database = Preconditions.checkNotNull(database, "database is 
null");
-        this.primaryKey = primaryKey;
+        this.primaryKey = ID;
     }
 
     @Override
diff --git 
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNodeTest.java
 
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNodeTest.java
index 478bc799b..c1e90bb21 100644
--- 
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNodeTest.java
+++ 
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/extract/MongoExtractNodeTest.java
@@ -35,7 +35,7 @@ public class MongoExtractNodeTest extends 
SerializeBaseTest<MongoExtractNode>  {
             new FieldInfo("age", new IntFormatInfo()));
         return new MongoExtractNode(
             "1", "test", fields,  null, null,
-            "id", "test", "localhost", "inlong", "password", "test"
+                "test", "localhost", "inlong", "password", "test"
         );
     }
 
diff --git 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
index 7ed02fcd7..e62a72bae 100644
--- 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
+++ 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
@@ -48,13 +48,10 @@ import java.util.stream.Collectors;
 public class MongoExtractFlinkSqlParseTest extends AbstractTestBase {
 
     private MongoExtractNode buildMongoNode() {
-        List<FieldInfo> fields = Arrays.asList(
-                new FieldInfo("name", new StringFormatInfo()),
-                new FieldInfo("_id", new StringFormatInfo()));
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("name", new 
StringFormatInfo()));
         return new MongoExtractNode("1", "mysql_input", fields,
-                null, null, "_id",
-                "test", "localhost:27017", "root", "inlong",
-                "test");
+                null, null, "test", "localhost:27017",
+                "root", "inlong", "test");
     }
 
     private KafkaLoadNode buildAllMigrateKafkaNode() {

Reply via email to