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());
+ }
+ }
+}