如何在puppeteer-cluster中正确实现多函数任务调度
问题根源
代码存在两个核心问题导致任务调度异常:
- 第一阶段链接提取任务执行完成后,你提前调用
await cluster.close()销毁了整个集群实例。关闭后的集群不再接收、调度新入队任务,后续入队的监控任务只有极个别可能在集群资源完全释放前抢到执行权,其余任务全部不会被调度。 - 第二阶段的页面监控任务通过
while(true)实现无限循环刷新,单任务会永久占用集群工作槽且永不释放。puppeteer-cluster的最大并行工作槽数量由maxConcurrency参数控制,工作槽被占满后后续任务会一直处于排队状态,永远不会被执行。
注:puppeteer-cluster原生支持为不同入队任务指定独立处理函数,你使用的多任务函数传参方式本身不存在API错误。
修复方案
按以下规则调整代码即可:
- 删除第一阶段任务结束后的
await cluster.close()调用,集群实例需要保留到所有任务入队完成,常驻监控场景不需要提前关闭集群。 - 由于监控任务是永久常驻不退出的,需要将
maxConcurrency设置为大于等于待监控URL的总数量,保证每个监控任务都能分配到独立工作槽执行。 - 提取到的
href属性可能是相对路径,需要拼接为完整绝对路径后再入队,避免page.goto()抛出路径错误。 - 常驻监控场景下任务永远不会执行完成,不需要调用
await cluster.idle()等待队列空闲,可以添加进程终止信号监听,在手动终止程序时再安全关闭集群、释放资源。 - 可选优化:在元素抓取逻辑外层增加错误捕获,避免页面结构变动导致任务频繁重试。
修正后核心代码
const { Cluster } = require('puppeteer-cluster'); const fs = require("fs"); const { addExtra } = require("puppeteer-extra"); const vanillaPuppeteer = require("puppeteer"); const StealthPlugin = require("puppeteer-extra-plugin-stealth"); var moment = require('moment'); var regexTemps = /(\d+)\s(\w+)$/; const urlsToCheck = []; const ENTRY_URL = 'https://www.apagewithsomelinks.com'; TZ = 'Europe/Paris' process.env.TZ = 'Europe/Paris' (async () => { const puppeteer = addExtra(vanillaPuppeteer); puppeteer.use(StealthPlugin()); const cluster = await Cluster.launch({ puppeteer, puppeteerOptions: { headless: false, args: ['--no-sandbox'], }, maxConcurrency: 20, // 根据实际待监控URL总数调整,需大于URL总数 concurrency: Cluster.CONCURRENCY_CONTEXT, monitor: false, skipDuplicateUrls: true, timeout:30000, retryLimit:10, }) cluster.on('taskerror', (err, data, willRetry) => { if (willRetry) { console.warn(`Encountered an error while crawling ${data}. ${err.message}\nThis job will be retried`); } else { console.error(`Failed to crawl ${data}: ${err.message}`); } }); // 进程终止时安全关闭集群 process.on('SIGINT', async () => { await cluster.close(); process.exit(0); }); const getElementOnPage = async ({ page, data: url }) => { console.log('=> Go to URL : ',url); await page.goto(url, {waitUntil: 'domcontentloaded'}); while (true) { try { console.log('=> Reload URL : ',page.url()) await page.reload({waitUntil: 'domcontentloaded'}); await page.waitForTimeout(1000); let allNews = await page.$$("article.news"); if (!allNews.length) { console.log(new Date(), 'No news element found, retry next loop'); await page.waitForTimeout(1000); continue; } let firstNews = allNews[0]; let info = await firstNews.$eval('.info span', s => s.textContent.trim()); console.log(new Date(), 'info : ',info); } catch (err) { console.log(new Date(), 'Crawl error:', err.message); } finally { await page.waitForTimeout(1000); } } }; const getListOfPagesToExplore = async ({ page, data: url }) => { console.log(new Date(), 'Get the list of deal pages to explore'); await page.goto(url, {waitUntil: 'domcontentloaded'}); await page.waitForTimeout(500); const hrefsToVisit = await page.$x('//a'); for( let hrefToVisit of hrefsToVisit ) { var link = await page.evaluate(el => el.getAttribute("href"), hrefToVisit); // 拼接相对路径为绝对路径 if (link && !link.startsWith('http')) { link = new URL(link, ENTRY_URL).href; } if (link) { console.log(new Date(), 'adding link to list : ', link); urlsToCheck.push(link); } } }; // 执行第一阶段链接提取 cluster.queue(ENTRY_URL, getListOfPagesToExplore); await cluster.idle(); console.log('Total urls to monitor:', urlsToCheck.length); // 入队所有第二阶段监控任务 for( let url of urlsToCheck ) { console.log('Push in queue : ',url); cluster.queue(url, getElementOnPage); } // 常驻任务不需要等待idle,保持进程运行即可 })();
内容的提问来源于stack exchange,提问作者user2178964
相关产品推荐
相关产品推荐

