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

NestJS(Bull)批量创建帖子最佳实践及扣费逻辑位置咨询

问题背景

我正在用NestJS、PostgreSQL和Prisma开发一款类Telegram的社交应用,现在有个场景需要批量创建数千条帖子,为了避免性能问题或崩溃,打算用队列方案,这是我第一次实现队列功能。

当前的代码实现:

批量添加任务的代码

for (const phone of inputPhones) {
  const params: CreatePostQueueDto = {
    channelId,
    title,
    content,
    phone,
    postGroupId: postGroup.id,
    postType,
    media,
    creatorId,
  };
  void this.QueueService.createPostJob(params);
}

队列服务代码

async createPostJob(data: CreatePostQueueDto) {
  await this.mainQueue.add('create_post', data);
}

队列处理器代码

@Processor('main')
@Injectable()
export class QueueProcessor {

  @Process('create_audience')

  @Process('create_post')
  async handleCreatePost(job: Job<CreatePostQueueDto>) {
    try {
      // create post
    } catch (error) {
      console.log(error);
      throw new BadRequestException(`Failed to process create_post: ${error}`);
    }
  }
}

我的疑问

  1. 在NestJS中用Bull批量创建大量帖子的最佳实践是什么?现有代码需要调整哪些地方来提升可靠性、性能或错误处理能力?
  2. 用户扣费逻辑(比如扣积分/余额)应该放在哪里?要求根据所有成功创建的帖子数量一次性扣费(部分帖子可能创建失败,得等全部任务完成才知道成功数),是在加队列任务前处理,还是在队列处理器里处理?

回答

问题1:Bull批量创建帖子的最佳实践与代码优化

1. 批量任务添加优化

你现在循环逐个调用createPostJob,每次单独和Redis通信,数千条数据会产生大量Redis请求,效率很低。应该改用Bull的批量添加任务方法,减少Redis交互次数:

// 先整理所有任务参数为数组
const jobDataList = inputPhones.map(phone => ({
  channelId,
  title,
  content,
  phone,
  postGroupId: postGroup.id,
  postType,
  media,
  creatorId,
}));
// 批量添加任务
await this.mainQueue.addBulk(jobDataList.map(data => ({
  name: 'create_post',
  data,
})));

2. 提升任务可靠性

  • 配置重试策略:创建帖子失败可能是临时问题(比如数据库连接波动),给任务设置重试次数和指数退避间隔,避免直接丢弃失败任务:
    // 批量添加时统一配置(单个任务添加也可用)
    await this.mainQueue.addBulk(jobDataList.map(data => ({
      name: 'create_post',
      data,
      opts: {
        attempts: 3, // 最多重试3次
        backoff: {
          type: 'exponential', // 指数退避,避免短时间重复冲击数据库
          delay: 1000, // 第一次重试间隔1秒,之后间隔翻倍
        },
      },
    })));
    
  • 设置任务超时:防止单个任务卡住占用资源,比如设置10秒超时:
    opts: {
      timeout: 10000, // 10秒后任务超时,标记为失败
    }
    

3. 错误处理改进

  • 队列处理器不需要抛出HTTP异常(比如BadRequestException),直接抛出错误让Bull按重试策略处理,同时记录详细错误日志:
    async handleCreatePost(job: Job<CreatePostQueueDto>) {
      try {
        // 执行Prisma创建帖子逻辑
        await this.prisma.post.create({
          data: {
            // 映射job.data到数据库字段
          },
        });
        this.logger.log(`Post created for phone: ${job.data.phone}`);
      } catch (error) {
        this.logger.error(`Failed to create post for phone ${job.data.phone}:`, error.stack);
        // 抛出错误,让Bull处理重试
        throw error;
      }
    }
    
  • 监听队列失败事件,统一处理最终失败的任务(比如存入失败表、通知管理员):
    // 在QueueService初始化时监听
    this.mainQueue.on('failed', async (job, err) => {
      this.logger.error(`Job ${job.id} failed permanently:`, err.stack);
      // 把失败任务数据存入数据库,方便后续排查和重试
      await this.prisma.failedJob.create({
        data: {
          jobId: job.id,
          type: 'create_post',
          data: JSON.stringify(job.data),
          error: err.message,
        },
      });
    });
    

