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() {