Skip to content
Open
Show file tree
Hide file tree
Changes from 11 commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
a0c7bec
feat(datastream-mongodb-to-firestore): implement high-throughput shad…
michaeltle-goog Aug 10, 2026
83ce0eb
fix(datastream-mongodb-to-firestore): fix timestamp sort key ordering…
michaeltle-goog Aug 11, 2026
8324c3f
style: apply spotless formatting
michaeltle-goog Aug 11, 2026
a4c6f78
fix(utils): add bool_ prefix to documentIdToString and add distinctne…
michaeltle-goog Aug 11, 2026
6983c6f
refactor: remove redundant read_method and _metadata_dlq_reconsumed c…
michaeltle-goog Aug 11, 2026
0137ad1
fix(review): address performance, thread-safety, defensive casting an…
michaeltle-goog Aug 11, 2026
440496b
perf(datastream-mongodb-to-firestore): eliminate CombineHotKeys windo…
michaeltle-goog Aug 11, 2026
cd1c08b
perf(datastream-mongodb-to-firestore): implement deterministic binary…
michaeltle-goog Aug 11, 2026
72bdcd7
perf(datastream-mongodb-to-firestore): decouple processing timestamps…
michaeltle-goog Aug 11, 2026
f4d5a8f
Revert "perf(datastream-mongodb-to-firestore): decouple processing ti…
michaeltle-goog Aug 11, 2026
ec9a813
refactor(datastream-mongodb-to-firestore): remove orderingStrategy pa…
michaeltle-goog Aug 11, 2026
9e4c0bf
Address review comments: 64-bit subSeconds coder, databaseName valida…
michaeltle-goog Aug 14, 2026
168bd30
Refactor ThrottledLogger, consolidate DatastreamConstants, standardiz…
michaeltle-goog Aug 14, 2026
add53a5
fix(datastream-mongodb-to-firestore): preserve customer document 'dat…
michaeltle-goog Aug 18, 2026
5df3db8
perf(datastream-mongodb-to-firestore): increase default write rate li…
michaeltle-goog Aug 20, 2026
6715eda
Revert "perf(datastream-mongodb-to-firestore): increase default write…
michaeltle-goog Aug 20, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,14 @@ public class DatastreamConstants {
// Event types
public static final String DELETE_EVENT = "DELETE";
public static final String UPDATE_EVENT = "UPDATE";
public static final String READ_EVENT = "READ";
public static final String EMPTY_EVENT = "";

// Read method metadata
public static final String EVENT_READ_METHOD_KEY = "_metadata_read_method";
public static final String READ_METHOD_BACKFILL = "backfill";
public static final String READ_METHOD_CDC = "cdc";

// Default shadow collection prefix
public static final String DEFAULT_SHADOW_COLLECTION_PREFIX = "shadow_";

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,13 @@
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.cloud.teleport.v2.transforms.MongoDbChangeEventContextCoder;
import com.google.cloud.teleport.v2.transforms.Utils;
import com.google.common.collect.ImmutableMap;
import java.io.IOException;
import java.io.Serializable;
import java.util.Objects;
import org.apache.beam.sdk.coders.DefaultCoder;
import org.bson.Document;
import org.bson.types.ObjectId;
import org.slf4j.Logger;
Expand All @@ -34,6 +37,7 @@
* MongoDB's implementation of ChangeEventContext that provides implementation for handling MongoDB
* change events.
*/
@DefaultCoder(MongoDbChangeEventContextCoder.class)
public class MongoDbChangeEventContext implements Serializable {

private static final Logger LOG = LoggerFactory.getLogger(MongoDbChangeEventContext.class);
Expand All @@ -51,6 +55,7 @@ public class MongoDbChangeEventContext implements Serializable {

private final JsonNode changeEvent;
private final JsonNode originalChangeEvent;
private final String shadowCollectionPrefix;
private final String dataCollection;
private final String shadowCollection;
private final Object documentId;
Expand All @@ -59,8 +64,8 @@ public class MongoDbChangeEventContext implements Serializable {
private final boolean isDeleteEvent;
private final boolean isUpdateEvent;
private final Document timestampDoc;
private boolean isDlqReconsumed;
private int retryCount;
private final boolean isDlqReconsumed;
private final int retryCount;

/** Gets the change type from the event metadata. */
private String getChangeType(JsonNode changeEvent) {
Expand All @@ -74,6 +79,46 @@ public String getChangeType() {
return getChangeType(this.changeEvent);
}

/** Determines if the event is a backfill snapshot event. */
public boolean isBackfillEvent() {
if (changeEvent.has(DatastreamConstants.EVENT_READ_METHOD_KEY)) {
String readMethod = changeEvent.get(DatastreamConstants.EVENT_READ_METHOD_KEY).asText();
if (DatastreamConstants.READ_METHOD_BACKFILL.equalsIgnoreCase(readMethod)) {
return true;
}
}
String changeType = getChangeType();
return DatastreamConstants.READ_EVENT.equalsIgnoreCase(changeType)
|| "BACKFILL".equalsIgnoreCase(changeType);
}

/** Determines if the event is a live CDC event. */
public boolean isCdcEvent() {
return !isBackfillEvent();
}

/** Gets epoch timestamp seconds. */
public long getTimestampSeconds() {
if (timestampDoc != null && timestampDoc.containsKey(TIMESTAMP_SECONDS_COL)) {
Object val = timestampDoc.get(TIMESTAMP_SECONDS_COL);
if (val instanceof Number) {
return ((Number) val).longValue();
}
}
return 0L;
}

/** Gets sub-second timestamp (wall nanoseconds for backfill, oplog increment for CDC). */
public long getTimestampSubSeconds() {
if (timestampDoc != null && timestampDoc.containsKey(TIMESTAMP_NANOS_COL)) {
Object val = timestampDoc.get(TIMESTAMP_NANOS_COL);
if (val instanceof Number) {
return ((Number) val).longValue();
}
}
return 0L;
}

/** Determines if the event is a delete event based on metadata. */
private boolean isDeleteEvent(JsonNode changeEvent) {
String changeType = getChangeType(changeEvent);
Expand Down Expand Up @@ -105,17 +150,28 @@ public MongoDbChangeEventContext(JsonNode payload, String shadowCollectionPrefix
public MongoDbChangeEventContext(
JsonNode payload, JsonNode originalPayload, String shadowCollectionPrefix)
throws JsonProcessingException {
this(
payload,
originalPayload,
shadowCollectionPrefix,
extractIsDlqReconsumed(payload),
extractRetryCount(payload));
}

private MongoDbChangeEventContext(
JsonNode payload,
JsonNode originalPayload,
String shadowCollectionPrefix,
boolean isDlqReconsumed,
int retryCount)
throws JsonProcessingException {
// Extracts the actual change event. If wrapped in a DLQ structure like {"changeEvent": {...}},
// it extracts the inner object.
this.changeEvent = Utils.extractInnerEvent(payload);
this.originalChangeEvent = Utils.extractInnerEvent(originalPayload);

this.retryCount =
changeEvent.has(DatastreamConstants.RETRY_COUNT)
? changeEvent.get(DatastreamConstants.RETRY_COUNT).asInt()
: payload.has(DatastreamConstants.RETRY_COUNT)
? payload.get(DatastreamConstants.RETRY_COUNT).asInt()
: 0;
this.shadowCollectionPrefix = shadowCollectionPrefix != null ? shadowCollectionPrefix : "";
this.isDlqReconsumed = isDlqReconsumed;
this.retryCount = retryCount;

// Extract collection name from the event
if (changeEvent.has(DatastreamConstants.EVENT_SOURCE_METADATA)) {
Expand All @@ -129,7 +185,7 @@ public MongoDbChangeEventContext(
throw new IllegalStateException("Invalid event record without _metadata_source.");
}

this.shadowCollection = shadowCollectionPrefix + this.dataCollection;
this.shadowCollection = this.shadowCollectionPrefix + this.dataCollection;

// Extract document id
if (changeEvent.has(DatastreamConstants.MONGODB_DOCUMENT_ID)) {
Expand All @@ -143,11 +199,13 @@ public MongoDbChangeEventContext(
this.documentId = docIdVal.asDouble();
} else if (docIdVal.isTextual()) {
this.documentId = docIdVal.asText();
} else if (docIdVal.isObject()) {
if (docIdVal.has(OID_FIELD_NAME) && docIdVal.get(OID_FIELD_NAME).isTextual()) {
} else if (docIdVal.isObject() || docIdVal.isArray()) {
if (docIdVal.isObject()
&& docIdVal.has(OID_FIELD_NAME)
&& docIdVal.get(OID_FIELD_NAME).isTextual()) {
this.documentId = new ObjectId(docIdVal.get(OID_FIELD_NAME).asText());
} else {
// Support for generic Object-typed IDs or other complex BSON types (e.g., Binary)
// Support for generic Document (Map), Array (List), Binary, and composite BSON _id types
Document wrapper = Document.parse("{ \"val\": " + docIdVal.toString() + " }");
this.documentId = wrapper.get("val");
}
Expand Down Expand Up @@ -182,8 +240,61 @@ public MongoDbChangeEventContext(
this.isUpdateEvent = isUpdateEvent(changeEvent);

this.jsonStringData = dataAsJsonString();
this.shadowDocument = generateShadowDocument();
this.isDlqReconsumed = isDlqReconsumed(changeEvent);
this.shadowDocument = null;
}

private static boolean extractIsDlqReconsumed(JsonNode payload) {
if (payload == null) {
return false;
}
JsonNode changeEvent = Utils.extractInnerEvent(payload);
if (changeEvent.has(DatastreamConstants.IS_DLQ_RECONSUMED)) {
return changeEvent
.get(DatastreamConstants.IS_DLQ_RECONSUMED)
.asText()
.equalsIgnoreCase("true");
}
if (payload.has(DatastreamConstants.IS_DLQ_RECONSUMED)) {
return payload.get(DatastreamConstants.IS_DLQ_RECONSUMED).asText().equalsIgnoreCase("true");
}
return false;
}

private static int extractRetryCount(JsonNode payload) {
if (payload == null) {
return 0;
}
JsonNode changeEvent = Utils.extractInnerEvent(payload);
if (changeEvent.has(DatastreamConstants.RETRY_COUNT)) {
return changeEvent.get(DatastreamConstants.RETRY_COUNT).asInt();
}
if (payload.has(DatastreamConstants.RETRY_COUNT)) {
return payload.get(DatastreamConstants.RETRY_COUNT).asInt();
}
return 0;
}

/**
* Reconstitutes a {@link MongoDbChangeEventContext} from its serialized components without Java
* reflection serialization.
*/
public static MongoDbChangeEventContext reconstitute(
String changeEventJson,
String originalChangeEventJson,
String shadowPrefix,
boolean isDlq,
int retryCount)
throws IOException {
if (changeEventJson == null) {
return null;
}
JsonNode changeEventNode = OBJECT_MAPPER.readTree(changeEventJson);
JsonNode originalChangeEventNode =
originalChangeEventJson != null
? OBJECT_MAPPER.readTree(originalChangeEventJson)
: changeEventNode;
return new MongoDbChangeEventContext(
changeEventNode, originalChangeEventNode, shadowPrefix, isDlq, retryCount);
}

/** Creates a shadow document for tracking event ordering. */
Expand Down Expand Up @@ -224,6 +335,18 @@ public JsonNode getOriginalChangeEvent() {
return originalChangeEvent;
}

public String getChangeEventJsonString() {
return changeEvent != null ? changeEvent.toString() : null;
}

public String getOriginalChangeEventJsonString() {
return originalChangeEvent != null ? originalChangeEvent.toString() : null;
}

public String getShadowCollectionPrefix() {
return shadowCollectionPrefix;
}

public String getDataCollection() {
return dataCollection;
}
Expand All @@ -245,6 +368,13 @@ public boolean isUpdateEvent() {
}

public Document getShadowDocument() {
if (this.shadowDocument == null && this.shadowCollection != null) {
try {
return generateShadowDocument();
} catch (JsonProcessingException e) {
LOG.warn("Failed to generate shadow document: {}", e.getMessage());
}
}
return shadowDocument;
}

Expand Down Expand Up @@ -291,8 +421,8 @@ public String toString() {
// Convert timestamp document to JSON
if (this.timestampDoc != null) {
ObjectNode timestampNode = OBJECT_MAPPER.createObjectNode();
timestampNode.put(TIMESTAMP_SECONDS_COL, this.timestampDoc.getLong(TIMESTAMP_SECONDS_COL));
timestampNode.put(TIMESTAMP_NANOS_COL, this.timestampDoc.getInteger(TIMESTAMP_NANOS_COL));
timestampNode.put(TIMESTAMP_SECONDS_COL, getTimestampSeconds());
timestampNode.put(TIMESTAMP_NANOS_COL, getTimestampSubSeconds());
jsonNode.set(TIMESTAMP_COL, timestampNode);
}

Expand All @@ -317,10 +447,39 @@ public String toString() {
}
}

@Override
public boolean equals(Object other) {
if (this == other) {
return true;
}
if (other instanceof MongoDbChangeEventContext) {
return Objects.equals(this.toString(), other.toString());
MongoDbChangeEventContext o = (MongoDbChangeEventContext) other;
return Objects.equals(this.dataCollection, o.dataCollection)
&& Objects.equals(this.shadowCollectionPrefix, o.shadowCollectionPrefix)
&& Objects.equals(this.documentId, o.documentId)
&& Objects.equals(this.timestampDoc, o.timestampDoc)
&& this.isDeleteEvent == o.isDeleteEvent
&& this.isUpdateEvent == o.isUpdateEvent
&& this.isDlqReconsumed == o.isDlqReconsumed
&& this.retryCount == o.retryCount
&& Objects.equals(this.changeEvent, o.changeEvent)
&& Objects.equals(this.originalChangeEvent, o.originalChangeEvent);
}
return false;
}

@Override
public int hashCode() {
return Objects.hash(
dataCollection,
shadowCollectionPrefix,
documentId,
timestampDoc,
isDeleteEvent,
isUpdateEvent,
isDlqReconsumed,
retryCount,
changeEvent,
originalChangeEvent);
}
}
Loading
Loading