Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
5 changes: 5 additions & 0 deletions CHANGLOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,10 @@
# Changelog

## 1.9.0

- Add `batchScanSize` option to scan the collection by `_id` ranges (avoid socket timeouts) - https://github.com/360Learning/mongo-bulk-data-migration/issues/39
- Add `hint` option to force the index used by the migration query

## 1.8.2 (2026-07-24)

- Fix bug: rollback progress correctly logged
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,8 @@ new MongoBulkDataMigration({ ..., options: { ... } })
- `bypassRollbackValidation` and `bypassUpdateValidation` (default: false): Will set validationLevel to "off" then to "moderate".
- `throttle` (default: 0): amount of time in ms to sleep between a bulk update. Use this to decrease database stress.
- `continueOnBulkWriteError` (default: false): will continue the migration on the error in a bulk.
- `batchScanSize` (default: none): scan the collection by consecutive `_id` ranges of `batchScanSize` documents (e.g. `100_000`) instead of a single query. Can be useful for long migration, can avoid socket timeout.
- `hint` (default: none): index to force for the migration query and its count

## 📕 Advanced usages

Expand Down
178 changes: 177 additions & 1 deletion __tests__/MongoBulkDataMigration.update.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,13 @@
import _ from 'lodash';
// import { ObjectId } from 'bson';
import { Collection, Db, Document, ObjectId, UpdateFilter } from 'mongodb';
import {
Collection,
Db,
Document,
ObjectId,
UpdateFilter,
WithId,
} from 'mongodb';
import { MongoBulkDataMigration, DELETE_OPERATION, FETCH_ALL } from '../src';
import { INITIAL_BULK_INFOS } from '../src/lib/AbstractBulkOperationResults';
import { LoggerInterface } from '../src/types';
Expand Down Expand Up @@ -522,6 +529,162 @@ describe('MongoBulkDataMigration', () => {
});
});

describe('options.batchScanSize', () => {
let ids: ObjectId[];
beforeEach(async () => {
const insertResult = await collection.insertMany(
Array.from({ length: 10 }, (_, i) => ({ value: i + 1 })),
);
ids = Object.values(insertResult.insertedIds);
});

it('should migrate all documents scanning the collection by _id ranges', async () => {
await db.command({ profile: 2 });
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
options: { batchScanSize: 3 },
update: { $set: { migrated: true } },
});

const updateResults = await dataMigration.update();

expect(updateResults).toEqual({
...INITIAL_BULK_INFOS,
nMatched: 10,
nModified: 10,
});
expect(await collection.countDocuments({ migrated: true })).toEqual(10);
const rangeQueries = await getProfiledFindQueries({
'command.hint': { _id: 1 },
});
expect(rangeQueries.map(({ docsExamined }) => docsExamined)).toEqual([
3, 3, 3, 1,
]);
});

it('should combine ranges with a query already filtering on _id', async () => {
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
query: { _id: { $in: [ids[0], ids[4], ids[9]] } },
options: { batchScanSize: 3 },
update: { $set: { migrated: true } },
});

await dataMigration.update();

const migratedIds = await collection
.find({ migrated: true })
.map(({ _id }) => _id)
.toArray();
expect(migratedIds).toEqual([ids[0], ids[4], ids[9]]);
});

it('should migrate documents matched by an aggregate pipeline', async () => {
const incUpdateStub = jest
.fn()
.mockReturnValue({ $set: { migrated: true } });
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
query: [{ $match: { value: { $mod: [2, 0] } } }],
projection: { value: 1 },
options: { batchScanSize: 3 },
update: incUpdateStub,
});

await dataMigration.update();

expect(incUpdateStub.mock.calls.map(([doc]) => doc)).toEqual(
[2, 4, 6, 8, 10].map((value) => ({ _id: ids[value - 1], value })),
);
});
});

describe('options.hint', () => {
beforeEach(async () => {
await collection.insertMany(
Array.from({ length: 10 }, (_, i) => ({ value: i + 1 })),
);
await collection.createIndex({ value: 1 }, { name: 'value_1' });
});

afterEach(async () => {
await collection.dropIndexes();
});

it('should use the hinted index for the migration query', async () => {
await db.command({ profile: 2 });
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
query: { value: { $gt: 5 } },
options: { hint: { _id: 1 } },
update: { $set: { migrated: true } },
});

await dataMigration.update();

const [migrationQuery, ...otherQueries] = await getProfiledFindQueries({
'command.filter': { value: { $gt: 5 } },
});
expect(otherQueries).toEqual([]);
expect(migrationQuery.command.hint).toEqual({ _id: 1 });
expect(migrationQuery.planSummary).toEqual('IXSCAN { _id: 1 }');
expect(await collection.countDocuments({ migrated: true })).toEqual(5);
});

