使用@nestjs-plus/rabbitmq构建Pub/Sub,需实现订阅者向发布者发送确认
Hey there! Since you already have a working Pub/Sub setup with @nestjs-plus/rabbitmq, adding a publisher acknowledgment flow is totally doable by leveraging RabbitMQ's built-in reply queue and correlation ID features. Let's break this down step by step.
Core Idea
We'll use a dedicated reply queue for sending status confirmations from subscribers back to publishers. Each published message will carry a unique correlationId to match the original message with its corresponding acknowledgment, ensuring the publisher knows exactly which request the confirmation refers to.
Step 1: Update the Publisher to Listen for Acknowledgments
First, modify your publisher service to generate unique correlation IDs for each message, declare a reply queue, and listen for incoming confirmations from subscribers.
import { Injectable } from '@nestjs/common'; import { AmqpConnection } from '@nestjs-plus/rabbitmq'; import { v4 as uuidv4 } from 'uuid'; @Injectable() export class PublisherService { // Map to track pending confirmations using correlation IDs private pendingConfirmations = new Map<string, (ack: any) => void>(); constructor(private readonly amqpConnection: AmqpConnection) { // Initialize the reply queue listener on service startup this.setupReplyQueue(); } private async setupReplyQueue() { // Declare a durable reply queue (adjust durability based on your needs) const { queue } = await this.amqpConnection.channel.assertQueue('publisher-replies', { durable: false, }); // Listen for messages on the reply queue await this.amqpConnection.channel.consume(queue, (msg) => { if (!msg?.properties?.correlationId) return; const correlationId = msg.properties.correlationId; const resolveFn = this.pendingConfirmations.get(correlationId); if (resolveFn) { // Parse and resolve the acknowledgment to the waiting publish promise const ackData = JSON.parse(msg.content.toString()); resolveFn(ackData); this.pendingConfirmations.delete(correlationId); // Acknowledge the confirmation message to RabbitMQ this.amqpConnection.channel.ack(msg); } }); } async publishMessage(exchange: string, routingKey: string, payload: any): Promise<any> { const correlationId = uuidv4(); // Return a promise that resolves when we receive the acknowledgment return new Promise((resolve) => { this.pendingConfirmations.set(correlationId, resolve); // Publish the message with reply queue and correlation ID metadata this.amqpConnection.publish(exchange, routingKey, payload, { replyTo: 'publisher-replies', correlationId: correlationId, }); }); } }
Step 2: Modify the Subscriber to Send Acknowledgments
Update your subscriber to extract the reply queue and correlation ID from incoming messages, then send a confirmation (success or failure) back to the publisher's reply queue.
import { Injectable } from '@nestjs/common'; import { RabbitSubscribe } from '@nestjs-plus/rabbitmq'; import { AmqpConnection } from '@nestjs-plus/rabbitmq'; @Injectable() export class SubscriberService { constructor(private readonly amqpConnection: AmqpConnection) {} @RabbitSubscribe({ exchange: 'your-target-exchange', routingKey: 'your-routing-key', queue: 'subscriber-processing-queue', }) async handleIncomingMessage(msg: any, context: any) { try { // Execute your core business logic here console.log('Processing message:', msg); // Extract reply queue and correlation ID from message properties const { replyTo, correlationId } = context.fields.properties; if (replyTo && correlationId) { // Build a success acknowledgment const successAck = { status: 'processed', timestamp: new Date().toISOString(), message: 'Message handled successfully', originalCorrelationId: correlationId, }; // Send the acknowledgment to the publisher's reply queue await this.amqpConnection.channel.sendToQueue( replyTo, Buffer.from(JSON.stringify(successAck)), { correlationId: correlationId } ); } // Acknowledge the original message to RabbitMQ (marks it as consumed) context.channel.ack(context.message); } catch (error) { console.error('Failed to process message:', error); const { replyTo, correlationId } = context.fields.properties; if (replyTo && correlationId) { // Build a failure acknowledgment const failureAck = { status: 'failed', timestamp: new Date().toISOString(), error: error.message, originalCorrelationId: correlationId, }; // Send failure status back to publisher await this.amqpConnection.channel.sendToQueue( replyTo, Buffer.from(JSON.stringify(failureAck)), { correlationId: correlationId } ); } // Reject the original message (adjust requeue flag based on your retry policy) context.channel.nack(context.message, false, false); } } }
Key Notes to Remember
- Correlation IDs: These are critical for matching incoming confirmations to their original published messages. Never skip using them—they prevent the publisher from mixing up acknowledgments.
- Reply Queue Durability: If you need confirmations to survive RabbitMQ restarts, set
durable: truewhen declaring the reply queue. For temporary, short-lived confirmations, usedurable: false. - RabbitMQ ACK vs. Business Acknowledgment:
context.channel.ack()tells RabbitMQ the message was consumed, while the reply queue message is a business-level confirmation for the publisher. They serve different purposes!
内容的提问来源于stack exchange,提问作者Shubham Thakare

