apache/seatunnel · warning
Cannot extract clusterTime from change stream event…
Error message
Cannot extract clusterTime from change stream event, fallback to current timestamp.
What it means
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.
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.
Example fix
// before
pipeline = Arrays.asList(match(...), project(exclude("clusterTime")))
// after
pipeline = Arrays.asList(match(...)) // do not project away clusterTime Defensive patterns
Strategy: fallback
Prevention
- 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.
When it happens
Trigger: 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.
Common situations: 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.
Understand the failure class
Background: "is required", "must be set", "missing required field": configuration validation errors across open-source libraries — this error's family across 36 libraries.
Related errors
- Change stream cursor has expired, trying to recreate cursor
- ILLEGAL_ARGUMENT
- ILLEGAL_ARGUMENT
- Non-heartbeat record has no documentKey field, this is…
- Resume token has expired, fallback to timestamp restart mode
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/7e6f513cd70aa9ba.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/fetch/MongodbStreamFetchTask.java:365
}
return null;
}
@Nonnull
private BsonDocument normalizeChangeStreamDocument(@Nonnull BsonDocument changeStreamDocument) {
// _id: primary key of change document.
BsonDocument normalizedDocument = normalizeKeyDocument(changeStreamDocument);
changeStreamDocument.put(ID_FIELD, normalizedDocument);
// ts_ms: It indicates the time at which the reader processed the event.
changeStreamDocument.put(TS_MS_FIELD, new BsonInt64(System.currentTimeMillis()));
// source
BsonDocument source = new BsonDocument();
source.put(SNAPSHOT_FIELD, new BsonString(FALSE_FALSE));
if (!changeStreamDocument.containsKey(CLUSTER_TIME_FIELD)) {
log.warn(
"Cannot extract clusterTime from change stream event, fallback to current timestamp.");
changeStreamDocument.put(CLUSTER_TIME_FIELD, currentBsonTimestamp());
}
// source.ts_ms
// It indicates the time that the change was made in the database. If the record is read
// from snapshot of the table instead of the change stream, the value is always 0.
BsonTimestamp clusterTime = changeStreamDocument.getTimestamp(CLUSTER_TIME_FIELD);
Instant clusterInstant = Instant.ofEpochSecond(clusterTime.getTime());
source.put(TS_MS_FIELD, new BsonInt64(clusterInstant.toEpochMilli()));
changeStreamDocument.put(SOURCE_FIELD, source);
return changeStreamDocument;
}
@Nonnull
private BsonDocument normalizeKeyDocument(@Nonnull BsonDocument changeStreamDocument) {
BsonDocument documentKey = changeStreamDocument.getDocument(DOCUMENT_KEY);View on GitHub (pinned to cf67b549a7)