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

使用@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.

Implementing Publisher Acknowledgments in Your NestJS RabbitMQ Setup

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: true when declaring the reply queue. For temporary, short-lived confirmations, use durable: 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:01:14