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

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new e6caf5a4b3 [Fix][Connector-V2] Preserve MongoDB CDC snapshot failure 
cause (#11880)
e6caf5a4b3 is described below

commit e6caf5a4b337803d866e877dab86541e6d15a364
Author: Goutam Adwant <[email protected]>
AuthorDate: Wed Aug 26 14:58:06 2026 +0000

    [Fix][Connector-V2] Preserve MongoDB CDC snapshot failure cause (#11880)
    
    Signed-off-by: goutamadwant <[email protected]>
---
 .../exception/MongodbConnectorException.java       |   5 ++
 .../mongodb/source/fetch/MongodbScanFetchTask.java |   3 +-
 .../source/fetch/MongodbScanFetchTaskTest.java     | 100 +++++++++++++++++++++
 3 files changed, 107 insertions(+), 1 deletion(-)

diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/exception/MongodbConnectorException.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/exception/MongodbConnectorException.java
index 2d2267e478..f16ae1f291 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/exception/MongodbConnectorException.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/exception/MongodbConnectorException.java
@@ -25,4 +25,9 @@ public class MongodbConnectorException extends 
SeaTunnelRuntimeException {
     public MongodbConnectorException(SeaTunnelErrorCode seaTunnelErrorCode, 
String errorMessage) {
         super(seaTunnelErrorCode, errorMessage);
     }
+
+    public MongodbConnectorException(
+            SeaTunnelErrorCode seaTunnelErrorCode, String errorMessage, 
Throwable cause) {
+        super(seaTunnelErrorCode, errorMessage, cause);
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTask.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTask.java
index d992a8f345..9bd9cabb70 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTask.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTask.java
@@ -156,7 +156,8 @@ public class MongodbScanFetchTask implements 
FetchTask<SourceSplitBase> {
                     ILLEGAL_ARGUMENT,
                     String.format(
                             "Execute snapshot read subtask for mongodb split 
%s fail",
-                            snapshotSplit));
+                            snapshotSplit),
+                    e);
         } finally {
             taskRunning = false;
         }
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTaskTest.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTaskTest.java
new file mode 100644
index 0000000000..9da1f0b8dd
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbScanFetchTaskTest.java
@@ -0,0 +1,100 @@
+/*
+ * 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.seatunnel.connectors.seatunnel.cdc.mongodb.source.fetch;
+
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.connectors.cdc.base.source.split.SnapshotSplit;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.config.MongodbSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.exception.MongodbConnectorException;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.source.dialect.MongodbDialect;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.source.offset.ChangeStreamOffset;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.utils.MongodbUtils;
+
+import org.bson.BsonTimestamp;
+import org.bson.RawBsonDocument;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import com.mongodb.MongoException;
+import com.mongodb.client.MongoClient;
+import com.mongodb.client.MongoCollection;
+import io.debezium.connector.base.ChangeEventQueue;
+import io.debezium.pipeline.DataChangeEvent;
+import io.debezium.relational.TableId;
+
+import static org.apache.seatunnel.api.table.type.BasicType.INT_TYPE;
+import static 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.utils.ChunkUtils.maxUpperBoundOfId;
+import static 
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.utils.ChunkUtils.minLowerBoundOfId;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.when;
+
+class MongodbScanFetchTaskTest {
+
+    @Test
+    void shouldPreserveSnapshotReadFailureCause() throws Exception {
+        TableId collectionId = new TableId("inventory", null, "products");
+        SnapshotSplit snapshotSplit =
+                new SnapshotSplit(
+                        "inventory.products:0",
+                        collectionId,
+                        new SeaTunnelRowType(
+                                new String[] {"_id"}, new 
SeaTunnelDataType<?>[] {INT_TYPE}),
+                        minLowerBoundOfId(),
+                        maxUpperBoundOfId());
+        MongodbScanFetchTask fetchTask = new 
MongodbScanFetchTask(snapshotSplit);
+
+        MongodbFetchTaskContext taskContext = 
mock(MongodbFetchTaskContext.class);
+        MongodbSourceConfig sourceConfig = mock(MongodbSourceConfig.class);
+        MongodbDialect dialect = mock(MongodbDialect.class);
+        MongoClient mongoClient = mock(MongoClient.class);
+        @SuppressWarnings("unchecked")
+        ChangeEventQueue<DataChangeEvent> queue = mock(ChangeEventQueue.class);
+        @SuppressWarnings("unchecked")
+        MongoCollection<RawBsonDocument> collection = 
mock(MongoCollection.class);
+        ChangeStreamOffset lowWatermark = new ChangeStreamOffset(new 
BsonTimestamp(1));
+        MongoException snapshotReadFailure = new MongoException(2, "invalid 
snapshot bounds");
+
+        when(taskContext.getSourceConfig()).thenReturn(sourceConfig);
+        when(taskContext.getDialect()).thenReturn(dialect);
+        when(taskContext.getQueue()).thenReturn(queue);
+        when(taskContext.getMongoClient()).thenReturn(mongoClient);
+        
when(dialect.displayCurrentOffset(sourceConfig)).thenReturn(lowWatermark);
+        when(collection.find()).thenThrow(snapshotReadFailure);
+
+        try (MockedStatic<MongodbUtils> mongodbUtils = 
mockStatic(MongodbUtils.class)) {
+            mongodbUtils
+                    .when(
+                            () ->
+                                    MongodbUtils.getMongoCollection(
+                                            mongoClient, collectionId, 
RawBsonDocument.class))
+                    .thenReturn(collection);
+
+            MongodbConnectorException actual =
+                    assertThrows(
+                            MongodbConnectorException.class, () -> 
fetchTask.execute(taskContext));
+
+            assertSame(snapshotReadFailure, actual.getCause());
+            assertFalse(fetchTask.isRunning());
+        }
+    }
+}

Reply via email to