引言
微服务架构的复杂度大半在「通信」:传输、可靠、鉴权、容错、观测。NestJS 用模块化 + 依赖注入 + 可插拔传输层,把微服务拆成一套有纪律的类型安全模式。本文覆盖 TCP/Redis/Kafka 传输层、客户端-服务端两种角色、请求-应答与事件驱动两种消息语义、BFF 网关聚合、JWT 鉴权与 zod 消息校验、重试超时熔断,最后落到日志/追踪/指标的三维可观测性。
前置:/typescript-nodejs-backend/(Node 服务端)、/typescript-decorators-metaprogramming/(装饰器元编程)、/typescript-typed-events-streams/(事件与消息模型)。
目录
- 1. NestJS 微服务架构总览
- 2. 模块化与依赖注入
- 3. 传输层:TCP、Redis 与 Kafka
- 4. 微服务客户端与服务端模式
- 5. 消息模式:请求-应答与事件驱动
- 6. 网关聚合:BFF 与 API 编排
- 7. 鉴权与消息校验
- 8. 重试、超时与熔断
- 9. 可观测性:日志、追踪与指标
- 10. 速查表与一句话记忆
- 延伸阅读
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 专题 — 事件流与消费者组
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。