首页 / NestJS 入门教程 / 微服务传输

NestJS 入门教程

微服务传输

本教程共 47 篇 · 第 40 篇 · 更新于 2026-08-09 · 约 15 分钟阅读

NestJS微服务RedisRabbitMQNATSKafkaMQTTgRPC

本节目标:掌握 NestJS 支持的六种传输层(Redis、RabbitMQ、NATS、Kafka、MQTT、gRPC),了解每种传输层的特点、配置方式和适用场景。

传输层概览

上一章我们学了微服务的基础知识。我们知道微服务之间通信靠的是传输层。NestJS 内置了六种传输层,除了默认的 TCP,还有 Redis、RabbitMQ、NATS、Kafka、MQTT 和 gRPC。

怎么选传输层?打个比方,传输层就像快递公司。TCP 是自建车队,Redis 是同城闪送,RabbitMQ 是顺丰快递,Kafka 是物流大仓,MQTT 是给物联网设备送件,gRPC 是国际快递。每种都有自己的特长。

传输层特点适用场景
TCP简单直接,无外部依赖开发测试、简单场景
Redis轻量级发布订阅缓存场景、简单消息
RabbitMQ功能丰富的消息队列企业级消息传递
NATS高性能、轻量级云原生、IoT
Kafka高吞吐、持久化大数据流处理
MQTT超低带宽、物联网IoT 设备通信
gRPC高性能 RPC、强类型跨语言微服务

Redis 传输层

Redis 传输层用的是发布订阅模式。消息发到频道里,谁订阅了谁就能收到。没人订阅的消息直接就没了,不保证送达。

安装

npm i --save ioredis

服务端配置

import { NestFactory } from '@nestjs/core';
import { Transport, MicroserviceOptions } from '@nestjs/microservices';
import { AppModule } from './app.module';

async function bootstrap() {
  const app = await NestFactory.createMicroservice<MicroserviceOptions>(
    AppModule,
    {
      transport: Transport.REDIS,
      options: {
        host: 'localhost',
        port: 6379,
      },
    },
  );
  await app.listen();
}
bootstrap();

Redis 传输层还支持 wildcards 选项,开启后可以用通配符订阅频道:

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.REDIS,
    options: {
      host: 'localhost',
      port: 6379,
      wildcards: true, // 开启通配符
    },
  },
);

开启后就能用 notifications.* 这样的模式来匹配频道了。

客户端配置

@Module({
  imports: [
    ClientsModule.register([
      {
        name: 'MATH_SERVICE',
        transport: Transport.REDIS,
        options: {
          host: 'localhost',
          port: 6379,
        },
      },
    ]),
  ],
})
export class AppModule {}

获取上下文

import { MessagePattern, Payload, Ctx, RedisContext } from '@nestjs/microservices';

@MessagePattern('notifications')
getNotifications(@Payload() data: number[], @Ctx() context: RedisContext) {
  console.log(`Channel: ${context.getChannel()}`);
}
Note

Redis 传输层适合轻量级场景。它不保证消息送达,如果消息丢失不能接受,就别用 Redis 做消息队列。

RabbitMQ 传输层

RabbitMQ 是最流行的消息队列之一。它支持消息确认、持久化、通配符路由等高级功能。如果你的系统对消息可靠性要求高,RabbitMQ 是个好选择。

安装

npm i --save amqplib amqp-connection-manager

服务端配置

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.RMQ,
    options: {
      urls: ['amqp://localhost:5672'],
      queue: 'cats_queue',
      queueOptions: {
        durable: false, // 队列是否持久化
      },
    },
  },
);

几个关键配置项:

配置项说明
urls连接地址数组
queue队列名称
noAck设为 false 开启手动确认
prefetchCount预取数量,控制并发
persistent消息是否持久化
wildcards是否启用通配符路由

手动确认消息

生产环境建议开启手动确认,确保消息处理完才删除:

// 配置
options: {
  urls: ['amqp://localhost:5672'],
  queue: 'cats_queue',
  noAck: false, // 关闭自动确认
}

// 处理器
@MessagePattern('notifications')
getNotifications(@Payload() data: number[], @Ctx() context: RmqContext) {
  const channel = context.getChannelRef();
  const originalMsg = context.getMessage();
  
  // 处理完业务逻辑后手动确认
  channel.ack(originalMsg);
}

发送带 headers 的消息

