diff --git a/CHANGLOG.md b/CHANGLOG.md index 092cf93..32af98e 100644 --- a/CHANGLOG.md +++ b/CHANGLOG.md @@ -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 diff --git a/README.md b/README.md index cef94b9..277eae0 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/__tests__/MongoBulkDataMigration.update.test.ts b/__tests__/MongoBulkDataMigration.update.test.ts index 92e4a86..62a3eab 100644 --- a/__tests__/MongoBulkDataMigration.update.test.ts +++ b/__tests__/MongoBulkDataMigration.update.test.ts @@ -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'; @@ -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 }]); @@ -756,6 +919,19 @@ describe('MongoBulkDataMigration', () => { }); }); + async function getProfiledFindQueries( + filter: Document = {}, + ): Promise[]> { + 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) diff --git a/package.json b/package.json index ba7faf9..53d5d50 100644 --- a/package.json +++ b/package.json @@ -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", diff --git a/src/MongoBulkDataMigration.ts b/src/MongoBulkDataMigration.ts index 9542dad..080f80f 100644 --- a/src/MongoBulkDataMigration.ts +++ b/src/MongoBulkDataMigration.ts @@ -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, ); @@ -156,14 +156,14 @@ export default class MongoBulkDataMigration< let updatePromises: Promise[] = []; let treatedDocumentsCount = 0; - let document = (await cursor.next()) as WithId | null; + let document = await nextDocument(); while (document !== null) { const bulkUpdateWrappedPromise = updatePromiseLimiter( this.buildBulkUpdater(document, bulkBackup, bulkMigration), ); updatePromises.push(bulkUpdateWrappedPromise); - document = (await cursor.next()) as WithId | null; + document = await nextDocument(); if (!document || updatePromises.length >= this.options.maxBulkSize) { await Promise.all(updatePromises); const backupRes = (await bulkBackup.execute()).getResults(); @@ -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, rollbackCollection: Collection, ) { const resolvedQuery = await this.resolveQuery(rollbackCollection); - const cursor = getCursor(resolvedQuery, this.migrationInfos); + const documents: AsyncIterable> = this.options.batchScanSize + ? this.iterateByIdRanges( + migrationCollection, + resolvedQuery, + this.options.batchScanSize, + ) + : (getCursor( + resolvedQuery, + this.migrationInfos, + this.options.hint, + ) as AsyncIterable>); const totalEntries = await getTotalEntries(resolvedQuery, this); - return { cursor, totalEntries }; + return { documents: documents[Symbol.asyncIterator](), totalEntries }; function getCursor( query: Filter | MongoPipeline, { projection }: MigrationInfos, + hint: DataMigrationOptions['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( @@ -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); } @@ -270,6 +290,61 @@ export default class MongoBulkDataMigration< } } + private async *iterateByIdRanges( + migrationCollection: Collection, + query: Filter | MongoPipeline, + batchScanSize: number, + ): AsyncGenerator> { + const { projection } = this.migrationInfos; + let lowerBound: unknown = undefined; + + do { + const upperBoundDocument = await migrationCollection + .find( + (lowerBound === undefined + ? {} + : { _id: { $gte: lowerBound } }) as Filter, + { 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>(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, { + projection, + hint: this.options.hint ?? { _id: 1 }, + }); + } + + lowerBound = upperBound; + if (lowerBound !== undefined) { + await this.throttle(); + } + } while (lowerBound !== undefined); + } + private buildBulkUpdater( document: WithId, bulkBackup: BackupBulk, diff --git a/src/types.ts b/src/types.ts index 0ef61c9..700e878 100644 --- a/src/types.ts +++ b/src/types.ts @@ -5,6 +5,7 @@ import type { UpdateFilter, ObjectId, Document, + Hint, } from 'mongodb'; import type { DELETE_OPERATION } from './lib/MigrationBulk'; import { DELETE_COLLECTION, FETCH_ALL } from './MongoBulkDataMigration'; @@ -13,10 +14,14 @@ import type { NO_UPDATE } from './MongoBulkDataMigration'; export type DataMigrationOptions = { /** Array filters to use in case of a migration on nested object in arrays */ arrayFilters: Document[]; + /** Limit documents to search and exec a new find() at every batch */ + batchScanSize?: number; /** Disable document validation temporarily on the rollback process */ bypassRollbackValidation: boolean; /** Disable document validation temporarily on the update process */ bypassUpdateValidation: boolean; + /** Index to force for the migration query */ + hint?: Hint; /** When counting drops performance before the migration _(un-indexed results or aggregation)_, turn this on */ dontCount: boolean; /** When set to true, for an update, MongoBulkWriteError (only) won't stop the update operation and be accumulated in the return response */