微服务传输
本教程共 47 篇 · 第 40 篇 · 更新于 2026-08-09 · 约 15 分钟阅读
本节目标:掌握 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()}`);
}
NoteRedis 传输层适合轻量级场景。它不保证消息送达,如果消息丢失不能接受,就别用 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!' };
}
TipRabbitMQ 的
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' ? '🐱' : '🐈';
}
TipNATS 的优雅关闭值得一提。配置
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();
}
TipKafka 的命名约定:服务端会自动给
clientId和groupId加-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 });
}
}
NotegRPC 客户端的方法名是小驼峰格式。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);
},
},
},
);
传输层对比总结
| 特性 | Redis | RabbitMQ | NATS | Kafka | MQTT | gRPC |
|---|---|---|---|---|---|---|
| 消息持久化 | 否 | 是 | 可选 | 是 | 否 | 否 |
| 消息确认 | 否 | 是 | 否 | 是 | 否 | 否 |
| 通配符 | 是 | 是 | 是 | 否 | 是 | 否 |
| 负载均衡 | 否 | 是 | 是 | 是 | 否 | 否 |
| 流式通信 | 否 | 否 | 否 | 是 | 否 | 是 |
| 跨语言 | 否 | 否 | 否 | 否 | 否 | 是 |
| 吞吐量 | 中 | 中 | 高 | 极高 | 低 | 高 |
Tip选择建议:
- 开发阶段或简单场景:用 TCP 就够了
- 需要消息可靠传递:RabbitMQ
- 大数据流处理、事件溯源:Kafka
- 高性能、云原生:NATS
- 物联网设备:MQTT
- 跨语言微服务调用:gRPC
- 已有 Redis 基础设施、不要求消息可靠:Redis
小结
本章详细介绍了 NestJS 支持的六种传输层。虽然传输层不同,但 NestJS 帮你屏蔽了底层差异。@MessagePattern()、@EventPattern()、ClientProxy 这些 API 在所有传输层上都是一样的用法。切换传输层只需要改配置,业务代码不用动。
这就是 NestJS 微服务的优雅之处:换传输层就像换快递公司,你的包裹(业务逻辑)完全不受影响。