RmqRecordBuilder 可以给消息加 headers 和优先级:

import { RmqRecordBuilder } from '@nestjs/microservices';

const record = new RmqRecordBuilder(':cat:')
  .setOptions({
    headers: { 'x-version': '1.0.0' },
    priority: 3,
  })
  .build();

this.client.send('replace-emoji', record).subscribe(...);

服务端读取 headers:

@MessagePattern('replace-emoji')
replaceEmoji(@Payload() data: string, @Ctx() context: RmqContext): string {
  const { properties: { headers } } = context.getMessage();
  return headers['x-version'] === '1.0.0' ? '🐱' : '🐈';
}

通配符路由

开启 wildcards: true 后,可以用 #(匹配多个词)和 *(匹配一个词)做路由:

// 匹配 cats、cats.meow、cats.meow.purr
@MessagePattern('cats.#')
getCats(@Payload() data: { message: string }) {
  return { message: 'Hello from cats service!' };
}
Tip

RabbitMQ 的 durable: true 配合 persistent: true 可以保证 RabbitMQ 重启后消息不丢。做金融、订单类系统一定要开。

NATS 传输层

NATS 是一个高性能的消息系统,用 Go 写的,非常轻量。它支持通配符订阅和分布式队列(负载均衡)。

安装

npm i --save nats

服务端配置

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.NATS,
    options: {
      servers: ['nats://localhost:4222'],
    },
  },
);

分布式队列

NATS 内置了负载均衡功能,配置 queue 就行:

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.NATS,
    options: {
      servers: ['nats://localhost:4222'],
      queue: 'cats_queue', // 同名队列的消费者会负载均衡
    },
  },
);

通配符订阅

@MessagePattern('time.us.*')
getDate(@Payload() data: number[], @Ctx() context: NatsContext) {
  console.log(`Subject: ${context.getSubject()}`); // 比如 "time.us.east"
  return new Date().toLocaleTimeString();
}

发送带 headers 的消息

NatsRecordBuilder

import * as nats from 'nats';
import { NatsRecordBuilder } from '@nestjs/microservices';

const headers = nats.headers();
headers.set('x-version', '1.0.0');

const record = new NatsRecordBuilder(':cat:').setHeaders(headers).build();
this.client.send('replace-emoji', record).subscribe(...);

服务端读取:

@MessagePattern('replace-emoji')
replaceEmoji(@Payload() data: string, @Ctx() context: NatsContext): string {
  const headers = context.getHeaders();
  return headers['x-version'] === '1.0.0' ? '🐱' : '🐈';
}
Tip

NATS 的优雅关闭值得一提。配置 gracefulShutdown: true 后,服务关闭时会先取消订阅再断开连接,避免消息丢失:

options: {
  servers: ['nats://localhost:4222'],
  gracefulShutdown: true,
  gracePeriod: 10000, // 等待 10 秒
}

Kafka 传输层

Kafka 是大数据领域的王者。高吞吐、消息持久化、支持回放,非常适合事件流处理。

安装

npm i --save kafkajs

服务端配置

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.KAFKA,
    options: {
      client: {
        brokers: ['localhost:9092'],
      },
      consumer: {
        groupId: 'my-consumer',
      },
    },
  },
);

Kafka 的配置分几块:client(连接配置)、consumer(消费者配置)、run(运行配置)、subscribe(订阅配置)、producer(生产者配置)。

消息处理器

@Controller()
export class HeroesController {
  @MessagePattern('hero.kill.dragon')
  killDragon(@Payload() message: KillDragonMessage): any {
    const dragonId = message.dragonId;
    return [{ id: 1, name: 'Mythical Sword' }];
  }
}

Kafka 还可以发送带 key 和 headers 的消息:

@MessagePattern('hero.kill.dragon')
killDragon(@Payload() message: KillDragonMessage) {
  return {
    key: message.heroId, // 消息 key,用于分区
    headers: { realm: 'Nest' },
    value: [{ id: 1, name: 'Mythical Sword' }],
  };
}

请求-响应的特殊处理

Kafka 用两个 topic 来实现请求-响应:一个发请求,一个收回复。客户端需要先订阅回复 topic:

@Injectable()
export class HeroesService implements OnModuleInit {
  constructor(@Inject('HERO_SERVICE') private client: ClientKafkaProxy) {}

