NestJS 微服务架构:模块化、消息与网关

系统覆盖用 NestJS 构建 TypeScript 微服务架构的完整实践:模块化与依赖注入、TCP/Redis/Kafka 传输层选型、客户端-服务端微服务模式、请求-应答与事件驱动消息模式、BFF 网关聚合与 API 编排、JWT 鉴权与消息校验、重试超时与熔断容错,以及日志追踪指标三维可观测性,帮助开发者用框架约束搭建类型安全、可观测、可容错的微服务网络。

引言

微服务架构的复杂度大半在「通信」:传输、可靠、鉴权、容错、观测。NestJS 用模块化 + 依赖注入 + 可插拔传输层,把微服务拆成一套有纪律的类型安全模式。本文覆盖 TCP/Redis/Kafka 传输层、客户端-服务端两种角色、请求-应答与事件驱动两种消息语义、BFF 网关聚合、JWT 鉴权与 zod 消息校验、重试超时熔断,最后落到日志/追踪/指标的三维可观测性。

前置:/typescript-nodejs-backend/(Node 服务端)、/typescript-decorators-metaprogramming/(装饰器元编程)、/typescript-typed-events-streams/(事件与消息模型)。

目录

1. NestJS 微服务架构总览

核心抽象:ClientProxy(客户端发送的代理,封装传输层细节)、@MessagePattern(请求-应答,有返回值)、@EventPattern(事件驱动,无返回)、Transport 枚举(TCP/REDIS/KAFKA)。

import { Controller } from "@nestjs/common";
import { MessagePattern, Payload } from "@nestjs/microservices";

@Controller()
export class UserController {
  @MessagePattern("user.get")
  async getUser(@Payload() dto: { id: number }) {
    return this.userService.findById(dto.id); // 返回值回传客户端
  }
}

工程纪律:消息 pattern 是服务间契约,放共享类型包(见 /typescript-api-type-generation/),避免两处手抄字符串。

2. 模块化与依赖注入

Module 划边界、DI 容器管依赖、测试可替换 provider:

@Module({
  imports: [ConfigModule.forRoot(), TypeOrmModule.forFeature([User])],
  controllers: [UserController],
  providers: [UserService, Logger],
  exports: [UserService],
})
export class UserModule {}

@Injectable()
export class UserService {
  constructor(
    private readonly repo: Repository<User>,
    private readonly logger: Logger,
  ) {}
}

要点:频繁依赖(Config、Logger)标 @Global() 减少重复 import;exports 显式暴露对外 provider;循环依赖用 forwardRef(能避免就避免,拖慢启动);用类令牌(@Inject(Connection))优于字符串令牌,类型检查更强。坑:TypeOrmModule.forFeature 忘 import 实体会报 Nest can't resolve dependencies,先看哪个 provider 缺依赖。

3. 传输层:TCP、Redis 与 Kafka

选型决定可靠性、吞吐与运维成本:

传输层可靠性适用
TCP中(需自管重连)内网服务请求-应答
Redis中(Pub/Sub 易丢消息)临时事件广播
Kafka高(持久化 + 重放)事件溯源/数据管道
// TCP 服务端
const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
  transport: Transport.TCP,
  options: { host: "0.0.0.0", port: 3001 },
});
await app.listen();

// Kafka
const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
  transport: Transport.KAFKA,
  options: {
    client: { brokers: ["kafka:9092"] },
    consumer: { groupId: "user-svc" },
  },
});

选型铁律:业务请求用 TCP;解耦事件流用 Kafka(消费者组 + 重放);临时广播用 Redis。别用 Redis 承载不能丢的领域事件——Pub/Sub 无人消费即丢失。

4. 微服务客户端与服务端模式

客户端通过 ClientProxyFactory 创建并注入:

@Injectable()
export class OrderService {
  constructor(@Inject("USER_SERVICE") private readonly userClient: ClientProxy) {}

  async getUser(id: number) {
    const user$ = this.userClient.send<UserDto>("user.get", { id }); // send 返回 Observable
    return await lastValueFrom(user$); // Observable → Promise
  }
}
providers: [
  {
    provide: "USER_SERVICE",
    useFactory: () =>
      ClientProxyFactory.create({ transport: Transport.TCP, options: { host: "user-svc", port: 3001 } }),
  },
]

要点:send 返回 Observable,可用 RxJS 管道做重试/超时再 lastValueFrom;emit 发送即返回(事件驱动);ClientProxy 做单例 provider 避免每请求建连接;超时必须有——send 在服务不可达时不会自动超时,务必 timeout(5000)。坑:send<UserDto> 的泛型只是「承诺」,入站端必须再做运行时校验(§7)。

5. 消息模式:请求-应答与事件驱动

// 请求-应答:同步语义,等结果
const order = await lastValueFrom(
  this.orderClient.send("order.create", dto).pipe(timeout(5000))
);

// 事件驱动:广播给所有消费者
this.userClient.emit("user.created", { id: userId });
@Controller()
export class NotificationConsumer {
  @EventPattern("user.created")
  async onUserCreated(@Payload() payload: { id: number }) {
    await this.notifyService.sendWelcome(payload.id); // 不返回内容
  }
}
维度请求-应答事件驱动
同步性等结果即返回
消费者一个多个(广播)
失败处理客户端重试消费者重试/死信
适用查询、命令确认通知、解耦、副作用

坑:请求-应答的 pattern 名被两个服务注册时,消息发给「第一个可用消费者」,行为不可预期——pattern 名全局唯一。

6. 网关聚合:BFF 与 API 编排

网关把多个微服务响应聚合给前端:

@Controller("orders")
export class OrderGatewayController {
  constructor(
    private readonly orderClient: ClientProxy,
    private readonly userClient: ClientProxy,
  ) {}

