NodeJS中使用fork时等待所有子进程完成的实现方案
问题
我编写了一个Node.js脚本,用于遍历URL列表下载PDF文件,每个URL对应5个PDF,因此使用fork创建子进程并行下载这些PDF。目前需要解决两个问题:
- 仅当当前URL对应的所有子进程(下载任务)全部执行完成后,才继续处理下一个URL。
child.js中执行stream.close()后,子进程并未退出。
原main.js代码
const puppeteer = require('puppeteer'); const fs = require('fs'); const fork = require('child_process').fork; const ls = fork("download.js"); var list = []; var links = []; var names = []; (async () => { const browser = await puppeteer.launch(({headless: false})); const page = await browser.newPage(); /* const list_arr = fs.readFileSync('link_list.csv').toString().split(","); for(l = 1; l < list_arr.length; l+2){ links[(l - 1) / 2] = list_arr[l]; names([l - 1] / 2) = list_arr[l - 1]; } for(let i = 0; i < links.length; i++){*/ // url = links[i]; // name = names[i]; url = 'https://www.responsibilityreports.com/Company/abb-ltd'; name = 'abcd'; await page.goto(url, { waitUntil: 'networkidle2' }); await page.waitForXPath('//li[@class="top_content_list"]/div[@class="left"]/span[@class="ticker_name"]'); let ticker_r = await page.$x('//li[@class="top_content_list"]/div[@class="left"]/span[@class="ticker_name"]'); let ticker = await (await ticker_r[0].getProperty('textContent')).jsonValue(); let cat = ticker[0].toLowerCase(); //console.log(ticker); await page.waitForXPath('//li[@class="top_content_list"]/div[@class="right"]'); let exchange_r = await page.$x('//li[@class="top_content_list"]/div[@class="right"]'); let exchange = (await (await exchange_r[0].getProperty('textContent')).jsonValue()).split('Exchange').pop().split('More').shift().replace(/\n/gm, '').trim(); //console.log(exchange); await page.waitForXPath('//div[@class="most_recent_content_block"]/span[@class="bold_txt"]'); let mr_res = await page.$x('//div[@class="most_recent_content_block"]/span[@class="bold_txt"]'); let mr_txt = await mr_res[0].getProperty('textContent'); let mr_text = await mr_txt.jsonValue(); console.log(mr_text); let [mr_year, ...mr_rest] = mr_text.split(' '); let mr_type = mr_rest.join(' ').trim(); mr_url = 'https://www.responsibilityreports.com/HostedData/ResponsibilityReports/PDF/'+ exchange + '_' + ticker + '_' + mr_year + '.pdf'; let mr_obj = { "year" : mr_year, "type" : mr_type.trim(), "url" : mr_url } list.push(mr_obj); await page.waitForXPath('//div[@class="archived_report_content_block"]/ul/li/div/span[@class="heading"]'); let ar_reps = await page.$x('//div[@class="archived_report_content_block"]/ul/li/div/span[@class="heading"]'); console.log(ar_reps.length); for(let k = 0; k < ar_reps.length; k++){ let ar_txt = await ar_reps[k].getProperty('textContent'); let ar_text = await ar_txt.jsonValue(); let [ar_year, ...ar_rest] = ar_text.split(' '); if(parseInt(ar_year) < 2017){ break; } let ar_type = ar_rest.join(' '); ar_url = 'https://www.responsibilityreports.com/HostedData/ResponsibilityReportArchive/' + cat + '/'+ exchange + '_' + ticker + '_' + ar_year + '.pdf'; let ar_obj = { "year" : ar_year, "type" : ar_type, "url" : ar_url } list.push(ar_obj); console.log(ar_year); } /*} for(f = 0; f < list.length; f++){ let url_s = list[f].url; let type_s = list[f].type; let year_s = list[f].year; ls.on('exit', (code)=>{ console.log('child_process exited with code ${code}'); }); ls.on('message', (msg) => { ls.send([url_s ,name,type_s,year_s]); console.log(msg); }); } await browser.close(); console.log('done'); })();
原child.js代码
const down_path = 'downloads/'; const https = require('https'); const fs = require('fs'); process.on('message', async (arr)=> { console.log("CHILD: url received from parent process", arr); url = arr[0]; name = arr[1]; type = arr[2]; year = arr[3]; await download(url,name,type,year); }); process.send('Executed'); async function download(url,name,type,year) { https.get(url, res => { const stream = fs.createWriteStream(down_path + name + '_' + type + '_' + year + '.pdf'); res.pipe(stream); stream.on('finish', async() => { console.log('done : ' + year); stream.close(); }); }); }
解决方案
一、实现按URL批次等待所有子进程完成
原代码仅创建了单个子进程,且消息监听逻辑无法批量控制任务完成状态。我们需要为每个下载任务创建独立子进程,并用Promise封装子进程生命周期,确保当前URL的所有下载任务完成后再处理下一个URL。
修改后的main.js
const puppeteer = require('puppeteer'); const fs = require('fs'); const fork = require('child_process').fork; // 封装子进程下载逻辑,返回Promise等待任务完成 function spawnDownloader(taskData) { return new Promise((resolve, reject) => { const child = fork('./download.js'); child.send(taskData); child.on('message', (msg) => { console.log(`子进程消息: ${msg}`); }); child.on('exit', (code) => { console.log(`子进程退出,代码: ${code}`); if (code === 0) { resolve(); } else { reject(new Error(`子进程异常退出,代码: ${code}`)); } }); child.on('error', (err) => { reject(err); }); }); } (async () => { const browser = await puppeteer.launch({ headless: false }); const page = await browser.newPage(); // 恢复CSV读取逻辑(修正原循环错误) const list_arr = fs.readFileSync('link_list.csv').toString().split(","); const links = []; const names = []; // 修正原循环的l+2错误,应该是l += 2 for (let l = 1; l < list_arr.length; l += 2) { const index = (l - 1) / 2; links[index] = list_arr[l]; names[index] = list_arr[l - 1]; } // 遍历每个URL,处理完成所有PDF后再进入下一个 for (let i = 0; i < links.length; i++) { const url = links[i]; const name = names[i]; const list = []; console.log(`开始处理URL: ${url}`); await page.goto(url, { waitUntil: 'networkidle2' }); // 获取ticker await page.waitForXPath('//li[@class="top_content_list"]/div[@class="left"]/span[@class="ticker_name"]'); const ticker_r = await page.$x('//li[@class="top_content_list"]/div[@class="left"]/span[@class="ticker_name"]'); const ticker = await (await ticker_r[0].getProperty('textContent')).jsonValue(); const cat = ticker[0].toLowerCase(); // 获取exchange await page.waitForXPath('//li[@class="top_content_list"]/div[@class="right"]'); const exchange_r = await page.$x('//li[@class="top_content_list"]/div[@class="right"]'); const exchange = (await (await exchange_r[0].getProperty('textContent')).jsonValue()) .split('Exchange').pop().split('More').shift().replace(/\n/gm, '').trim(); // 处理最新报告 await page.waitForXPath('//div[@class="most_recent_content_block"]/span[@class="bold_txt"]'); const mr_res = await page.$x('//div[@class="most_recent_content_block"]/span[@class="bold_txt"]'); const mr_txt = await mr_res[0].getProperty('textContent'); const mr_text = await mr_txt.jsonValue(); const [mr_year, ...mr_rest] = mr_text.split(' '); const mr_type = mr_rest.join(' ').trim(); const mr_url = `https://www.responsibilityreports.com/HostedData/ResponsibilityReports/PDF/${exchange}_${ticker}_${mr_year}.pdf`; list.push({ year: mr_year, type: mr_type, url: mr_url, name: name }); // 处理归档报告 await page.waitForXPath('//div[@class="archived_report_content_block"]/ul/li/div/span[@class="heading"]'); const ar_reps = await page.$x('//div[@class="archived_report_content_block"]/ul/li/div/span[@class="heading"]'); for (let k = 0; k < ar_reps.length; k++) { const ar_txt = await ar_reps[k].getProperty('textContent'); const ar_text = await ar_txt.jsonValue(); const [ar_year, ...ar_rest] = ar_text.split(' '); if (parseInt(ar_year) < 2017) { break; } const ar_type = ar_rest.join(' '); const ar_url = `https://www.responsibilityreports.com/HostedData/ResponsibilityReportArchive/${cat}/${exchange}_${ticker}_${ar_year}.pdf`; list.push({ year: ar_year, type: ar_type, url: ar_url, name: name }); console.log(`添加归档报告: ${ar_year}`); } // 并行启动所有下载任务,等待全部完成 console.log(`开始并行下载${list.length}个PDF`); const downloadPromises = list.map(task => spawnDownloader([task.url, task.name, task.type, task.year])); await Promise.all(downloadPromises); console.log(`当前URL所有PDF下载完成,进入下一个URL`); } await browser.close(); console.log('所有任务完成'); })();
关键修改点
- 新增
spawnDownloader函数,将子进程的创建、消息监听和退出逻辑封装为Promise,方便批量等待。 - 修正原CSV读取循环的错误(
l+2改为l += 2),并使用块级变量避免全局污染。 - 将每个URL的PDF任务收集后,用
Promise.all等待所有子进程完成,再进入下一个URL的处理。
二、修复子进程不退出的问题
原child.js存在两个核心问题:download函数未正确等待流完成,且下载结束后未主动退出子进程。
修改后的child.js
const down_path = 'downloads/'; const https = require('https'); const fs = require('fs'); const path = require('path'); // 确保下载目录存在 fs.mkdirSync(down_path, { recursive: true }); process.on('message', async (arr) => { console.log("CHILD: 收到父进程任务", arr); const [url, name, type, year] = arr; try { await download(url, name, type, year); process.send(`下载完成: ${year}`); process.exit(0); // 主动退出子进程 } catch (err) { console.error(`下载失败: ${year}`, err); process.send(`下载失败: ${year}`); process.exit(1); // 异常退出 } }); // 将download改为返回Promise,确保等待流完成 function download(url, name, type, year) { return new Promise((resolve, reject) => { const filePath = path.join(down_path, `${name}_${type}_${year}.pdf`); const stream = fs.createWriteStream(filePath); https.get(url, (res) => { // 处理HTTP错误状态码 if (res.statusCode !== 200) { reject(new Error(`HTTP错误: ${res.statusCode}`)); stream.destroy(); return; } res.pipe(stream); stream.on('finish', () => { console.log(`下载完成: ${year}`); resolve(); }); stream.on('error', (err) => { reject(err); }); }).on('error', (err) => { reject(err); stream.destroy(); }); }); }
关键修改点
- 将
download函数改为返回Promise,监听流的finish和error事件,确保下载过程的异步等待。 - 下载完成或出错后,主动调用
process.exit()退出子进程,避免子进程挂起。 - 添加HTTP状态码检查和目录创建逻辑,增强鲁棒性。
内容的提问来源于stack exchange,提问作者DarkZeus
相关产品推荐
相关产品推荐