  onModuleInit() {
    this.client.subscribeToResponseOf('hero.kill.dragon');
  }
}
Note

如果你只用事件模式(@EventPattern + emit),不需要 subscribeToResponseOf。只有用请求-响应模式才需要。

手动提交 offset

@EventPattern('user.created')
async handleUserCreated(@Payload() data: any, @Ctx() context: KafkaContext) {
  // 处理业务逻辑...

  const { offset } = context.getMessage();
  const partition = context.getPartition();
  const topic = context.getTopic();
  const consumer = context.getConsumer();
  await consumer.commitOffsets([{ topic, partition, offset }]);
}

关闭自动提交:

options: {
  client: { brokers: ['localhost:9092'] },
  run: { autoCommit: false },
}

慢处理场景的 heartbeat

如果处理一条消息很慢,需要定期发心跳防止会话超时:

@MessagePattern('hero.kill.dragon')
async killDragon(@Payload() message: any, @Ctx() context: KafkaContext) {
  const heartbeat = context.getHeartbeat();

  await doWorkPart1();
  await heartbeat(); // 告诉 Kafka 我还活着
  await doWorkPart2();
}
Tip

Kafka 的命名约定:服务端会自动给 clientIdgroupId-server 后缀,客户端加 -client 后缀,避免冲突。

MQTT 传输层

MQTT 是物联网领域的标准协议。超低带宽、低延迟,适合设备间通信。

安装

npm i --save mqtt

服务端配置

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.MQTT,
    options: {
      url: 'mqtt://localhost:1883',
    },
  },
);

通配符

MQTT 支持两种通配符:+(单层)和 #(多层):

// 匹配 sensors/room1/temperature/1、sensors/room2/temperature/2 等
@MessagePattern('sensors/+/temperature/+')
getTemperature(@Ctx() context: MqttContext) {
  console.log(`Topic: ${context.getTopic()}`);
}

QoS(服务质量)

MQTT 有三个级别的 QoS:

  • QoS 0:最多一次,可能丢失
  • QoS 1:至少一次,可能重复
  • QoS 2:恰好一次,最安全

全局设置:

options: {
  url: 'mqtt://localhost:1883',
  subscribeOptions: { qos: 2 },
}

也可以按模式单独设置:

@EventPattern('critical-events', { extras: { qos: 2 } })
handleCriticalEvent(@Payload() data: any) {
  // QoS 2,保证恰好一次
}

@EventPattern('metrics', { extras: { qos: 0 } })
handleMetrics(@Payload() data: any) {
  // QoS 0,允许丢失
}

发送带属性的消息

MqttRecordBuilder

import { MqttRecordBuilder } from '@nestjs/microservices';

const record = new MqttRecordBuilder(':cat:')
  .setProperties({ userProperties: { 'x-version': '1.0.0' } })
  .setQoS(1)
  .build();

client.send('replace-emoji', record).subscribe(...);

gRPC 传输层

gRPC 跟前面几个都不一样。它用的是 Protocol Buffers 做序列化,性能极高,还支持跨语言调用。如果你的微服务用不同语言写的,gRPC 是最佳选择。

安装

npm i --save @grpc/grpc-js @grpc/proto-loader

定义 proto 文件

gRPC 的第一步是写 .proto 文件,定义服务接口:

// hero/hero.proto
syntax = "proto3";

package hero;

service HeroesService {
  rpc FindOne (HeroById) returns (Hero) {}
}

message HeroById {
  int32 id = 1;
}

message Hero {
  int32 id = 1;
  string name = 2;
}

这个文件定义了:一个叫 HeroesService 的服务,有个 FindOne 方法,接收 HeroById,返回 Hero

Tip

nest-cli.json 里配置 assets.proto 文件自动复制到 dist 目录:

{
  "compilerOptions": {
    "assets": ["**/*.proto"],
    "watchAssets": true
  }
}

服务端配置

import { join } from 'path';

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.GRPC,
    options: {
      package: 'hero',
      protoPath: join(__dirname, 'hero/hero.proto'),
    },
  },
);

实现服务

@GrpcMethod() 装饰器来实现 proto 文件里定义的方法:

import { GrpcMethod } from '@nestjs/microservices';

@Controller()
export class HeroesController {
  @GrpcMethod('HeroesService', 'FindOne')
  findOne(data: HeroById): Hero {
    const items = [
      { id: 1, name: 'John' },
      { id: 2, name: 'Doe' },
    ];
    return items.find(({ id }) => id === data.id);
  }
}

