{"record":{"id":"7e6f513cd70aa9ba","repo":"apache/seatunnel","slug":"cannot-extract-clustertime-from-change-stream-even","errorCode":null,"errorMessage":"Cannot extract clusterTime from change stream event, fallback to current timestamp.","messagePattern":"Cannot extract clusterTime from change stream event, fallback to current timestamp\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbStreamFetchTask.java","lineNumber":365,"sourceCode":"        }\n        return null;\n    }\n\n    @Nonnull\n    private BsonDocument normalizeChangeStreamDocument(@Nonnull BsonDocument changeStreamDocument) {\n        // _id: primary key of change document.\n        BsonDocument normalizedDocument = normalizeKeyDocument(changeStreamDocument);\n        changeStreamDocument.put(ID_FIELD, normalizedDocument);\n\n        // ts_ms: It indicates the time at which the reader processed the event.\n        changeStreamDocument.put(TS_MS_FIELD, new BsonInt64(System.currentTimeMillis()));\n\n        // source\n        BsonDocument source = new BsonDocument();\n        source.put(SNAPSHOT_FIELD, new BsonString(FALSE_FALSE));\n\n        if (!changeStreamDocument.containsKey(CLUSTER_TIME_FIELD)) {\n            log.warn(\n                    \"Cannot extract clusterTime from change stream event, fallback to current timestamp.\");\n            changeStreamDocument.put(CLUSTER_TIME_FIELD, currentBsonTimestamp());\n        }\n\n        // source.ts_ms\n        // It indicates the time that the change was made in the database. If the record is read\n        // from snapshot of the table instead of the change stream, the value is always 0.\n        BsonTimestamp clusterTime = changeStreamDocument.getTimestamp(CLUSTER_TIME_FIELD);\n        Instant clusterInstant = Instant.ofEpochSecond(clusterTime.getTime());\n        source.put(TS_MS_FIELD, new BsonInt64(clusterInstant.toEpochMilli()));\n        changeStreamDocument.put(SOURCE_FIELD, source);\n\n        return changeStreamDocument;\n    }\n\n    @Nonnull\n    private BsonDocument normalizeKeyDocument(@Nonnull BsonDocument changeStreamDocument) {\n        BsonDocument documentKey = changeStreamDocument.getDocument(DOCUMENT_KEY);","sourceCodeStart":347,"sourceCodeEnd":383,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbStreamFetchTask.java#L347-L383","documentation":"MongodbStreamFetchTask.normalizeChangeStreamDocument builds the source block of a change-event record. The MongoDB change stream event did not carry a clusterTime field, so the task logs this warning and substitutes the current BSON timestamp for source.ts_ms. The record is still emitted; only the event-time precision is degraded.","triggerScenarios":"A change stream document arriving from the MongoDB watch cursor lacks the clusterTime field — e.g. events produced by some proxy/sharded-router setups, resume-token-replayed events, or synthesized events that bypass the normal oplog entry wrapping.","commonSituations":"Reading change streams through mongos or a middleware that strips metadata; MongoDB versions/deployments where event structure differs; custom aggregation pipelines on $changeStream that drop clusterTime via $project/$replaceRoot stages.","solutions":["Check any $changeStream aggregation pipeline for $project/$replaceRoot stages that remove clusterTime and keep the field.","Verify events come directly from MongoDB (not a proxy/Debezium-style re-emit) that preserves clusterTime.","If the fallback timestamp is acceptable, treat this as informational; no action needed.","Upgrade the connector/MongoDB driver if events are consistently missing clusterTime on a standard replica set."],"exampleFix":"// before\npipeline = Arrays.asList(match(...), project(exclude(\"clusterTime\")))\n// after\npipeline = Arrays.asList(match(...)) // do not project away clusterTime","handlingStrategy":"fallback","validationCode":null,"typeGuard":null,"tryCatchPattern":null,"preventionTips":["Do not use $project/$replaceRoot in $changeStream pipelines in ways that drop clusterTime.","Consume change streams directly from the replica set/sharded cluster, not through proxies.","Monitor source.ts_ms drift to detect repeated fallback usage."],"tags":["mongodb","cdc","change-stream","timestamp"],"backgroundTag":"missing-required-config-field","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}