首页 / NestJS 入门教程 / 微服务基础

NestJS 入门教程

微服务基础

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

NestJS微服务ClientProxyMessagePatternEventPattern传输层TCP

本节目标:搞懂 NestJS 微服务的核心概念,学会创建微服务、定义消息处理器、使用客户端通信,以及处理超时和异常。

什么是微服务

你可能听过一句话:“把一个大应用拆成很多小应用,每个小应用各管各的。“这就是微服务的基本思想。

打个比方,单体应用就像一个大食堂,炒菜、打饭、收银全在一个窗口。微服务就像美食广场,每个档口只做一种菜,互不干扰。

在 NestJS 里,微服务本质上就是一个不用 HTTP 做通信的应用。它用的是别的传输层,比如 TCP、Redis、NATS、RabbitMQ、Kafka、gRPC 等。

Note

NestJS 里微服务和普通应用用的还是同一套东西:依赖注入、装饰器、管道、守卫、拦截器,全都通用。区别只在于传输层不同。

安装依赖

开始之前,先装好微服务相关的包:

npm i --save @nestjs/microservices

这个包提供了所有微服务需要的装饰器、类和接口。

创建微服务

创建一个微服务,跟创建普通 HTTP 应用很像,只是把 NestFactory.create() 换成 NestFactory.createMicroservice()

// main.ts
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.TCP,
      options: {
        host: 'localhost',
        port: 3001,
      },
    },
  );
  await app.listen();
}
bootstrap();

createMicroservice() 的第二个参数是一个配置对象,有两个关键字段:

字段作用
transport指定传输层类型,比如 Transport.TCPTransport.REDIS
options传输层的具体配置,不同传输层配置不同

TCP 传输层常用的配置项:

配置项说明
host连接地址
port连接端口
retryAttempts重试次数,默认 0
retryDelay重试间隔(毫秒),默认 0
Tip

不指定 transport 的话,默认用 TCP。所以如果你就想用 TCP,其实可以省略这个字段。

混合应用

有时候你不想把应用完全变成微服务,而是想让它同时支持 HTTP 和微服务。这种场景叫”混合应用”。

比如你的主应用还是 HTTP 服务,但里面有几个接口需要跟其他微服务通信。这时候用 connectMicroservice() 方法:

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

async function bootstrap() {
  // 先创建普通 HTTP 应用
  const app = await NestFactory.create(AppModule);

  // 再挂一个微服务
  app.connectMicroservice<MicroserviceOptions>({
    transport: Transport.TCP,
    options: { port: 3001 },
  });

  // 先启动所有微服务
  await app.startAllMicroservices();
  // 再启动 HTTP 服务
  await app.listen(3000);
}
bootstrap();

注意顺序:先 startAllMicroservices(),再 listen()。如果你挂了多个微服务,startAllMicroservices() 会一次性把它们全启动。

消息模式(Pattern)

微服务之间通信靠的是模式匹配。模式就是一个普通值,可以是字符串,也可以是一个对象。

打个比方,你去快递柜取件,取件码就是”模式”。快递柜根据取件码找到对应的包裹给你。微服务也一样,根据模式找到对应的处理函数。

NestJS 支持两种通信方式:

  • 请求-响应:发了消息要等回复
  • 事件:发了消息不等回复,像发朋友圈一样

请求-响应模式

请求-响应是最常用的通信方式。用 @MessagePattern() 装饰器来定义消息处理器:

import { Controller } from '@nestjs/common';
import { MessagePattern } from '@nestjs/microservices';

@Controller()
export class MathController {
  @MessagePattern({ cmd: 'sum' })
  accumulate(data: number[]): number {
    return (data || []).reduce((a, b) => a + b);
  }
}

这里 { cmd: 'sum' } 就是消息模式。当客户端发来匹配这个模式的消息时,accumulate() 就会被调用。

Note

@MessagePattern() 只能用在控制器里。放在 Provider 里不会生效,NestJS 会直接忽略它。

异步响应

消息处理器支持异步,你可以用 async/await

@MessagePattern({ cmd: 'sum' })
async accumulate(data: number[]): Promise<number> {
  return (data || []).reduce((a, b) => a + b);
}

也可以返回 Observable,这样会多次响应,每发一个值就响应一次:

import { Observable, from } from 'rxjs';

@MessagePattern({ cmd: 'sum' })
accumulate(data: number[]): Observable<number> {
  return from([1, 2, 3]);
}

上面这个例子会响应三次,分别返回 1、2、3。

获取请求上下文

