You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在MongoDB集合中批量存储多接口返回数据(Node.js环境)

Solution for Batch Storing API Results in MongoDB Single Collection

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 apiUniqueId with 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_SIZE to 500 or lower to avoid hitting MongoDB's batch size limits.
  • API Fetch Logic: Swap out the fetchApiData function 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: bulkWrite is 10-100x faster than individual insertOne/updateOne calls for large datasets.
  • Idempotency: Using upsert: true means 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 08:38:51