4. 性能优化

  • 设置处理器并发数:根据数据库承载能力调整并发数,避免同时创建太多帖子导致数据库压力过大:
    @Processor('main', {
      concurrency: 10, // 同时处理10个任务,可根据实际情况调整
    })
    
  • 数据库层面优化:每个任务内的创建操作尽量高效,比如避免不必要的关联查询;如果同批次帖子有共同的关联数据,可提前查询好传入任务,减少重复查询。

问题2:用户扣费逻辑的位置

这种需要统计全部任务成功数量后一次性扣费的场景,既不能在任务添加前处理(不知道最终成功数),也不能在单个任务处理器里处理(无法统计总数),推荐两种方案:

方案1:使用Bull的批次任务功能

Bull支持将一批任务归为同一个批次,监听批次完成事件,统计成功数量后执行扣费:

// 1. 创建批次
const batchId = await this.mainQueue.createBatch();
// 2. 批量添加任务到该批次
await this.mainQueue.addBulk(jobDataList.map(data => ({
  name: 'create_post',
  data,
  batchId,
})));
// 3. 监听批次完成事件
this.mainQueue.on(`batch:completed:${batchId}`, async (batch) => {
  // 获取批次内所有成功完成的任务
  const successJobs = await this.mainQueue.getJobs(['completed'], 0, -1, true, batchId);
  const successCount = successJobs.length;
  
  if (successCount > 0) {
    // 执行扣费逻辑,比如按成功数量扣除积分
    await this.prisma.user.update({
      where: { id: creatorId },
      data: {
        points: { decrement: successCount * POINTS_PER_POST },
      },
    });
  }

  // 可选:处理失败任务的通知
  const failedJobs = await this.mainQueue.getJobs(['failed'], 0, -1, true, batchId);
  if (failedJobs.length > 0) {
    // 通知用户有多少帖子创建失败
  }
});

方案2:手动记录批次状态到数据库

如果不想用Bull的批次功能,可以先创建批次记录,每个任务完成后更新状态,最后触发扣费:

  1. 创建批次记录:
const postBatch = await this.prisma.postBatch.create({
  data: {
    creatorId,
    totalTasks: inputPhones.length,
    successCount: 0,
    failedCount: 0,
    status: 'PROCESSING',
  },
});
  1. 在任务处理器中更新批次状态并触发扣费:
async handleCreatePost(job: Job<CreatePostQueueDto>) {
  try {
    await this.prisma.post.create({/* ... */});
    // 原子更新成功数
    await this.prisma.postBatch.update({
      where: { id: job.data.postBatchId }, // 需将postBatchId传入任务数据
      data: { successCount: { increment: 1 } },
    });
  } catch (error) {
    // 原子更新失败数
    await this.prisma.postBatch.update({
      where: { id: job.data.postBatchId },
      data: { failedCount: { increment: 1 } },
    });
    throw error;
  } finally {
    // 检查是否所有任务已完成
    const batch = await this.prisma.postBatch.findUnique({
      where: { id: job.data.postBatchId },
    });
    if (batch.successCount + batch.failedCount === batch.totalTasks) {
      await this.prisma.postBatch.update({
        where: { id: job.data.postBatchId },
        data: { status: 'COMPLETED' },
      });
      // 执行扣费
      if (batch.successCount > 0) {
        await this.prisma.user.update({
          where: { id: batch.creatorId },
          data: { points: { decrement: batch.successCount * POINTS_PER_POST } },
        });
      }
    }
  }
}

注:Prisma的increment是原子操作,不会出现并发更新问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 14:14:57