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

  1. Check any $changeStream aggregation pipeline for $project/$replaceRoot stages that remove clusterTime and keep the field.
  2. Verify events come directly from MongoDB (not a proxy/Debezium-style re-emit) that preserves clusterTime.
  3. If the fallback timestamp is acceptable, treat this as informational; no action needed.
  4. 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

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


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)