如何让RabbitMQ消费者监听器每次执行结果都被控制器获取处理?
问题
我有一个消费RabbitMQ队列的函数,希望将其每次执行的结果保存下来,以运行其他代码(例如存入数据库)。目前的问题是,结果仅能被保存一次:监听器可以正常运行并持续获取队列中的事件信息,函数内部代码也会随着新事件加入队列而执行,但无法将每次执行的结果重新赋值给变量。
现有代码
调用消费者的控制器
async run() { const eventData = await this.eventManager.consume(QueuesToConsume.USER_CREATED) await this.createUserUseCase.run(eventData); }
RabbitMQ消费者
async consume(queue: string): Promise<DomainEvent> { let eventData: DomainEvent; return new Promise<DomainEvent>(async (resolve, reject) => { await this.channel.consume(queue, async (msg: Message) => { console.log(`Message: \n ${Buffer.from(msg.content)} \n received successfully!`) await this.channel.ack(msg) eventData = JSON.parse(Buffer.from(msg.content).toString('utf8')) console.log('Message acknowledged successfully') resolve(eventData); }).catch(err => { console.log(`Error consuming the message: \n ${err}`) reject(err) }); }) }
当前代码无法正常工作:控制器中的eventData无法获取每一次响应,createUserUseCase仅能执行一次。请问如何修改才能让eventData获取消费者返回的每一次结果?
解决方案
问题核心在于当前的consume函数只返回单个Promise,而Promise一旦resolve就会完成,无法多次返回结果。要实现持续消费并处理每一条消息,需要调整架构:
1. 修改消费者方法,接受回调处理每条消息
将consume改为接收回调函数的形式,让每条消息触发一次业务逻辑:
async consume(queue: string, messageHandler: (event: DomainEvent) => Promise<void>): Promise<void> { await this.channel.consume(queue, async (msg: Message) => { if (!msg) return; console.log(`Message: \n ${Buffer.from(msg.content)} \n received successfully!`) try { const eventData = JSON.parse(Buffer.from(msg.content).toString('utf8')); // 调用外部传入的处理逻辑 await messageHandler(eventData); // 业务处理成功后再确认消息 await this.channel.ack(msg); console.log('Message acknowledged successfully'); } catch (err) { console.log(`Error processing message: \n ${err}`); // 处理失败时拒绝消息,可根据需求配置是否重新入队 await this.channel.nack(msg, false, false); } }).catch(err => { console.log(`Error consuming the queue: \n ${err}`); throw err; }); }
2. 调整控制器调用逻辑
在控制器中传入处理消息的回调,确保每条消息都会触发createUserUseCase.run:
async run() { await this.eventManager.consume(QueuesToConsume.USER_CREATED, async (eventData) => { await this.createUserUseCase.run(eventData); }); console.log(`Started consuming queue: ${QueuesToConsume.USER_CREATED}`); }
关键调整说明
- 移除单个Promise返回逻辑,改用回调实现多消息处理,确保每条消息都能触发业务代码
- 调整消息确认时机:业务逻辑执行成功后再ack,避免未完成处理就确认消息导致数据丢失
- 添加异常捕获,处理消息解析或业务执行失败的场景,防止单个消息异常导致消费者崩溃
内容的提问来源于stack exchange,提问作者Roger González Hermosa
相关产品推荐
相关产品推荐