it('should use the hinted index for an aggregate pipeline', async () => {
await db.command({ profile: 2 });
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
query: [{ $match: { value: { $gt: 5 } } }],
options: { hint: { value: 1 }, dontCount: true },
update: { $set: { migrated: true } },
});

await dataMigration.update();

const [aggregateQuery] = await getProfiledFindQueries({
op: 'command',
'command.aggregate': COLLECTION,
});
expect(aggregateQuery.command.hint).toEqual({ value: 1 });
expect(aggregateQuery.planSummary).toEqual('IXSCAN { value: 1 }');
expect(await collection.countDocuments({ migrated: true })).toEqual(5);
});

it('should override the default _id hint of batchScanSize range queries', async () => {
await db.command({ profile: 2 });
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
query: { value: { $gt: 5 } },
options: { hint: { value: 1 }, batchScanSize: 5, dontCount: true },
update: { $set: { migrated: true } },
});

await dataMigration.update();

const rangeQueries = await getProfiledFindQueries({
'command.hint': { value: 1 },
});
expect(rangeQueries.map(({ planSummary }) => planSummary)).toEqual([
'IXSCAN { value: 1 }',
'IXSCAN { value: 1 }',
]);
expect(await collection.countDocuments({ migrated: true })).toEqual(5);
});

it('should reject when the hinted index does not exist', async () => {
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
options: { hint: 'unknown_index', dontCount: true },
update: { $set: { migrated: true } },
});

await expect(dataMigration.update()).rejects.toThrow(
'hint provided does not correspond to an existing index',
);
});
});

