View a markdown version of this page

$changeStreamSplitLargeEvent - Amazon DocumentDB

$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()