有时候你需要获取更多请求信息,比如 NATS 的原始主题、Kafka 的消息头。这时候用 @Payload()@Ctx() 装饰器:

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

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

@Payload() 还可以传属性名来提取特定字段:

@MessagePattern({ cmd: 'get_user' })
getUser(@Payload('id') id: number) {
  return this.userService.findOne(id);
}

事件模式

有些场景你不需要回复。比如用户注册成功后,通知其他服务”有新用户了”。这种场景用事件模式更合适。

@EventPattern() 装饰器定义事件处理器:

import { Controller } from '@nestjs/common';
import { EventPattern } from '@nestjs/microservices';

@Controller()
export class UserController {
  @EventPattern('user_created')
  async handleUserCreated(data: Record<string, unknown>) {
    console.log('新用户创建:', data);
  }
}
Tip

同一个事件模式可以注册多个处理器,它们会并行执行。比如用户创建后,既要发欢迎邮件,又要初始化用户积分,两个处理器可以同时干活。

请求-响应 vs 事件,怎么选?

场景推荐方式
需要拿到返回值请求-响应
只是通知,不需要回复事件
用 Kafka 等流式传输层事件(更契合)
需要确认消息被接收请求-响应

请求-响应会创建两个逻辑通道(一个发数据,一个等回复),有一定开销。如果不需要返回值,用事件模式更轻量。

客户端(ClientProxy)

服务端写好了,客户端怎么调用呢?NestJS 提供了 ClientProxy 类。

通过 ClientsModule 注册

最常用的方式是通过 ClientsModule.register() 注册客户端:

import { Module } from '@nestjs/common';
import { ClientsModule, Transport } from '@nestjs/microservices';

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

name 是注入令牌,后面用它来注入客户端实例。transportoptions 跟服务端配置类似。

注入并使用

注册好之后,在需要的地方用 @Inject() 注入:

import { Injectable, Inject } from '@nestjs/common';
import { ClientProxy } from '@nestjs/microservices';

@Injectable()
export class AppService {
  constructor(
    @Inject('MATH_SERVICE') private client: ClientProxy,
  ) {}

  // 请求-响应:用 send()
  accumulate(): Observable<number> {
    const pattern = { cmd: 'sum' };
    const payload = [1, 2, 3];
    return this.client.send<number>(pattern, payload);
  }

  // 事件:用 emit()
  notifyUserCreated() {
    this.client.emit('user_created', { name: '张三' });
  }
}
Note

send() 返回的是冷 Observable,你必须订阅它才会真正发送消息。emit() 返回的是热 Observable,不管你有没有订阅,它都会立刻发送事件。

动态配置

如果连接信息要从 ConfigService 读取,用 registerAsync()

@Module({
  imports: [
    ClientsModule.registerAsync([
      {
        imports: [ConfigModule],
        name: 'MATH_SERVICE',
        useFactory: async (configService: ConfigService) => ({
          transport: Transport.TCP,
          options: {
            host: configService.get('MATH_HOST'),
            port: configService.get('MATH_PORT'),
          },
        }),
        inject: [ConfigService],
      },
    ]),
  ],
})
export class AppModule {}

其他注入方式

除了 ClientsModule,还有两种方式可以获取 ClientProxy

方式一:用 ClientProxyFactory

@Module({
  providers: [
    {
      provide: 'MATH_SERVICE',
      useFactory: (configService: ConfigService) => {
        return ClientProxyFactory.create({
          transport: Transport.TCP,
          options: { port: 3001 },
        });
      },
      inject: [ConfigService],
    },
  ],
})

方式二:用 @Client() 装饰器

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

@Client({ transport: Transport.TCP })
client: ClientProxy;
Tip

@Client() 装饰器虽然简单,但不推荐用。因为它不好测试,也不好共享实例。优先用 ClientsModule 的方式。

懒连接

ClientProxy 是懒加载的。它不会一创建就连接,而是在第一次调用时才建立连接,之后复用。

如果你想让应用等连接建立好再开始工作,可以在 onApplicationBootstrap 生命周期钩子里手动连接:

async onApplicationBootstrap() {
  await this.client.connect();
}

如果连接失败,connect() 会抛出错误。

监听状态和事件

你可以订阅客户端的 status 流来获取连接状态变化:

this.client.status.subscribe((status: TcpStatus) => {
  console.log('连接状态:', status);
});

TCP 传输层会发出 connecteddisconnected 事件。

也可以监听内部事件,比如错误事件:

this.client.on('error', (err) => {
  console.error('微服务通信出错:', err);
});

超时处理

