$changeStreamSplitLargeEvent
This aggregation stage is not supported by Elastic clusters.
This stage is available starting with Amazon DocumentDB engine version 8.0.2.
The $changeStreamSplitLargeEvent aggregation stage is used within a $changeStream pipeline to split events exceeding 16 MB into smaller fragments. It returns each fragment sequentially via the change stream cursor.
Parameters
None. The input to the $changeStreamSplitLargeEvent stage should be an empty document.
Example – MongoDB Shell
The following example demonstrates using the $changeStreamSplitLargeEvent stage to split an oversized change event into multiple fragments.
Query example
// Insert a document
db.inventory.insertOne({ _id: 1, item: "Widget", payload: "x".repeat(9 * 1024 * 1024) })
// Open change stream with the $changeStreamSplitLargeEvent stage
var changeStream = db.inventory.aggregate([
{ $changeStream: { fullDocument: "updateLookup" } },
{ $changeStreamSplitLargeEvent: {} }
]);
// Update the document
db.inventory.updateOne({ _id: 1 }, { $set: { payload: "y".repeat(9 * 1024 * 1024) } })
// Read the change event
if (changeStream.hasNext()) {
print(tojson(changeStream.next()));
}
Output
{
_id: { _data: '...' },
splitEvent: { fragment: 1, of: 2},
ns: { db: 'test', coll: 'inventory' },
clusterTime: Timestamp(4, 1789154779),
documentKey: { _id: 1 },
fullDocument: { _id: 1, item: 'Widget', payload: yyyyyyyy...(9437184 chars) },
operationType: 'update',
}
{
_id: { _data: '...' },
splitEvent: { fragment: 2, of: 2},
updateDescription: { updatedFields: { payload: yyyyyyyy...(9437184 chars) }, removedFields: [] }
}
Code examples
To view a code example for using the $changeStreamSplitLargeEvent aggregation stage, choose the tab for the language that you want to use:
- Node.js
-
const { MongoClient } = require('mongodb');
async function example() {
const client = await MongoClient.connect('mongodb://<username>:<password>@<cluster-endpoint>:27017/?tls=true&tlsCAFile=global-bundle.pem&replicaSet=rs0&readPreference=secondaryPreferred&retryWrites=false');
const db = client.db('test');
const collection = db.collection('inventory');
// Insert a large document first
await collection.insertOne({ _id: 1, item: 'Widget', payload: 'x'.repeat(9 * 1024 * 1024) });
// Open change stream with $changeStreamSplitLargeEvent
const changeStream = collection.watch(
[{ $changeStreamSplitLargeEvent: {} }],
{ fullDocument: 'updateLookup' }
);
changeStream.on('change', (change) => {
console.log('Change detected:', change);
});
// Update to trigger a split event — fullDocument + updateDescription > 16 MB
setTimeout(async () => {
console.log('Triggering update...');
await collection.updateOne({ _id: 1 }, { $set: { payload: 'y'.repeat(9 * 1024 * 1024) } });
}, 1000);
// Keep connection open to receive changes
// In production, handle cleanup appropriately
}
example();
- Python
-
from pymongo import MongoClient
import threading
import time
def example():
client = MongoClient('mongodb://<username>:<password>@<cluster-endpoint>:27017/?tls=true&tlsCAFile=global-bundle.pem&replicaSet=rs0&readPreference=secondaryPreferred&retryWrites=false')
db = client['test']
collection = db['inventory']
# Insert large document first
collection.drop()
collection.insert_one({'_id': 1, 'item': 'Widget', 'payload': 'x' * (9 * 1024 * 1024)})
# Open change stream with $changeStreamSplitLargeEvent
change_stream = collection.watch(
[{'$changeStreamSplitLargeEvent': {}}],
full_document='updateLookup'
)
# Update in separate thread after delay
def update_doc():
collection.update_one({'_id': 1}, {'$set': {'payload': 'y' * (9 * 1024 * 1024)}})
threading.Thread(target=update_doc).start()
# Watch for changes
for change in change_stream:
print('Change detected:', change)
client.close()
example()