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

JNSimba pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris-flink-connector.git


The following commit(s) were added to refs/heads/master by this push:
     new 79ae9b03 [Fix](mongodb-cdc) Fix fullDocument JSON string parsing error 
(#671)
79ae9b03 is described below

commit 79ae9b03fc993ea7167b0aefa8547a211d7058fc
Author: liziyan <[email protected]>
AuthorDate: Tue Aug 18 14:02:24 2026 +0800

    [Fix](mongodb-cdc) Fix fullDocument JSON string parsing error (#671)
    
    Debezium may serialize fullDocument as a JSON string (TextNode) instead
    of an ObjectNode. Added a check to parse the textual fullDocument before
    extracting row data, preventing ClassCastException or deserialization 
errors.
---
 .../tools/cdc/mongodb/serializer/MongoJsonDebeziumDataChange.java | 8 ++++++++
 1 file changed, 8 insertions(+)

diff --git 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/tools/cdc/mongodb/serializer/MongoJsonDebeziumDataChange.java
 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/tools/cdc/mongodb/serializer/MongoJsonDebeziumDataChange.java
index 12fa8581..fe13f293 100644
--- 
a/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/tools/cdc/mongodb/serializer/MongoJsonDebeziumDataChange.java
+++ 
b/flink-doris-connector/flink-doris-connector-flink1/src/main/java/org/apache/doris/flink/tools/cdc/mongodb/serializer/MongoJsonDebeziumDataChange.java
@@ -24,6 +24,7 @@ import com.fasterxml.jackson.databind.JsonNode;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import com.fasterxml.jackson.databind.node.NullNode;
 import org.apache.doris.flink.cfg.DorisOptions;
+import org.apache.doris.flink.exception.DorisRuntimeException;
 import org.apache.doris.flink.sink.writer.ChangeEvent;
 import org.apache.doris.flink.sink.writer.serializer.DorisRecord;
 import 
org.apache.doris.flink.sink.writer.serializer.jsondebezium.CdcDataChange;
@@ -127,6 +128,13 @@ public class MongoJsonDebeziumDataChange extends 
CdcDataChange implements Change
     @Override
     public Map<String, Object> extractAfterRow(JsonNode recordRoot) {
         JsonNode dataNode = recordRoot.get(FIELD_DATA);
+        if (dataNode != null && dataNode.isTextual()) {
+            try {
+                dataNode = objectMapper.readTree(dataNode.asText());
+            } catch (IOException e) {
+                throw new DorisRuntimeException("Failed to parse fullDocument 
JSON", e);
+            }
+        }
         return JsonNodeExtractUtil.extractAfterRow(dataNode, objectMapper);
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to