分布式系统里,远程服务随时可能挂掉。如果不设超时,你的请求可能会一直等下去。

用 RxJS 的 timeout 操作符就能搞定:

import { timeout } from 'rxjs/operators';

this.client
  .send<number>(pattern, data)
  .pipe(timeout(5000))
  .subscribe(
    (result) => console.log('结果:', result),
    (err) => console.error('超时了:', err),
  );

上面这段代码设置了 5 秒超时。如果 5 秒内微服务没响应,就会抛出超时错误。

Tip

踩坑经验:生产环境一定要设超时。不然一个微服务挂了,调用它的服务全部阻塞,最终整个系统雪崩。

RPC 异常处理

微服务里抛异常用的是 RpcException,不是普通的 HTTP 异常:

import { Controller } from '@nestjs/common';
import { MessagePattern, RpcException } from '@nestjs/microservices';

@Controller()
export class UserController {
  @MessagePattern({ cmd: 'get_user' })
  getUser(id: number) {
    const user = this.userService.findOne(id);
    if (!user) {
      throw new RpcException('用户不存在');
    }
    return user;
  }
}

可以用 RPC 异常过滤器来统一处理:

import { Catch, RpcExceptionFilter, ArgumentsHost } from '@nestjs/common';
import { RpcException } from '@nestjs/microservices';
import { Observable, throwError } from 'rxjs';

@Catch(RpcException)
export class RpcExceptionFilter implements RpcExceptionFilter<RpcException> {
  catch(exception: RpcException, host: ArgumentsHost): Observable<any> {
    return throwError(() => ({
      statusCode: exception.getError()['statusCode'],
      message: exception.message,
    }));
  }
}

动态配置微服务

如果微服务的配置需要从 ConfigService 动态获取(而不是写死在代码里),可以用 AsyncMicroserviceOptions

import { ConfigService } from '@nestjs/config';
import { AsyncMicroserviceOptions, Transport } from '@nestjs/microservices';

async function bootstrap() {
  const app = await NestFactory.createMicroservice<AsyncMicroserviceOptions>(
    AppModule,
    {
      useFactory: (configService: ConfigService) => ({
        transport: Transport.TCP,
        options: {
          host: configService.get<string>('HOST'),
          port: configService.get<number>('PORT'),
        },
      }),
      inject: [ConfigService],
    },
  );
  await app.listen();
}
bootstrap();

这样配置信息就可以从环境变量、配置文件等地方读取,不用硬编码。

TLS 加密

如果微服务之间走的是公网通信,一定要开启 TLS 加密。NestJS 的 TCP 传输层内置了 TLS 支持。

服务端配置私钥和证书:

import * as fs from 'fs';

const app = await NestFactory.createMicroservice<MicroserviceOptions>(
  AppModule,
  {
    transport: Transport.TCP,
    options: {
      tlsOptions: {
        key: fs.readFileSync('./keys/server.key', 'utf8'),
        cert: fs.readFileSync('./keys/server.cert', 'utf8'),
      },
    },
  },
);

客户端配置 CA 证书:

ClientsModule.register([
  {
    name: 'MATH_SERVICE',
    transport: Transport.TCP,
    options: {
      tlsOptions: {
        ca: [fs.readFileSync('./keys/ca.cert', 'utf-8')],
      },
    },
  },
])

请求作用域

微服务里也支持请求作用域。如果你需要每个请求都有独立的处理器实例,可以用 Scope.REQUEST

import { Injectable, Scope, Inject } from '@nestjs/common';
import { CONTEXT, RequestContext } from '@nestjs/microservices';

@Injectable({ scope: Scope.REQUEST })
export class CatsService {
  constructor(@Inject(CONTEXT) private ctx: RequestContext) {}
}

RequestContext 包含两个属性:

interface RequestContext<T = any> {
  pattern: string | Record<string, any>; // 消息模式
  data: T; // 消息数据
}

小结

本章覆盖了 NestJS 微服务的核心知识点:

  • 微服务就是不用 HTTP 传输的 NestJS 应用
  • createMicroservice() 创建纯微服务,用 connectMicroservice() 创建混合应用
  • @MessagePattern() 处理请求-响应消息
  • @EventPattern() 处理事件消息
  • ClientProxysend() 发消息等回复,emit() 发事件不等回复
  • send() 是冷 Observable 要订阅才发,emit() 是热 Observable 立刻发
  • 生产环境一定要设超时,防止级联故障
  • 微服务里用 RpcException 处理异常
  • 公网通信要开 TLS 加密

下一章我们来看 NestJS 支持的各种传输层的具体用法。