Puppeteer-Cluster未按预期并行执行,如何实现多任务并行截图?
Puppeteer-Cluster 无法并行执行任务的问题排查与修复
问题描述
我使用puppeteer-cluster创建了Puppeteer Worker集群,配置如下:
const cluster = await Cluster.launch({ concurrency: Cluster.CONCURRENCY_PAGE, puppeteerOptions: { userDataDir: path.join(__dirname,'user_data/1'), headless: false, args: ['--no-sandbox'] }, maxConcurrency: maxCon, monitor: true, skipDuplicateUrls: true, timeout:40000000, retryLimit:5, });
随后通过循环将一批URL加入队列,任务为捕获网站截图。脚本可正常运行,但未按预期并行执行,而是串行处理——浏览器逐个切换标签页完成截图后再处理下一个。请问如何调整才能实现多任务并行?
完整代码
const puppeteer = require('puppeteer'); const { Cluster } = require('puppeteer-cluster'); const fs = require('fs'); const path = require('path'); var pdfkit = require('pdfkit'); //const zip = require('./zip_files'); //const cfolder = require('./create_folders'); const site = 'scribd.com'; const docType = ['pdf', 'word', 'spreadsheet']; const t_out = 10000; const wait = ms => new Promise(res => setTimeout(res, ms)); const scrnDir = 'screenshots'; const docDir = 'documents'; const zipDir = 'zips'; var data_1 = ['Exporter']; var data_2 = []; (async() => { const browser = await puppeteer.launch({ headless: false, userDataDir: path.join(__dirname,'user_data/main'), }); const page = (await browser.pages())[0]; for(let i = 0; i < data_1.length; i++){ // for(let j = 0; j < data_2.length; j++){ var numFiles = 1000000; let folder = data_1[i].replace(/\n/gm, '').replace(/\r/gm, ''); let searchTerm = data_1[i].replace(/\n/gm, '').replace(/\r/gm, ''); for(let pageNum = 1; pageNum < 2/*(Math.ceil(numFiles/42) +1)*/ && pageNum < 236; pageNum++){ //maxPageNum = 235 let docType = 'pdf'; let query = 'https://www.scribd.com/search?query='+searchTerm+'&content_type=documents&page='+pageNum+'&filetype='+docType; await page.goto(query, {waitUntil : 'networkidle2'}); //await cfolder.createFolder(docDir, searchTerm); fs.appendFileSync('progress/query.txt', query + '\n'); if(pageNum == 1){ let numFiles = await fileCount(page); } let docPages = await page.waitForXPath('//section[@data-testid="search-results"]', { timeout: t_out }).then(async() => { let searchResults = await page.$x('//section[@data-testid="search-results"]'); await searchResults[0].waitForXPath('//div/ul/li'); let docPages = await searchResults[0].$x('//div/ul/li'); return docPages; }).catch( e => { console.log('getLinks Error'); console.log(e); }); await save(browser, searchTerm, docPages); } //await zip.zipFolder(docDir + '/' + folder, zipDir + '/' + searchTerm + '.zip'); // } } })(); async function save(browser, searchTerm, docPages){ //let docPage = await browser.newPage(); let maxCon = 3; const cluster = await Cluster.launch({ concurrency: Cluster.CONCURRENCY_PAGE, puppeteerOptions: { userDataDir: path.join(__dirname,'user_data/1'), headless: false, args: ['--no-sandbox'] }, maxConcurrency: maxCon, monitor: true, skipDuplicateUrls: true, timeout:40000000, retryLimit:5, }); await cluster.task(async ({ page, data: {url, title} }) => { let docPage = page; await docPage.goto(url, {waitUntil: 'networkidle2'}); //await cfolder.createFolder(scrnDir, title); await docPage.evaluate('document.querySelector(".nav_and_banners_fixed").remove()'); await docPage.evaluate('document.querySelector(".recommender_list_wrapper").remove()'); await docPage.evaluate('document.querySelector(".auto__doc_page_app_page_body_fixed_viewport_bottom_components").remove()'); //await autoScroll(docPage); //await docPage.evaluate('document.querySelector(".wrapper__doc_page_webpack_doc_page_body_document_useful").remove()'); await docPage.addStyleTag({content: '.wrapper__doc_page_webpack_doc_page_body_document_useful{visibility: hidden}'}) await docPage.waitForXPath('//span[@class="page_of"]'); let numOfPagesR = await docPage.$x('//span[@class="page_of"]'); let numOfPages = parseInt((await (await numOfPagesR[0].getProperty('textContent')).jsonValue()).split('of ').pop()); console.log(numOfPages); //const pages = await docPage.$x('//*[@class="newpage"]'); let imgs = []; for(let j = 0; j < numOfPages; j++){ let sel = '//*[@id="page' + (j+1) + '"]'; let pages = await docPage.$x(sel); await pages[0].screenshot({ path: scrnDir + '/' + title +j+'.jpg' }); imgs[j] = title + j +'.jpg'; } //await createPdf(searchTerm, title, imgs); }); cluster.on('taskerror', (err, data) => { console.log(` Error crawling ${data}: ${err.message}`); }); for(let i = 0; i < 6/**docPages.length*/; i++){ await docPages[i].waitForXPath('//article/a'); let urlR = await docPages[i].$x('//article/a'); let url = await (await urlR[0].getProperty('href')).jsonValue(); await docPages[i].waitForXPath('//p[@data-e2e="title"]'); let titleR = await docPages[i].$x('//p[@data-e2e="title"]'); let title = await (await titleR[0].getProperty('textContent')).jsonValue(); cluster.queue({url : url, title : title}); //console.log(title); } await cluster.idle(); //docPage.close(); } async function fileCount(page){ await page.waitForXPath('//div[@class="_7a1igU"]', { timeout: t_out }).then(async() => { let fileCountR = await page.$x('//div[@class="_7a1igU"]'); let fileCountS = await (await fileCountR[0].getProperty('textContent')).jsonValue(); let numFiles = parseInt(fileCountS.split('of ').pop().split(' results').shift().replace(/,/g, '')); console.log('Total Files : ' + numFiles); return numFiles; }).catch( e => { console.log('File Count Error'); console.log(e); }); } async function getLinks(page){ } async function createPdf(searchTerm, title, images){ //await cfolder.createFolder(docDir, searchTerm); let pdf = new pdfkit({ autoFirstPage: false }); let writeStream = fs.createWriteStream(docDir+ '/' + searchTerm + '/' + title + '.pdf'); pdf.pipe(writeStream); for(let i = 0; i < images.length; i++){ let img = pdf.openImage('./' + scrnDir + '/' + title + '/' + images[i]); pdf.addPage({size: [img.width, img.height]}); pdf.image(img, 0, 0); } pdf.end(); await new Promise(async (resolve) => { writeStream.on('close', ()=>{ console.log('PDF Created succesfully'); resolve(); }); }); }
注:
const zip = require('./zip_files');和const cfolder = require('./create_folders');为最终代码所需,与当前问题无关。
问题原因与修复方案
核心问题分析
- 集群实例重复创建:每次调用
save函数都会新建一个Cluster实例,集群无法复用资源,且每个小集群只能串行处理自身队列的任务。 - 主流程强制串行:主函数中通过
await save(...)等待每个save调用完成才继续循环,即使集群内部支持并行,外部的阻塞也会让整体任务变成串行执行。 - 任务逻辑重复绑定:每次创建集群都重新绑定
cluster.task处理函数,属于冗余操作,且可能导致调度冲突。
具体修复步骤
1. 全局初始化集群,仅创建一次
将Cluster.launch移到主函数最开始,只初始化一个集群实例,所有任务统一加入该集群队列。
2. 移除save函数,统一任务处理逻辑
把原save中的任务处理逻辑绑定到全局集群上,避免重复绑定和实例创建。
3. 主流程取消阻塞,批量加入任务
主函数中遍历搜索结果后,直接将任务加入集群队列,无需等待任务处理完成,让集群自行调度并行执行。
修改后的代码示例
const puppeteer = require('puppeteer'); const { Cluster } = require('puppeteer-cluster'); const fs = require('fs'); const path = require('path'); var pdfkit = require('pdfkit'); const site = 'scribd.com'; const docType = ['pdf', 'word', 'spreadsheet']; const t_out = 10000; const wait = ms => new Promise(res => setTimeout(res, ms)); const scrnDir = 'screenshots'; const docDir = 'documents'; const zipDir = 'zips'; var data_1 = ['Exporter']; var data_2 = []; (async() => { // 1. 全局初始化集群,仅创建一次 let maxCon = 3; const cluster = await Cluster.launch({ concurrency: Cluster.CONCURRENCY_PAGE, puppeteerOptions: { userDataDir: path.join(__dirname,'user_data/1'), headless: false, args: ['--no-sandbox'] }, maxConcurrency: maxCon, monitor: true, skipDuplicateUrls: true, timeout:40000000, retryLimit:5, }); // 2. 统一绑定任务处理逻辑 await cluster.task(async ({ page, data: {url, title, searchTerm} }) => { await page.goto(url, {waitUntil: 'networkidle2'}); // 优化DOM操作:添加可选链避免元素不存在报错 await page.evaluate('document.querySelector(".nav_and_banners_fixed")?.remove()'); await page.evaluate('document.querySelector(".recommender_list_wrapper")?.remove()'); await page.evaluate('document.querySelector(".auto__doc_page_app_page_body_fixed_viewport_bottom_components")?.remove()'); await page.addStyleTag({content: '.wrapper__doc_page_webpack_doc_page_body_document_useful{visibility: hidden}'}) try { await page.waitForXPath('//span[@class="page_of"]', { timeout: t_out }); let numOfPagesR = await page.$x('//span[@class="page_of"]'); let numOfPages = parseInt((await (await numOfPagesR[0].getProperty('textContent')).jsonValue()).split('of ').pop()); console.log(`文档 ${title} 总页数:${numOfPages}`); let imgs = []; for(let j = 0; j < numOfPages; j++){ let sel = '//*[@id="page' + (j+1) + '"]'; let pages = await page.$x(sel); if(pages.length > 0){ await pages[0].screenshot({ path: `${scrnDir}/${title}${j}.jpg` }); imgs[j] = `${title}${j}.jpg`; } } //await createPdf(searchTerm, title, imgs); } catch(e) { console.log(`处理文档 ${title} 出错:${e.message}`); } }); cluster.on('taskerror', (err, data) => { console.log(` Error crawling ${data.url}: ${err.message}`); }); // 初始化主浏览器用于搜索 const browser = await puppeteer.launch({ headless: false, userDataDir: path.join(__dirname,'user_data/main'), }); const page = (await browser.pages())[0]; for(let i = 0; i < data_1.length; i++){ let searchTerm = data_1[i].replace(/\n/gm, '').replace(/\r/gm, ''); for(let pageNum = 1; pageNum < 2 && pageNum < 236; pageNum++){ let docType = 'pdf'; let query = `https://www.scribd.com/search?query=${searchTerm}&content_type=documents&page=${pageNum}&filetype=${docType}`; await page.goto(query, {waitUntil : 'networkidle2'}); fs.appendFileSync('progress/query.txt', query + '\n'); if(pageNum == 1){ await fileCount(page); } try { await page.waitForXPath('//section[@data-testid="search-results"]', { timeout: t_out }); let searchResults = await page.$x('//section[@data-testid="search-results"]'); await searchResults[0].waitForXPath('//div/ul/li'); let docPages = await searchResults[0].$x('//div/ul/li'); // 3. 批量加入任务,不阻塞主流程 for(let i = 0; i < 6/*docPages.length*/; i++){ await docPages[i].waitForXPath('//article/a'); let urlR = await docPages[i].$x('//article/a'); let url = await (await urlR[0].getProperty('href')).jsonValue(); await docPages[i].waitForXPath('//p[@data-e2e="title"]'); let titleR = await docPages[i].$x('//p[@data-e2e="title"]'); let title = await (await titleR[0].getProperty('textContent')).jsonValue(); cluster.queue({url, title, searchTerm}); } } catch( e ) { console.log('getLinks Error'); console.log(e); } } } // 等待所有任务完成后关闭集群和浏览器 await cluster.idle(); await cluster.close(); await browser.close(); })(); async function fileCount(page){ try { await page.waitForXPath('//div[@class="_7a1igU"]', { timeout: t_out }); let fileCountR = await page.$x('//div[@class="_7a1igU"]'); let fileCountS = await (await fileCountR[0].getProperty('textContent')).jsonValue(); let numFiles = parseInt(fileCountS.split('of ').pop().split(' results').shift().replace(/,/g, '')); console.log('Total Files : ' + numFiles); return numFiles; } catch( e ) { console.log('File Count Error'); console.log(e); } } async function createPdf(searchTerm, title, images){ let pdf = new pdfkit({ autoFirstPage: false }); let writeStream = fs.createWriteStream(docDir+ '/' + searchTerm + '/' + title + '.pdf'); pdf.pipe(writeStream); for(let i = 0; i < images.length; i++){ let img = pdf.openImage('./' + scrnDir + '/' + title + '/' + images[i]); pdf.addPage({size: [img.width, img.height]}); pdf.image(img, 0, 0); } pdf.end(); await new Promise(async (resolve) => { writeStream.on('close', ()=>{ console
相关产品推荐
相关产品推荐

