Implement insertPrices() bulk operation for daily price records
All checks were successful
CI/CD Pipeline - Apartment API / Scan Dependencies (pull_request) Successful in 13s
CI/CD Pipeline - Apartment API / Run Tests (pull_request) Successful in 9m47s
CI/CD Pipeline - Apartment API / Send Webhook Notification (pull_request) Successful in 2s
CI/CD Pipeline - Apartment API / Run Linting (pull_request) Successful in 9m37s
CI/CD Pipeline - Apartment API / Build & Push Image (pull_request) Has been skipped
CI/CD Pipeline - Apartment API / Deploy to Production (pull_request) Has been skipped

Add insertPrices() to scraperService that performs bulk upsert of
daily unit price records to the prices collection. Key behaviors:

- Uses bulkWrite with updateOne/upsert ops for atomic batch writes
- Compound key (unit_code + date_checked) ensures idempotent re-runs
- Filters out units with null prices before writing
- Stores price, last_updated, and data_source fields per record
- Uses ordered:false for parallel write performance
- Logs warnings for empty/all-null input without hitting the database
- Propagates bulkWrite errors with detailed logging

Includes 30 tests covering filtering, upsert behavior, composite key
logic, field validation, idempotency, error handling, and logging.
This commit is contained in:
2026-02-05 22:20:33 -07:00
parent 4863dc1824
commit 406672c16b
2 changed files with 646 additions and 0 deletions

View File

@ -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,