SCRAPE-7: Implement insertPrices() bulk operation (#12)
Co-authored-by: Stephen Minakian <stephenminakian@gmail.com> Co-committed-by: Stephen Minakian <stephenminakian@gmail.com>
This commit is contained in:
@ -430,11 +430,75 @@ async function upsertUnits(db, units, logger) {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Insert price records for today using bulkWrite with upsert.
|
||||
* Uses unit_code + date_checked as the composite key to prevent
|
||||
* duplicate price records for the same unit on the same day.
|
||||
* Re-running on the same day updates existing records (idempotent).
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Array<Object>} units - Array of unit objects
|
||||
* @param {string} date - Date in YYYY-MM-DD format
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<Object>} Insert result with insertedCount
|
||||
*/
|
||||
async function insertPrices(db, units, date, logger) {
|
||||
const collection = db.collection(config.COLLECTIONS.PRICES);
|
||||
const now = new Date().toISOString();
|
||||
|
||||
// Filter out units without prices
|
||||
const unitsWithPrices = units.filter(u => u.price !== null);
|
||||
|
||||
const priceRecords = unitsWithPrices.map(unit => ({
|
||||
unit_code: unit.unit_code,
|
||||
date_checked: date,
|
||||
price: unit.price,
|
||||
last_updated: now,
|
||||
data_source: 'web_scraper'
|
||||
}));
|
||||
|
||||
if (priceRecords.length === 0) {
|
||||
logger.warn('No price records to insert');
|
||||
return { insertedCount: 0 };
|
||||
}
|
||||
|
||||
// Use updateOne with upsert to handle re-runs on same day
|
||||
const operations = priceRecords.map(record => ({
|
||||
updateOne: {
|
||||
filter: {
|
||||
unit_code: record.unit_code,
|
||||
date_checked: record.date_checked
|
||||
},
|
||||
update: { $set: record },
|
||||
upsert: true
|
||||
}
|
||||
}));
|
||||
|
||||
try {
|
||||
const result = await collection.bulkWrite(operations, { ordered: false });
|
||||
|
||||
logger.info('Prices inserted', {
|
||||
inserted: result.upsertedCount,
|
||||
updated: result.modifiedCount
|
||||
});
|
||||
|
||||
return { insertedCount: result.upsertedCount + result.modifiedCount };
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Failed to insert prices', {
|
||||
errorType: error.name,
|
||||
errorMessage: error.message
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
fetchPage,
|
||||
parseUnits,
|
||||
convertDataTypes,
|
||||
upsertUnits,
|
||||
insertPrices,
|
||||
// Export helpers for testing
|
||||
parseInteger,
|
||||
parsePositiveInteger,
|
||||
|
||||
Reference in New Issue
Block a user