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

JavaScript大数量新闻异步处理任务优化方案咨询

问题解答

当前代码在大规模数据下的稳定性结论

当前代码采用串行同步处理每个新闻条目,当articles数量达到100或1000时,无法稳定运行,核心问题包括:

  • 总处理时间线性暴涨:每个条目包含爬取、GPT调用、图片生成等多个耗时操作,串行执行会导致总耗时随条目数成倍增加,进程长时间占用期间一旦意外退出(服务器重启、网络波动),未完成任务全部丢失
  • 容错能力极差:单个条目处理失败(爬取超时、API调用报错)会直接中断整个循环,后续所有条目无法继续处理
  • 内存溢出风险:大量未释放的爬取内容、API返回数据累积在内存中,可能触发内存泄漏或溢出

优化方案

1. 引入任务队列(如Bull.js)实现异步解耦

将每个新闻条目的处理逻辑封装为独立任务,放入队列异步执行,是大规模场景下最可靠的优化方案,核心优势:

  • 任务持久化:进程重启后未完成的任务不会丢失,可从断点继续处理
  • 失败自动重试:可配置重试次数、间隔(如指数退避),避免单个任务失败影响全局
  • 并发可控:根据服务器性能和第三方API限流规则,设置合理并发数,防止请求过载被封禁
  • 可监控可追溯:直观查看任务进度、失败原因,便于排查问题

示例改造代码:

// 初始化Bull队列(依赖Redis)
const Queue = require('bull');
const newsProcessingQueue = new Queue('news-processing', 'redis://localhost:6379');

// 定义独立的任务处理函数
newsProcessingQueue.process(async (job) => {
  const article = job.data;
  try {
    const { title, sourceUrl, id, publishedAt } = article;
    // 爬取内容
    const markup = await scraper(sourceUrl);
    // GPT处理
    const data = await askGpt(markup);
    // DALL·E生成图片
    const generatedImageUrl = await generateImg(data?.imageDescription);
    // 图片上传S3
    const s3ImageUrl = await generateImgUrl(generatedImageUrl, title, id);
    // 推送至Strapi
    await createPost(
      data?.title,
      data?.abstract,
      data?.content,
      s3ImageUrl,
      publishedAt,
      data?.categories
    );
    return `Article ${id} processed successfully`;
  } catch (error) {
    // 抛出错误触发重试机制
    throw new Error(`Failed to process article ${id}: ${error.message}`);
  }
});

// 修改原函数:仅负责拉取数据并加入队列
const fetchAndProcessNews = async (queryString, from) => {
  // 可调整批量拉取的size,建议配合API分页
  const query = { queryString, from, size: 100 };
  try {
    const { articles } = await searchApi.getNews(query);
    if (articles?.length > 0) {
      // 批量添加任务到队列,配置重试规则
      await Promise.all(articles.map(article => 
        newsProcessingQueue.add(article, {
          attempts: 3, // 最多重试3次
          backoff: { type: 'exponential', delay: 1000 } // 指数退避重试(每次间隔翻倍)
        })
      ));
      console.log(`Added ${articles.length} articles to processing queue`);
    } else {
      console.log('No articles found');
    }
  } catch (error) {
    console.error('Error fetching news:', error.message);
  }
};

2. 并行处理+分批控制(轻量替代方案)

如果暂时不想引入队列,可通过分批并行处理控制并发数,避免一次性发起过多请求:

// 分批处理函数,控制每批并行数
const batchProcess = async (items, batchSize) => {
  const results = [];
  for (let i = 0; i < items.length; i += batchSize) {
    const batch = items.slice(i, i + batchSize);
    // 并行处理当前批次的条目
    const batchResults = await Promise.all(
      batch.map(async (article) => {
        try {
          // 单个条目的处理逻辑(同原流程)
          const { title, sourceUrl, id, publishedAt } = article;
          const markup = await scraper(sourceUrl);
          const data = await askGpt(markup);
          const generatedImageUrl = await generateImg(data?.imageDescription);
          const s3ImageUrl = await generateImgUrl(generatedImageUrl, title, id);
          await createPost(
            data?.title,
            data?.abstract,
            data?.content,
            s3ImageUrl,
            publishedAt,
            data?.categories
          );
          return { id, status: 'success' };
        } catch (error) {
          console.error(`Failed to process article ${article.id}:`, error.message);
          return { id, status: 'failed', error: error.message };
        }
      })
    );
    results.push(...batchResults);
    // 批次间添加短暂延迟,避免触发API限流
    await new Promise(resolve => setTimeout(resolve, 1000));
  }
  return results;
};

// 修改原函数调用分批处理
const fetchAndProcessNews = async (queryString, from) => {
  const query = { queryString, from, size: 100 };
  try {
    const { articles } = await searchApi.getNews(query);
    if (articles?.length > 0) {
      console.log('Processing news in batches...');
      // 每次并行处理5个条目,可根据服务器性能调整
      const results = await batchProcess(articles, 5);
      console.log('Processing completed:', results);
    } else {
      console.log('No articles found');
    }
  } catch (error) {
    console.error('Error fetching news:', error.message);
  }
};

3. 基础稳定性增强措施

  • 独立错误捕获:每个条目处理逻辑单独捕获错误,记录失败条目到日志/数据库,后续可手动重试
  • API超时控制:对第三方API调用添加超时(如Promise.race结合setTimeout),避免单个请求阻塞流程
  • 分页拉取源数据:如果API支持分页,不要一次性拉取1000条,而是分批(如每次50条)获取,降低单次请求的内存占用

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 16:33:12