SCRAPE-6: Implement upsertUnits() bulk operation (#11)
Co-authored-by: Stephen Minakian <stephenminakian@gmail.com> Co-committed-by: Stephen Minakian <stephenminakian@gmail.com>
This commit is contained in:
@ -369,10 +369,72 @@ function convertDataTypes(unit) {
|
||||
};
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Database Operations
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Upsert unit records to database using bulkWrite.
|
||||
* Each unit is matched by unit_code as the unique key.
|
||||
* Sets last_scraped and data_source on every update.
|
||||
* Sets first_seen only on initial insert via $setOnInsert.
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Array<Object>} units - Array of unit objects
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<Object>} Bulk write result
|
||||
*/
|
||||
async function upsertUnits(db, units, logger) {
|
||||
const collection = db.collection(config.COLLECTIONS.UNITS);
|
||||
const now = new Date().toISOString();
|
||||
|
||||
const operations = units.map(unit => ({
|
||||
updateOne: {
|
||||
filter: { unit_code: unit.unit_code },
|
||||
update: {
|
||||
$set: {
|
||||
...unit,
|
||||
last_scraped: now,
|
||||
data_source: 'web_scraper'
|
||||
},
|
||||
$setOnInsert: {
|
||||
first_seen: now
|
||||
}
|
||||
},
|
||||
upsert: true
|
||||
}
|
||||
}));
|
||||
|
||||
if (operations.length === 0) {
|
||||
logger.warn('No units to upsert');
|
||||
return { modifiedCount: 0, upsertedCount: 0 };
|
||||
}
|
||||
|
||||
try {
|
||||
const result = await collection.bulkWrite(operations, { ordered: false });
|
||||
|
||||
logger.info('Units upserted', {
|
||||
matched: result.matchedCount,
|
||||
modified: result.modifiedCount,
|
||||
upserted: result.upsertedCount
|
||||
});
|
||||
|
||||
return result;
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Failed to upsert units', {
|
||||
errorType: error.name,
|
||||
errorMessage: error.message
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
fetchPage,
|
||||
parseUnits,
|
||||
convertDataTypes,
|
||||
upsertUnits,
|
||||
// Export helpers for testing
|
||||
parseInteger,
|
||||
parsePositiveInteger,
|
||||
|
||||
Reference in New Issue
Block a user