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

NodeJS中使用fork时等待所有子进程完成的实现方案

问题

我编写了一个Node.js脚本,用于遍历URL列表下载PDF文件,每个URL对应5个PDF,因此使用fork创建子进程并行下载这些PDF。目前需要解决两个问题:

  1. 仅当当前URL对应的所有子进程(下载任务)全部执行完成后,才继续处理下一个URL。
  2. 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('所有任务完成');
})();

关键修改点

  1. 新增spawnDownloader函数,将子进程的创建、消息监听和退出逻辑封装为Promise,方便批量等待。
  2. 修正原CSV读取循环的错误(l+2改为l += 2),并使用块级变量避免全局污染。
  3. 将每个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();
        });
    });
}

关键修改点

  1. 将download函数改为返回Promise,监听流的finish和error事件,确保下载过程的异步等待。
  2. 下载完成或出错后,主动调用process.exit()退出子进程,避免子进程挂起。
  3. 添加HTTP状态码检查和目录创建逻辑,增强鲁棒性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:01:07