describe('NO_UPDATE update action', () => {
it('should not update documents when the update function returns NO_UPDATE', async () => {
await collection.insertMany([{ key: 1 }, { key: 2 }, { key: 3 }]);
Expand Down Expand Up @@ -756,6 +919,19 @@ describe('MongoBulkDataMigration', () => {
});
});

async function getProfiledFindQueries(
filter: Document = {},
): Promise<WithId<Document>[]> {
await db.command({ profile: 0 });
const profileCollection = db.collection('system.profile');
const queries = await profileCollection
.find({ ns: `${db.databaseName}.${COLLECTION}`, op: 'query', ...filter })
.sort({ ts: 1 })
.toArray();
await profileCollection.drop();
return queries;
}

function extractLogsPayload(logText: string) {
return loggerMock.info.mock.calls
.filter(([_, msg]) => msg === logText)
Expand Down
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@360-l/mongo-bulk-data-migration",
"version": "1.8.2",
"version": "1.9.0",
"description": "MongoDB bulk data migration for node scripts",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
Expand Down
97 changes: 86 additions & 11 deletions src/MongoBulkDataMigration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,7 @@ export default class MongoBulkDataMigration<
}

await this.lowerValidationLevel('update');
const { cursor, totalEntries } = await this.getCursorAndCount(
const { documents, totalEntries } = await this.getDocumentsAndCount(
migrationCollection,
rollbackCollection,
);
Expand All @@ -156,14 +156,14 @@ export default class MongoBulkDataMigration<
let updatePromises: Promise<any>[] = [];

let treatedDocumentsCount = 0;
let document = (await cursor.next()) as WithId<TSchema> | null;
let document = await nextDocument();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[nit] Wouldn't it be better to have a clearer interface for this function?

Suggested change
let document = await nextDocument();
let document = await nextDocument(documents);

while (document !== null) {
const bulkUpdateWrappedPromise = updatePromiseLimiter(
this.buildBulkUpdater(document, bulkBackup, bulkMigration),
);
updatePromises.push(bulkUpdateWrappedPromise);

document = (await cursor.next()) as WithId<TSchema> | null;
document = await nextDocument();
if (!document || updatePromises.length >= this.options.maxBulkSize) {
await Promise.all(updatePromises);
const backupRes = (await bulkBackup.execute()).getResults();
Expand Down Expand Up @@ -205,29 +205,45 @@ export default class MongoBulkDataMigration<
);
await this.restoreValidationLevel('update');
return bulkMigration.getResults();

async function nextDocument() {
const { value, done } = await documents.next();
return done ? null : value;
}
}

private async getCursorAndCount(
private async getDocumentsAndCount(
migrationCollection: Collection<TSchema>,
rollbackCollection: Collection<TSchema>,
) {
const resolvedQuery = await this.resolveQuery(rollbackCollection);

const cursor = getCursor(resolvedQuery, this.migrationInfos);
const documents: AsyncIterable<WithId<TSchema>> = this.options.batchScanSize
? this.iterateByIdRanges(
migrationCollection,
resolvedQuery,
this.options.batchScanSize,
)
: (getCursor(
resolvedQuery,
this.migrationInfos,
this.options.hint,
) as AsyncIterable<WithId<TSchema>>);
const totalEntries = await getTotalEntries(resolvedQuery, this);
return { cursor, totalEntries };
return { documents: documents[Symbol.asyncIterator](), totalEntries };

function getCursor(
query: Filter<TSchema> | MongoPipeline,
{ projection }: MigrationInfos<TSchema>,
hint: DataMigrationOptions<TSchema>['hint'],
) {
if (isPipeline(query)) {
const pipelineWithProjection = query.concat(
_.isEmpty(projection) ? [] : [{ $project: projection }],
);
return migrationCollection.aggregate(pipelineWithProjection);
return migrationCollection.aggregate(pipelineWithProjection, { hint });
}
return migrationCollection.find(query, { projection });
return migrationCollection.find(query, { projection, hint });
}

async function getTotalEntries(
Expand All @@ -250,14 +266,18 @@ export default class MongoBulkDataMigration<
try {
if (isPipeline(query)) {
const pipelineComputeTotal = query.concat({ $count: 'totalEntries' });
const cursorComputeTotal =
migrationCollection.aggregate(pipelineComputeTotal);
const cursorComputeTotal = migrationCollection.aggregate(
pipelineComputeTotal,
{ hint: that.options.hint },
);
const total = (await cursorComputeTotal.next()) as unknown as {
totalEntries: number;
} | null;
return total === null ? 0 : total.totalEntries;
}
return migrationCollection.countDocuments(query);
return migrationCollection.countDocuments(query, {
hint: that.options.hint,
});
} finally {
clearTimeout(countTakingTooLongTimeout);
}
Expand All @@ -270,6 +290,61 @@ export default class MongoBulkDataMigration<
}
}

private async *iterateByIdRanges(
migrationCollection: Collection<TSchema>,
query: Filter<TSchema> | MongoPipeline,
batchScanSize: number,
): AsyncGenerator<WithId<TSchema>> {
const { projection } = this.migrationInfos;
let lowerBound: unknown = undefined;

do {
const upperBoundDocument = await migrationCollection
.find(
(lowerBound === undefined
? {}
: { _id: { $gte: lowerBound } }) as Filter<TSchema>,
{ projection: { _id: 1 } },
)
.sort({ _id: 1 })
.skip(batchScanSize)
.limit(1)
.next();
const upperBound = upperBoundDocument?._id;
const idRange = {
...(lowerBound !== undefined && { $gte: lowerBound }),
...(upperBound !== undefined && { $lt: upperBound }),
};

if (Array.isArray(query)) {
const rangeMatch: MongoPipeline = _.isEmpty(idRange)
? []
: [{ $match: { _id: idRange } }];
const pipeline = rangeMatch
.concat(query)
.concat(_.isEmpty(projection) ? [] : [{ $project: projection }]);
yield* migrationCollection.aggregate<WithId<TSchema>>(pipeline, {
hint: this.options.hint ?? { _id: 1 },
});
} else {
const rangeQuery = _.isEmpty(idRange)
? query
: '_id' in query
? { $and: [query, { _id: idRange }] }
: { ...query, _id: idRange };
yield* migrationCollection.find(rangeQuery as Filter<TSchema>, {
projection,
hint: this.options.hint ?? { _id: 1 },
});
}

lowerBound = upperBound;
if (lowerBound !== undefined) {
await this.throttle();
}
} while (lowerBound !== undefined);
}

private buildBulkUpdater(
document: WithId<TSchema>,
bulkBackup: BackupBulk<TSchema>,
Expand Down
Loading
Loading