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

如何高效比对百万级爬虫数据与MySQL库并避免连接超时?

百万级数据批量去重插入优化方案

问题背景

我在MySQL中存储了500万条数据,表结构如下:

id (primary_key), name, phone_number
1, John Doe, 12346789

工具持续爬取百万级新数据,每次将1万条数据传入processDataInChunks函数,需要比对数据库条目,插入不存在的记录。原Sequelize实现代码如下:

// data = chunk of data collected by the scraper
// callback = function to call when new entry is detected
function processDataInChunks(data, callback) {
   data.map(function (entry) { // loop through the data array where entry is nth element of the array
      db.findAll({phone_number: entry.phone_number}) // db variable represents the SQL table
         .then(function (rows) { // call this function with rows as argument when the query is successful
            if (!rows.length > 0) { // if phone number is not in the database
               db.create({ // create entry in the database
                  name: entry.name,
                  phone_number: entry.phone_number
               }).then(function () {
                  callback(entry);
                  console.log(`Found a new phone number: ${entry}`)
               }).catch(err=>console.log(err))

            }
         }).catch(err=>console.log(err))
   })
}

运行时出现ConnectionAcquireTimeoutError(连接池耗尽),改用async/await后仍耗时极长,求最优实现方案。


核心优化思路

原代码的问题在于单条数据循环发起数据库请求,1万条数据会瞬间创建1万个异步请求,直接耗尽连接池,同时大量IO往返导致效率极低。优化方向是批量操作+数据库端存在性判断,减少数据库交互次数。

1. 先给phone_number添加唯一索引

这是所有优化的基础,既可以加速存在性查询,又能避免重复插入:

ALTER TABLE your_table_name ADD UNIQUE INDEX idx_phone_number (phone_number);

2. 方案一:批量查询+批量插入(Sequelize封装方式)

先一次性查询当前批次中所有已存在的手机号,过滤出不存在的记录后批量插入,仅需2次数据库请求:

async function processDataInChunks(data, callback) {
  try {
    // 1. 提取当前批次所有手机号
    const phoneNumbers = data.map(entry => entry.phone_number);
    
    // 2. 批量查询已存在的手机号,仅返回phone_number字段
    const existingEntries = await db.findAll({
      attributes: ['phone_number'],
      where: {
        phone_number: phoneNumbers
      }
    });
    
    // 3. 转换为Set方便快速查找
    const existingPhones = new Set(existingEntries.map(item => item.phone_number));
    
    // 4. 过滤出不存在的记录
    const newEntries = data.filter(entry => !existingPhones.has(entry.phone_number));
    
    if (newEntries.length === 0) {
      console.log("No new entries to insert");
      return;
    }
    
    // 5. 批量插入新记录
    await db.bulkCreate(newEntries, {
      // 可选:如果需要触发模型钩子,设置为true,否则false更高效
      hooks: false
    });
    
    // 6. 批量触发回调
    newEntries.forEach(entry => callback(entry));
    console.log(`Inserted ${newEntries.length} new entries`);
  } catch (err) {
    console.error("Error processing data chunk:", err);
  }
}

3. 方案二:用MySQL原生语法直接处理(最快方式)

利用MySQL的INSERT IGNORE或INSERT ... ON DUPLICATE KEY UPDATE,让数据库直接处理存在性判断,仅需1次数据库请求,效率最高:

async function processDataInChunks(data, callback) {
  try {
    // 1. 构造批量插入的SQL语句(用Sequelize的query方法执行原生SQL)
    const values = data.map(entry => `('${entry.name}', '${entry.phone_number}')`).join(',');
    
    // 使用INSERT IGNORE:如果手机号已存在则忽略该条插入
    const [result] = await db.query(`
      INSERT IGNORE INTO your_table_name (name, phone_number)
      VALUES ${values}
    `);
    
    // 获取实际插入的行数
    console.log(`Inserted ${result.affectedRows} new entries`);
    
    // 若需要精准触发callback,可先过滤出未存在的记录(同方案一的过滤逻辑)
    // newEntries.forEach(entry => callback(entry));
  } catch (err) {
    console.error("Error processing data chunk:", err);
  }
}

注意:如果需要同步更新已存在记录的name字段,改用以下语句:

INSERT INTO your_table_name (name, phone_number)
VALUES ${values}
ON DUPLICATE KEY UPDATE name = VALUES(name);

4. 连接池配置优化

调整Sequelize的连接池参数,适配批量操作的需求:

const sequelize = new Sequelize('database', 'username', 'password', {
  host: 'localhost',
  dialect: 'mysql',
  pool: {
    max: 20, // 最大连接数,根据服务器配置调整,不要超过MySQL的max_connections
    min: 5, // 最小空闲连接数
    acquire: 30000, // 连接超时时间(毫秒)
    idle: 10000 // 连接空闲超时时间(毫秒)
  }
});

5. 并发控制(可选)

如果批量操作仍有压力,可将1万条数据拆分为更小的批次(比如每1000条一批),限制并发数:

async function processDataInChunks(data, callback) {
  const batchSize = 1000;
  const batches = [];
  
  // 拆分数据为小批次
  for (let i = 0; i < data.length; i += batchSize) {
    batches.push(data.slice(i, i + batchSize));
  }
  
  // 串行处理每个小批次,避免并发过高
  for (const batch of batches) {
    await processSingleBatch(batch, callback);
  }
}

// 复用方案一或方案二的批量处理逻辑
async function processSingleBatch(batch, callback) {
  // 这里放方案一或方案二的代码
}

内容的提问来源于stack exchange,提问作者n00b

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 05:27:15