@GrpcMethod() 的两个参数分别对应 proto 文件里的服务名和方法名。如果你省略第二个参数,NestJS 会根据方法名自动匹配(把 findOne 转成 FindOne)。如果两个参数都省略,还会根据类名匹配服务名。

客户端调用

gRPC 客户端用的是 ClientGrpc,不是 ClientProxy

imports: [
  ClientsModule.register([
    {
      name: 'HERO_PACKAGE',
      transport: Transport.GRPC,
      options: {
        package: 'hero',
        protoPath: join(__dirname, 'hero/hero.proto'),
      },
    },
  ]),
]

使用时通过 getService() 获取具体服务实例:

@Injectable()
export class AppService implements OnModuleInit {
  private heroesService: HeroesService;

  constructor(@Inject('HERO_PACKAGE') private client: ClientGrpc) {}

  onModuleInit() {
    this.heroesService = this.client.getService<HeroesService>('HeroesService');
  }

  getHero(): Observable<any> {
    return this.heroesService.findOne({ id: 1 });
  }
}
Note

gRPC 客户端的方法名是小驼峰格式。proto 文件里的 FindOne,在客户端调用时要用 findOne

gRPC 流式通信

gRPC 支持流式通信,适合聊天、实时数据推送等场景。

在 proto 文件里定义流式方法:

service HelloService {
  rpc BidiHello(stream HelloRequest) returns (stream HelloResponse);
  rpc LotsOfGreetings(stream HelloRequest) returns (HelloResponse);
}

@GrpcStreamMethod() 处理流式请求:

import { GrpcStreamMethod } from '@nestjs/microservices';

@GrpcStreamMethod()
bidiHello(messages: Observable<any>): Observable<any> {
  const subject = new Subject();

  messages.subscribe({
    next: (message) => {
      console.log(message);
      subject.next({ reply: 'Hello, world!' });
    },
    complete: () => subject.complete(),
  });

  return subject.asObservable();
}

也可以用 @GrpcStreamCall() 直接操作原生流:

@GrpcStreamCall()
bidiHello(requestStream: any) {
  requestStream.on('data', (message) => {
    console.log(message);
    requestStream.write({ reply: 'Hello, world!' });
  });
}

gRPC 元数据

gRPC 的 metadata 类似于 HTTP 的 headers,用来传递调用信息(认证 token、请求 ID 等):

@GrpcMethod('HeroesService', 'FindOne')
findOne(data: HeroById, metadata: Metadata, call: ServerUnaryCall<any, any>): Hero {
  // 读取客户端发来的 metadata
  const token = metadata.get('authorization');

  // 发送 metadata 给客户端
  const serverMetadata = new Metadata();
  serverMetadata.add('Set-Cookie', 'yummy_cookie=choco');
  call.sendMetadata(serverMetadata);

  return items.find(({ id }) => id === data.id);
}

gRPC Reflection

开发阶段可以开启 gRPC Reflection,让客户端自动发现服务接口(类似 REST 的 OpenAPI 文档):

npm i --save @grpc/reflection
import { ReflectionService } from '@grpc/reflection';

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    options: {
      onLoadPackageDefinition: (pkg, server) => {
        new ReflectionService(pkg).addToServer(server);
      },
    },
  },
);

传输层对比总结

特性RedisRabbitMQNATSKafkaMQTTgRPC
消息持久化可选
消息确认
通配符
负载均衡
流式通信
跨语言
吞吐量极高
Tip

选择建议:

  • 开发阶段或简单场景:用 TCP 就够了
  • 需要消息可靠传递:RabbitMQ
  • 大数据流处理、事件溯源:Kafka
  • 高性能、云原生:NATS
  • 物联网设备:MQTT
  • 跨语言微服务调用:gRPC
  • 已有 Redis 基础设施、不要求消息可靠:Redis

小结

本章详细介绍了 NestJS 支持的六种传输层。虽然传输层不同,但 NestJS 帮你屏蔽了底层差异。@MessagePattern()@EventPattern()ClientProxy 这些 API 在所有传输层上都是一样的用法。切换传输层只需要改配置,业务代码不用动。

这就是 NestJS 微服务的优雅之处:换传输层就像换快递公司,你的包裹(业务逻辑)完全不受影响。