  @Get(":id")
  async getOrderDetail(@Param("id") id: string) {
    const [order, user] = await Promise.allSettled([
      lastValueFrom(this.orderClient.send("order.get", { id })),
      lastValueFrom(this.userClient.send("user.get", { id })),
    ]);
    return { order: order.status === "fulfilled" ? order.value : null, buyer: user };
  }
}

网关职责:编排(并行聚合/串行依赖)、裁剪(只透传前端需要的字段)、错误归一(微服务错误码翻译成 HTTP 状态码)、鉴权入口(JWT 在此校验一次)。坑:Promise.all 一个服务失败整个请求 500——非关键依赖用 Promise.allSettled 降级,保主路径可用。

7. 鉴权与消息校验

微服务边界上,鉴权与校验是两道必须主动设防的闸门。

// 网关侧:全局 JWT 守卫
@Injectable()
export class JwtAuthGuard implements CanActivate {
  canActivate(ctx: ExecutionContext): boolean {
    const req = ctx.switchToHttp().getRequest();
    const token = req.headers.authorization?.replace(/^Bearer /, "");
    const payload = this.jwt.verify<{ sub: string }>(token ?? "");
    req.userId = payload.sub;
    return true;
  }
}

消息校验——跨服务消息是「外网输入」,必须运行时校验:

import { z } from "zod";
const CreateUserSchema = z.object({
  email: z.string().email(),
  name: z.string().min(1).max(100),
  roles: z.array(z.enum(["admin", "user"])).default(["user"]),
});

@MessagePattern("user.create")
async createUser(@Payload() raw: unknown) {
  const dto = CreateUserSchema.parse(raw); // 类型安全 + 运行期安全
  return this.userService.create(dto);
}

关键认知:TS 类型在编译后消失,send<UserDto> 不构成运行期保障。跨服务边界 = 编译期类型(共享包)+ 运行期 Schema(zod),缺一不可。

8. 重试、超时与熔断

可靠性靠「客户端主动容错」,三件套:

async function callOrder(dto: OrderDto) {
  return await lastValueFrom(
    this.orderClient.send("order.create", dto).pipe(
      timeout(3000),                     // 1. 超时
      retry({ count: 2, delay: 200 }),   // 2. 有限重试
      catchError((err) => {              // 3. 熔断降级
        this.circuitBreaker.recordFailure();
        if (this.circuitBreaker.isOpen()) throw new ServiceUnavailable("order-svc");
        throw err;
      })
    )
  );
}

熔断状态机:CLOSED(正常转发)→ 失败率超阈值 → OPEN(直接短路)→ 冷却后 → HALF_OPEN(放探测请求)→ 成功回 CLOSED / 失败回 OPEN。注意:写操作重试要幂等(Idempotency-Key);退避用指数 + 抖动防雪崩;超时 3s × 重试 2 次,别无限重试。坑:重试打在已超时的慢请求上会堆积连接——用 retryWhen + delayWhen 限总时长。

9. 可观测性:日志、追踪与指标

排查微服务问题靠三维观测:

// 1. 结构化日志:JSON + 服务名 + traceId
app.useLogger(MyStructuredLogger);

// 2. 追踪:入口生成 traceId,跨服务透传
import { middleware as requestContext } from "cls-hooked";
app.use(requestContext("req"));
const trace = req.header("x-trace-id") ?? crypto.randomUUID();
this.userClient.send("user.get", { id, traceId: trace });

// 3. 指标:计数器 + 直方图
const orderLatency = new Histogram({ name: "order_create_duration_seconds", labelNames: ["result"] });
const start = Date.now();
try { await callOrder(dto); orderLatency.observe({ result: "ok" }, (Date.now() - start) / 1000); }
catch { orderLatency.observe({ result: "error" }, (Date.now() - start) / 1000); throw err; }

清单:日志用结构化 JSON 带 service/traceId/level,禁止散装 console.log;traceId 入口生成并透传所有下游;指标至少覆盖请求量/延迟 p95/p99/错误率/下游状态;每个服务暴露 /health 与 /ready,网关只打流量到就绪实例。坑:日志没 traceId 等于没日志——上下文中间件必须全局最外层注册。

10. 速查表与一句话记忆

场景做法
模块划分Module 边界 + exports 显式
传输层TCP 请求 / Kafka 事件 / Redis 广播
客户端ClientProxy 单例 + send/emit + timeout
消息校验zod parse 入站载荷
网关Promise.allSettled 聚合 + 字段裁剪
鉴权网关 JWT + 微服务端也校验
容错timeout + 幂等重试 + 熔断
可观测JSON 日志 + traceId 透传 + 指标

一句话记忆:NestJS 微服务 = 模块化 DI 分边界 + 传输层按可靠度选型(TCP/Kafka/Redis)+ 请求-应答与事件两种消息语义 + 网关并行聚合 + 边界双重校验(类型 + Schema)+ 超时重试熔断三件套 + 日志追踪指标三维观测。

延伸阅读

  • /typescript-nodejs-backend/ — Node 服务端与进程模型
  • /typescript-decorators-metaprogramming/ — 装饰器与元编程原理
  • /typescript-typed-events-streams/ — 事件模型与消息总线
  • /typescript-runtime-validation-typesafe/ — 入站消息运行时校验
  • /typescript-error-handling-result/ — 服务错误与 Result 建模
  • Kafka 专题 — 事件流与消费者组

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「typescript」更多文章

  1. TypeScript 应用安全加固:依赖、注入与敏感信息防护
  2. TypeScript Monorepo 工程化:pnpm、Turborepo 与多包协作
  3. Node.js Worker Threads:TypeScript 并行计算实战