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}`); } } }
我的疑问
- 在NestJS中用Bull批量创建大量帖子的最佳实践是什么?现有代码需要调整哪些地方来提升可靠性、性能或错误处理能力?
- 用户扣费逻辑(比如扣积分/余额)应该放在哪里?要求根据所有成功创建的帖子数量一次性扣费(部分帖子可能创建失败,得等全部任务完成才知道成功数),是在加队列任务前处理,还是在队列处理器里处理?
回答
问题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的批次功能,可以先创建批次记录,每个任务完成后更新状态,最后触发扣费:
- 创建批次记录:
const postBatch = await this.prisma.postBatch.create({ data: { creatorId, totalTasks: inputPhones.length, successCount: 0, failedCount: 0, status: 'PROCESSING', }, });
- 在任务处理器中更新批次状态并触发扣费:
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
相关产品推荐
相关产品推荐

