AMQPClient类RPC调用复用首个correlationId问题排查
AMQPClient RPC调用重复使用旧correlationId问题
我正在开发AMQPClient类来封装RPC调用,首次调用完全正常,但第二次调用(在第一次完成后发起)时,onMessage函数里始终复用第一次生成的correlationId,新生成的ID根本没生效。试过把const correlationId = Math.random().toString(36).slice(2)移到Promise外面、用匿名函数包裹onMessage回调、把correlationId当参数传递等方法,都没解决问题。
核心RPC方法代码
async RPC<T>(queue: string, message: string): Promise<T> { if (!this.channel) { throw new Error('Channel not initialized') } const replyTo = `${queue}.reply` await this.channel.assertQueue(replyTo) await this.channel.assertQueue(queue) return new Promise<T>((resolve) => { const correlationId = Math.random().toString(36).slice(2) console.log('generated correlationId: ', correlationId) const onMessage = (message: ConsumeMessage | null) => { console.log( correlationId, message?.properties.correlationId, correlationId === message?.properties.correlationId ) if (message && message.properties.correlationId === correlationId) { resolve(JSON.parse(message.content.toString())) this.channel?.removeListener('message', onMessage) } } this.channel?.consume(replyTo, onMessage, { noAck: true }) this.channel?.sendToQueue(queue, Buffer.from(message), { correlationId, replyTo }) }) }
运行输出
generated correlationId: lwfvgqym5ya lwfvgqym5ya lwfvgqym5ya true generated correlationId: 1m09k9jk2xm lwfvgqym5ya 1m09k9jk2xm false
完整代码
import client, { Channel, Connection, ConsumeMessage } from 'amqplib' class AMQPClient { private channel?: Channel constructor(private readonly amqpUrl: string) { client.connect(this.amqpUrl).then((connection) => { connection.createChannel().then((channel) => { this.channel = channel }) process.on('SIGINT', () => this.close(connection)) process.on('SIGTERM', () => this.close(connection)) }) } async RPC<T>(queue: string, message: string): Promise<T> { if (!this.channel) { throw new Error('Channel not initialized') } const replyTo = `${queue}.reply` await this.channel.assertQueue(replyTo) await this.channel.assertQueue(queue) return new Promise<T>((resolve) => { const correlationId = Math.random().toString(36).slice(2) console.log('generated correlationId: ', correlationId) const onMessage = (message: ConsumeMessage | null) => { console.log( correlationId, message?.properties.correlationId, correlationId === message?.properties.correlationId ) if (message && message.properties.correlationId === correlationId) { resolve(JSON.parse(message.content.toString())) this.channel?.removeListener('message', onMessage) } } this.channel?.consume(replyTo, (msg) => onMessage(msg), { noAck: true }) this.channel?.sendToQueue(queue, Buffer.from(message), { correlationId, replyTo }) }) } close(connection: Connection) { connection.close() process.exit(0) } } const amqpClient = new AMQPClient(process.env.AMQP_URL || 'amqp://localhost') export { amqpClient, AMQPClient }
调用代码
this.amqpClient.RPC<MerchantStatus>( 'getMerchantStatus', JSON.stringify({ merchantId: 'test' }) )
问题原因与解决方案
问题根源
你用了this.channel?.removeListener('message', onMessage)来移除监听器,但这根本没用——amqplib的consume方法添加的消费者,不是通过Node.js EventEmitter的on('message')绑定的,所以removeListener找不到对应的监听器,导致第一次调用的消费者一直留在队列上。第二次调用时,新消息会触发所有已存在的消费者,第一个消费者拿着旧的correlationId,就出现了日志里的不匹配情况。
修复代码
把移除监听器的方式改成用consume返回的consumerTag调用channel.cancel(),这才是amqplib官方推荐的移除消费者的方法:
async RPC<T>(queue: string, message: string): Promise<T> { if (!this.channel) { throw new Error('Channel not initialized') } const replyTo = `${queue}.reply` await this.channel.assertQueue(replyTo) await this.channel.assertQueue(queue) return new Promise<T>(async (resolve) => { const correlationId = Math.random().toString(36).slice(2) console.log('generated correlationId: ', correlationId) const onMessage = async (message: ConsumeMessage | null) => { if (message && message.properties.correlationId === correlationId) { console.log( correlationId, message.properties.correlationId, correlationId === message.properties.correlationId ) resolve(JSON.parse(message.content.toString())) // 用consumerTag取消对应的消费者 if (consumerTag) { await this.channel?.cancel(consumerTag) } } } // 保存consume返回的consumerTag const consumeResult = await this.channel.consume(replyTo, onMessage, { noAck: true }) const consumerTag = consumeResult?.consumerTag this.channel.sendToQueue(queue, Buffer.from(message), { correlationId, replyTo }) }) }
额外优化建议
- 可以给reply队列设置
exclusive: true或者autoDelete: true,避免队列长期存在堆积消息:await this.channel.assertQueue(replyTo, { exclusive: true, autoDelete: true }) - 给Promise添加reject逻辑,处理超时或者异常情况,避免Promise一直pending。
内容的提问来源于stack exchange,提问作者Lincon Dias
相关产品推荐
相关产品推荐

