在MongoDB集合中批量存储多接口返回数据(Node.js环境)
Great question! Handling large batch inserts/updates into MongoDB efficiently is key here, especially when dealing with 100k+ total records. Let's walk through a practical, efficient solution that focuses on batch processing after each API call, with built-in safeguards for duplicates and performance.
1. Core Strategy Overview
Instead of inserting records one by one (which is slow and resource-heavy), we'll use MongoDB's bulkWrite API for each API response. This method lets us process hundreds of operations in a single request, drastically improving speed. We'll also use upsert logic to avoid duplicate entries if your API returns overlapping data across param calls.
2. Step-by-Step Implementation
We'll use the native MongoDB driver (since it's lightweight and direct for bulk operations) and include placeholder code for your API fetch (replace this with your existing request-based script).
First, Install Dependencies
npm install mongodb axios # Use axios as a stand-in; replace with your request module
Full Code Example
const { MongoClient } = require('mongodb'); const axios = require('axios'); // Replace this with your existing request-based logic // MongoDB Configuration const MONGO_URI = 'mongodb://localhost:27017'; // Update with your URI const DB_NAME = 'your_target_db'; const COLLECTION_NAME = 'single_combined_collection'; // API Configuration const API_BASE_URL = 'https://your-api-endpoint.com/results'; const PARAMS = ['param1', 'param2', 'param3']; // Your three distinct params // Chunk size for bulk operations (adjust based on document size; 1000 is safe for most cases) const CHUNK_SIZE = 1000; // Helper: Connect to MongoDB async function getMongoCollection() { const client = new MongoClient(MONGO_URI); await client.connect(); console.log('✅ Connected to MongoDB'); return client.db(DB_NAME).collection(COLLECTION_NAME); } // Helper: Fetch data from API (replace this with your existing request code) async function fetchApiData(param) { try { const response = await axios.get(API_BASE_URL, { params: { filter: param } }); return response.data; // Assume this returns an array of 30k+ objects } catch (err) { console.error(`❌ Failed to fetch data for param: ${param}`, err.message); throw err; // Re-throw to halt or handle upstream } } // Helper: Split large arrays into smaller chunks function chunkArray(arr, size) { const chunks = []; for (let i = 0; i < arr.length; i += size) { chunks.push(arr.slice(i, i + size)); } return chunks; } // Helper: Batch upsert documents to MongoDB async function batchUpsert(collection, documents) { const chunks = chunkArray(documents, CHUNK_SIZE); for (const chunk of chunks) { const operations = chunk.map(doc => ({ updateOne: { filter: { apiUniqueId: doc.id }, // Use a unique field from your API data update: { $set: doc }, upsert: true // Insert if new, update if exists } })); try { const result = await collection.bulkWrite(operations); console.log(`📝 Processed chunk: ${result.upsertedCount} inserted, ${result.modifiedCount} updated`); } catch (err) { console.error('❌ Bulk operation failed', err.message); throw err; } } } // Main workflow async function run() { let collection; try { collection = await getMongoCollection(); // Create a unique index to enforce duplicate prevention (run once, then comment out) await collection.createIndex({ apiUniqueId: 1 }, { unique: true }); // Process each param sequentially for (const param of PARAMS) { console.log(`\n🔄 Starting processing for param: ${param}`); const apiData = await fetchApiData(param); await batchUpsert(collection, apiData); console.log(`✅ Completed processing for param: ${param}`); } console.log('\n🎉 All data successfully stored in MongoDB!'); } catch (err) { console.error('\n❌ Workflow failed:', err.message); } finally { if (collection) { await collection.client.close(); console.log('🔌 MongoDB connection closed'); } } } run();
3. Key Details to Customize
- Unique Identifier: Replace
apiUniqueIdwith the actual unique field from your API responses (e.g.,userId,recordId). This ensures we don't insert duplicates. - Chunk Size: If your documents are very large (100KB+ each), reduce
CHUNK_SIZEto 500 or lower to avoid hitting MongoDB's batch size limits. - API Fetch Logic: Swap out the
fetchApiDatafunction with your existing request-based script—keep the return value as an array of documents. - Error Handling: Adjust the try/catch blocks to fit your needs (e.g., retry failed API calls, skip faulty chunks instead of halting).
4. Performance & Reliability Tips
- Bulk Operations:
bulkWriteis 10-100x faster than individualinsertOne/updateOnecalls for large datasets. - Idempotency: Using
upsert: truemeans you can rerun the script safely without creating duplicates. - Indexing: The unique index we create ensures data integrity and speeds up the filter checks during upserts.
- Logging: The included logs help you track progress and debug any failures quickly.
内容的提问来源于stack exchange,提问作者filipbarak

