如何高效比对百万级爬虫数据与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
相关产品推荐
相关产品推荐

