Node.js 微服务架构实战:服务拆分、通信、治理与分布式事务

Node.js 微服务架构完整指南:从单体到微服务的拆分策略、DDD 限界上下文、REST/gRPC/GraphQL 通信协议对比、服务发现、API 网关、BFF 模式、熔断重试、分布式链路追踪、Saga 分布式事务与 NestJS gRPC 微服务实战。

微服务架构已成为现代后端系统的主流范式。在 Node.js 生态中,从 Express 单体应用到 NestJS 模块化微服务集群的演进,不仅是技术栈的升级,更是组织边界与工程哲学的重构。然而,微服务并非免费的午餐——它用部署复杂性换取了团队并行度和技术多样性,用网络延迟换取了独立可扩展性。本文从工程实践出发,系统梳理 Node.js 微服务架构的核心议题:如何拆分服务、如何通信、如何治理失败、如何追踪跨服务请求、如何保证分布式事务一致性,最终给出一个基于 NestJS 与 gRPC 的可运行微服务示例。


1. 微服务与单体架构的权衡

1.1 单体架构的演进困境

单体架构是绝大多数项目的起始形态。所有业务逻辑、数据访问、API 层打包在同一个代码库和同一个部署单元中。这种模式在项目初期具有显著优势:开发简单、调试方便、事务一致、部署单一。然而,随着团队规模扩大和功能复杂度增长,单体架构会遇到结构性瓶颈。

单体架构的典型痛点:

痛点表现影响
代码耦合一个模块的修改可能意外破坏另一模块回归风险高,测试负担重
技术锁死全栈共享同一技术选型难以引入更适合特定领域的技术
部署冲突多团队共享代码库,发布需协调部署频率被迫降低
扩展粒度只能整体横向扩展,无法针对热点服务单独扩容资源浪费严重
认知负荷新成员需要理解整个代码库才能修改单个功能onboarding 周期长

当一个团队的代码库超过数十万行,构建时间超过十分钟,任何改动都需要回归全量测试时,就到了认真考虑拆分的时候。

1.2 微服务的核心收益与代价

微服务架构将单一应用拆分为一组小型服务,每个服务运行在自己的进程中,围绕业务能力构建,通过轻量级通信机制互相协作。独立部署、独立扩展、独立技术选型是微服务的三大核心特征。

微服务的收益:

  • 团队自治:每个服务由小型跨职能团队(通常 2-8 人,符合"两个披萨团队"原则)全权负责,从设计到运维端到端 ownership
  • 技术多样性:不同的服务可以使用最适合其业务领域的技术栈。支付服务用 Java 保证事务安全,实时推送服务用 Node.js 利用事件驱动优势
  • 弹性隔离:单个服务的故障不会导致整个系统崩溃。订单服务故障时,用户仍可浏览商品目录
  • 独立扩展:高并发服务可以单独扩容,低负载服务保持精简配置,优化云资源成本

微服务的代价:

  • 分布式复杂性:网络调用引入延迟、超时、重试、幂等性、部分失败等全新问题域
  • 数据一致性:跨服务的事务无法使用单机 ACID 事务保证,必须引入 Saga 等最终一致性模式
  • 运维负担:服务数量从 1 个变成几十个甚至上百个,监控、日志、链路追踪、配置管理的复杂度指数级上升
  • 调试困难:一个用户请求可能经过 5-10 个服务,定位问题需要在分布式链路中追踪线索

关键判断:微服务的收益往往体现在组织层面而非纯技术层面。如果你的系统复杂度尚未达到让多个团队有效并行工作的程度,过早拆分微服务会得不偿失。Martin Fowler 的观点仍然成立:“只有当你的单体应用已大到让团队感到痛苦时,才考虑微服务。”


2. 服务拆分策略:DDD 限界上下文

2.1 按业务边界拆分,而非技术层

最常见的反模式是按技术层拆分服务:一个 User Service、一个 Order Service、一个 Payment Service,每个服务都包含 Controller、Service、Repository 三层。这种拆分只是将单体内的包结构变成了网络边界,服务之间仍然紧密耦合,却没有获得微服务的任何好处。

正确的拆分依据是业务能力(Business Capability),而非技术层次。每个微服务应对应一个内聚的业务领域,拥有独立的数据存储和业务规则。

2.2 领域驱动设计(DDD)与限界上下文

领域驱动设计提供了微服务拆分最系统的理论框架。核心概念包括:

  • 限界上下文(Bounded Context):模型有明确边界的语义环境。同一个词在不同上下文中含义不同。例如"Customer"在销售上下文中是"潜在客户",在售后上下文中是"已购用户",在物流上下文中是"收件人"。将它们强行建模为同一个实体是微服务耦合的根源
  • 聚合(Aggregate):一致性边界内的实体和值对象的集合。订单聚合包含订单头、订单行、配送地址,它们在同一事务中保持强一致
  • 领域事件(Domain Event):聚合状态变化后发布的不可变事实。订单已创建、库存已扣减、支付已完成都是领域事件

拆分实践原则:

  1. 高内聚低耦合:一个限界上下文内的概念紧密相关,上下文之间通过定义良好的接口通信
  2. 独立部署单元验证:问自己"这个上下文能否独立演进和部署?“如果团队 A 修改它需要团队 B 同步修改,说明边界划分有误
  3. 数据库独立:每个微服务拥有独立的数据库 schema。严禁跨服务直接访问数据库,否则耦合度退化为分布式单体
  4. 团队边界对齐:微服务的边界应与康威定律下的团队沟通结构匹配。强行让两个频繁沟通的上下文分属不同服务只会增加不必要的网络开销
电商系统限界上下文划分示例:

┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐
│   商品目录上下文   │  │   订单上下文     │  │   库存上下文     │
│  Product Catalog │  │     Order       │  │   Inventory     │
│  - Product       │  │  - OrderHeader  │  │  - Stock        │
│  - Category      │  │  - OrderLine    │  │  - Reservation  │
│  - Search Index  │  │  - Status       │  │  - Warehouse    │
└────────┬────────┘  └────────┬────────┘  └────────┬────────┘
         │                    │                    │
         └────────────────────┴────────────────────┘
                              │
                    ┌─────────┴─────────┐
                    │    事件总线         │
                    │  (Domain Events)  │
                    └───────────────────┘

┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐
│   支付上下文     │  │   用户上下文     │  │   物流上下文     │
│    Payment      │  │     User        │  │   Fulfillment   │
│  - Payment      │  │  - Profile      │  │  - Shipment     │
│  - Refund       │  │  - Auth         │  │  - Tracking     │
│  - Gateway      │  │  - Preference   │  │  - Carrier API  │
└─────────────────┘  └─────────────────┘  └─────────────────┘

3. 服务间通信:REST、gRPC、GraphQL 对比

微服务间的通信机制是架构选型中最关键的决策之一。不同的协议在性能、类型安全、易用性、工具生态方面各有侧重。

3.1 三种协议的深度对比

维度REST over HTTP/1.1gRPC over HTTP/2GraphQL over HTTP
协议层HTTP/1.1 + JSONHTTP/2 + Protocol BuffersHTTP/1.1 或 HTTP/2 + JSON
序列化JSON(文本)Protobuf(二进制)JSON(文本)
性能中等(JSON 解析 + 文本传输)高(二进制 + 头部压缩 + 多路复用)中等(灵活查询带来额外解析成本)
类型安全弱(依赖 OpenAPI/Swagger 补充)强(编译时类型检查)强(Schema 驱动)
流式支持需 SSE / WebSocket原生支持 unary / server-stream / client-stream / bidirectional需 Subscription(WebSocket)
浏览器兼容原生支持需 gRPC-Web 转接代理原生支持
调试工具cURL、Postman、Swagger UIgrpcurl、BloomRPCGraphiQL、Playground
适用场景对外公开 API、简单 CRUD内部服务间高频通信BFF 层、聚合多服务数据

3.2 REST:简单、通用、浏览器原生

REST 仍是 Node.js 生态中最广泛使用的通信方式。其优势在于简单直观,任何语言和环境都能通过 HTTP 发起请求。Express、Fastify、NestJS 都提供完备的 REST 支持。

// NestJS REST 控制器示例
@Controller('orders')
export class OrderController {
  constructor(private readonly orderService: OrderService) {}

  @Post()
  async createOrder(@Body() dto: CreateOrderDto): Promise<Order> {
    return this.orderService.create(dto);
  }

  @Get(':id')
  async getOrder(@Param('id') id: string): Promise<Order> {
    return this.orderService.findById(id);
  }
}

REST 的挑战在于:

  • 缺乏原生类型契约,依赖文档和人为遵守
  • HTTP/1.1 的队头阻塞问题在高并发场景下明显
  • 版本控制历来是痛点(路径版本 /v2,还是媒体类型?)

3.3 gRPC:高性能的内部通信首选

gRPC 是 Google 开源的 RPC 框架,基于 HTTP/2 和 Protocol Buffers。在内部微服务通信中,gRPC 的性能优势无可争议。

gRPC 的核心优势:

  • HTTP/2 多路复用:单个 TCP 连接上并行传输多个请求,避免连接池爆炸
  • Protobuf 高效序列化:比 JSON 体积小 3-5 倍,序列化速度快一个数量级
  • 强类型契约.proto 文件定义服务接口,客户端和服务端通过代码生成保持类型一致
  • 四种调用模式:Unary(单次请求响应)、Server Streaming、Client Streaming、Bidirectional Streaming
  • 原生流式支持:适合实时数据推送和双向通信场景
// order.proto —— gRPC 服务定义
syntax = "proto3";

package ecommerce;

service OrderService {
  rpc CreateOrder (CreateOrderRequest) returns (Order);
  rpc GetOrder (GetOrderRequest) returns (Order);
  rpc StreamOrderUpdates (StreamOrderRequest) returns (stream OrderUpdate);
}

message CreateOrderRequest {
  string user_id = 1;
  repeated OrderItem items = 2;
  string currency = 3;
}

message OrderItem {
  string product_id = 1;
  int32 quantity = 2;
  int64 unit_price_cents = 3;
}

message Order {
  string order_id = 1;
  string user_id = 2;
  OrderStatus status = 3;
  int64 total_cents = 4;
  repeated OrderItem items = 5;
  string created_at = 6;
}

enum OrderStatus {
  PENDING = 0;
  PAID = 1;
  SHIPPED = 2;
  DELIVERED = 3;
  CANCELLED = 4;
}

message GetOrderRequest {
  string order_id = 1;
}

message StreamOrderRequest {
  string order_id = 1;
}

message OrderUpdate {
  string order_id = 1;
  OrderStatus new_status = 2;
  string timestamp = 3;
}

gRPC 的局限在于浏览器直接支持较弱(需要 grpc-web + 代理)、二进制格式对调试不够直观、.proto 文件的版本管理需要清晰的治理流程。

3.4 选型建议矩阵

对外公共 API(第三方集成) → REST + OpenAPI
内部服务间同步调用        → gRPC(NestJS 支持良好)
聚合层 / BFF(前端数据拼装) → GraphQL
异步解耦 / 事件驱动        → 消息队列(RabbitMQ / Kafka)

4. 服务发现:Consul、etcd、Kubernetes DNS

微服务的 IP 和端口是动态变化的,特别是在容器编排环境中。服务发现解决了"订单服务如何找到库存服务当前实例"这一基础问题。

4.1 服务发现模式

客户端发现:客户端直接查询注册中心获取服务实例列表,自己实现负载均衡。Eureka + Ribbon 是典型代表。

服务端发现:客户端请求发送到负载均衡器或 API 网关,由中间层负责服务解析和流量分发。Kubernetes Service、AWS ALB 属于此类。

4.2 主流方案对比

方案架构一致性模型健康检查适用场景
Consul分布式 + RaftCP(强一致)HTTP/TCP/脚本多云环境、需要精细服务治理
etcd分布式 + RaftCP需配合工具Kubernetes 核心依赖、配置存储
ZooKeeper分布式 + ZABCPTCP传统大数据生态(Kafka、Hadoop)
Kubernetes DNS内置-Kubelet 探针纯 K8s 环境,零额外依赖
AWS CloudMap托管AP/CP 可选AWS 集成纯 AWS 环境

Consul + Node.js 示例:

// consul-service-registrar.js
const Consul = require('consul');

const consul = new Consul({ host: 'consul-server', port: 8500 });
const serviceId = `order-service-${process.env.HOSTNAME || Date.now()}`;

async function register() {
  await consul.agent.service.register({
    id: serviceId,
    name: 'order-service',
    tags: ['nodejs', 'v1'],
    port: 3000,
    check: {
      http: 'http://localhost:3000/health',
      interval: '10s',
      timeout: '5s',
      deregistercriticalserviceafter: '30s',
    },
  });
  console.log(`Registered ${serviceId} with Consul`);
}

async function deregister(signal) {
  await consul.agent.service.deregister(serviceId);
  console.log(`Deregistered ${serviceId} (${signal})`);
  process.exit(0);
}

process.on('SIGINT', () => deregister('SIGINT'));
process.on('SIGTERM', () => deregister('SIGTERM'));

register();

Kubernetes DNS 方式(推荐用于容器化环境):

在 Kubernetes 中,Service 对象自动提供 DNS 记录。order-service.default.svc.cluster.local 解析到所有健康的 Pod IP,无需额外配置:

apiVersion: v1
kind: Service
metadata:
  name: order-service
spec:
  selector:
    app: order-service
  ports:
    - port: 80
      targetPort: 3000

Node.js 应用只需通过环境变量或配置中心获取服务名,通过 DNS 解析获取实例地址。这是最简洁的服务发现方式,没有额外的运维组件。


5. API Gateway 模式:Kong、Ambassador、Express Gateway

5.1 为什么需要 API Gateway

API Gateway 是微服务集群的统一入口,承担请求路由、认证鉴权、限流熔断、日志记录、协议转换等横切关注点。没有网关时,每个微服务都需要自行处理 JWT 验证、CORS、Rate Limiting,造成逻辑重复和配置碎片化。

无网关:  客户端 → [认证、限流、路由逻辑散布在各个服务]

有网关:  客户端 → API Gateway → 路由到各微服务
                  (认证、限流、日志、熔断集中处理)

5.2 主流网关对比

网关底层插件生态配置方式适用场景
KongOpenResty / Nginx + Lua丰富(100+ 插件)Admin API / Declarative Config大规模生产环境
TraefikGo中等标签驱动(Docker/K8s 原生)云原生 / 容器化首选
AmbassadorEnvoy Proxy丰富Kubernetes CRDK8s 原生部署
Express GatewayNode.js + Express基础YAML + Admin APINode.js 团队维护、轻量场景

Kong 限流配置示例:

# kong-declarative.yml
_format_version: "3.0"
services:
  - name: order-service
    url: http://order-service:3000
    routes:
      - name: order-routes
        paths:
          - /orders
    plugins:
      - name: rate-limiting
        config:
          minute: 100
          policy: redis
          redis_host: redis-cluster
      - name: jwt
        config:
          uri_param_names: []
          cookie_names: []
          key_claim_name: iss
          secret_is_base64: false
          claims_to_verify:
            - exp

5.3 网关的演进:从智能到薄网关

早期的网关承担了过多的业务逻辑,成为又一个分布式瓶颈。现代微服务架构倾向于**“薄网关 + Sidecar”**模式:网关只做横向关注点(TLS 终结、路由、限流),认证授权下沉到服务层(委托给 Identity Provider),流量治理通过 Sidecar 代理(如 Istio/Envoy)实现。这种分层避免了网关成为全业务的耦合点。


6. BFF(Backend for Frontend)模式

6.1 模式动机

微服务粒度是按业务领域划分的,但前端的数据需求是按页面和组件组织的。一个电商首页可能需要同时调用用户服务(获取收货地址)、商品服务(获取推荐列表)、促销服务(获取优惠券)和库存服务(获取到货时间)。如果让前端直接调用这四个微服务,不仅请求数量多,每个服务返回的数据格式也未必适合前端消费。

BFF 在客户端与微服务之间引入一层专用后端,为特定前端(Web App、iOS、Android、小程序)定制数据聚合和转换逻辑。

传统方式:              BFF 模式:

Web App ──┬─→ 用户服务   Web App ──→ Web BFF ──┬─→ 用户服务
          ├─→ 商品服务                          ├─→ 商品服务
          ├─→ 促销服务   iOS App ──→ iOS BFF ──┤─→ 促销服务
          └─→ 库存服务                          └─→ 库存服务

                        Android ──→ Android BFF ──→ ...

6.2 BFF 的职责边界

BFF 不应包含业务逻辑,它的职责仅限于:

  1. 数据聚合:将多个下游服务的响应合并为单个前端友好的数据结构
  2. 协议适配:将内部 gRPC 转换为外部 GraphQL 或 REST
  3. 视图裁剪:根据前端设备特性裁剪字段(移动端返回精简版,桌面端返回完整版)
  4. 缓存优化:在 BFF 层缓存热点数据,减少下游调用
// Web BFF —— 聚合首页数据的 NestJS 服务
@Controller('home')
export class HomeBffController {
  constructor(
    private userClient: UserGrpcClient,
    private productClient: ProductGrpcClient,
    private promoClient: PromoGrpcClient,
  ) {}

  @Get('dashboard')
  async getDashboard(@Req() req: Request) {
    const userId = req.user.sub;

    // 并行调用多个下游服务
    const [profile, recommendations, coupons] = await Promise.all([
      this.userClient.getProfile(userId),
      this.productClient.getRecommendations(userId),
      this.promoClient.getActiveCoupons(userId),
    ]);

    // 聚合为前端需要的结构
    return {
      greeting: `欢迎回来,${profile.name}`,
      defaultAddress: profile.addresses[0],
      recommendations: recommendations.map(r => ({
        id: r.id,
        title: r.title,
        price: this.formatPrice(r.price),
        image: r.images[0],
        inStock: r.stock > 0,
      })),
      coupons: coupons.filter(c => !c.expired).map(c => ({
        code: c.code,
        description: c.description,
        discount: `${c.discountPercent}%`,
      })),
    };
  }

  private formatPrice(cents: number): string {
    return ${(cents / 100).toFixed(2)}`;
  }
}

7. 熔断器、重试、超时模式

分布式系统中,故障是常态而非意外。下游服务的暂时不可用不应导致上游服务雪崩。以下三个模式是微服务弹性设计的基石。

7.1 熔断器(Circuit Breaker)

熔断器监控对下游服务的调用,当失败率达到阈值时"熔断”,后续请求直接快速失败,避免大量请求堆积等待超时的资源耗尽。经过冷却时间后,熔断器进入半开状态,允许少量探测请求通过后恢复原态。

// 基于 opossum 的熔断器封装
import CircuitBreaker from 'opossum';

const options = {
  timeout: 3000,           // 3 秒超时
  errorThresholdPercentage: 50,  // 失败率 50% 触发熔断
  resetTimeout: 30000,     // 30 秒后尝试恢复
  volumeThreshold: 10,     // 至少 10 次请求才计算失败率
};

const breaker = new CircuitBreaker(
  (paymentData) => paymentService.charge(paymentData),
  options
);

breaker.on('open', () => console.warn('Payment circuit opened'));
breaker.on('halfOpen', () => console.warn('Payment circuit half-open'));
breaker.on('close', () => console.info('Payment circuit closed'));

// 使用
async function processPayment(data: PaymentData) {
  try {
    return await breaker.fire(data);
  } catch (err) {
    if (breaker.opened) {
      // 熔断中:返回降级响应或排队稍后处理
      return { status: 'queued', message: '支付服务暂时不可用,已加入队列' };
    }
    throw err;
  }
}
状态行为
Closed(关闭)请求正常通过,记录失败率
Open(开启)请求直接快速失败,不走下游
Half-Open(半开)允许少量探测请求,成功则关闭,失败则重新开启

7.2 重试(Retry)与抖动(Jitter)

瞬时网络抖动导致的偶发失败,通过重试往往可以恢复。但盲目重试会加剧下游压力,尤其在下游已过载时。

import { backOff } from 'exponential-backoff';

async function reliableFetch(url: string) {
  return backOff(
    () => fetch(url),
    {
      numOfAttempts: 5,
      startingDelay: 100,
      maxDelay: 5000,
      timeMultiple: 2,  // 指数退避:100ms → 200ms → 400ms → ...
      jitter: 'full',   // 随机抖动避免惊群效应
      retry: (err, attempt) => {
        // 仅对可重试错误重试
        return err.status >= 500 || err.code === 'ECONNRESET';
      },
    }
  );
}

重试的关键原则:

  • 只对幂等操作重试(GET、PUT with idempotency-key)
  • GET 请求天然幂等,POST 订单创建需配合幂等键
  • 引入抖动的指数退避,防止所有客户端在同一时刻重试造成"惊群"

7.3 超时(Timeout)

没有超时设置的请求是一颗定时炸弹。一个未设置超时的请求可能在下游挂起数分钟,占满线程池或事件循环。级联超时的总时长应逐层递减,防止外层超时已过、内层仍在执行的无用功。

// Axios 请求超时配置
const client = axios.create({
  timeout: 2000,           // 2 秒请求超时
  timeoutErrorMessage: 'Service response timeout',
});

// 级联超时策略:网关层 5s → BFF 层 4s → 服务层 3s → 数据库层 2s

8. 分布式链路追踪:Jaeger / Zipkin

8.1 为什么需要链路追踪

单个用户请求可能触发十几个微服务调用。当请求延迟异常或调用失败时,仅凭各服务的独立日志无法还原完整的故事线。链路追踪为每个请求分配唯一的 Trace ID,并在每次跨服务调用时传递,将分散的日志和指标关联为一个完整的调用链。

8.2 OpenTelemetry + Jaeger 集成

OpenTelemetry 已成为链路追踪的事实标准。以下是 NestJS 应用集成 OTel 并导出到 Jaeger 的完整配置:

// tracing.ts —— NestJS 项目的 OpenTelemetry 初始化
import { NodeSDK } from '@opentelemetry/sdk-node';
import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-grpc';
import { getNodeAutoInstrumentations } from '@opentelemetry/auto-instrumentations-node';
import { Resource } from '@opentelemetry/resources';
import { SemanticResourceAttributes } from '@opentelemetry/semantic-conventions';

const sdk = new NodeSDK({
  resource: new Resource({
    [SemanticResourceAttributes.SERVICE_NAME]: 'order-service',
    [SemanticResourceAttributes.SERVICE_VERSION]: '1.0.0',
    [SemanticResourceAttributes.DEPLOYMENT_ENVIRONMENT]: 'production',
  }),
  traceExporter: new OTLPTraceExporter({
    url: process.env.JAEGER_COLLECTOR_URL || 'http://jaeger:4317',
  }),
  instrumentations: [
    getNodeAutoInstrumentations({
      '@opentelemetry/instrumentation-http': { enabled: true },
      '@opentelemetry/instrumentation-grpc': { enabled: true },
    }),
  ],
});

sdk.start();

process.on('SIGTERM', () => sdk.shutdown());

手动创建 Span:

import { trace } from '@opentelemetry/api';

const tracer = trace.getTracer('order-service');

async function createOrder(dto: CreateOrderDto) {
  return tracer.startActiveSpan('order.create', {
    attributes: { 'order.user_id': dto.userId },
  }, async (span) => {
    try {
      const order = await db.orders.insert(dto);
      span.setAttribute('order.id', order.id);
      return order;
    } catch (err) {
      span.recordException(err);
      throw err;
    } finally {
      span.end();
    }
  });
}

8.3 在代码中传递 Trace Context

跨服务调用时,Trace Context 通过 HTTP headers 或 gRPC metadata 传递。OTel 的自动埋点库已在 HTTP 和 gRPC 客户端中处理了注入逻辑,手动场景可使用 W3C Trace Context 标准:

// HTTP 请求中注入 trace context(自动,库已处理)
// 手动注入 gRPC metadata
import { propagation, context } from '@opentelemetry/api';
import { Metadata } from '@grpc/grpc-js';

function injectTraceContext(meta: Metadata): Metadata {
  const carrier = {};
  propagation.inject(context.active(), carrier);
  for (const [key, value] of Object.entries(carrier)) {
    meta.add(key, String(value));
  }
  return meta;
}

Jaeger UI 中的典型链路视图:

Trace: a3f7d2e8b1c9
├── [order-service] POST /orders           45ms
│   ├── [order-service] validate input      2ms
│   ├── [inventory-service] checkStock     18ms
│   ├── [payment-service] authorize        120ms  ⚠️ 延迟异常
│   │   ├── [payment-service] validate card  5ms
│   │   └── [payment-service] gateway call  112ms  ← 根因
│   └── [notification-service] sendEmail    8ms

9. Saga 模式:分布式事务处理

9.1 分布式事务的挑战

微服务架构中,一个业务操作往往涉及多个服务的本地数据库事务。例如创建订单需要:扣减库存(库存服务)、扣款(支付服务)、创建订单记录(订单服务)。传统的两阶段提交(2PC)在多服务场景中协调成本极高,锁定时间长,可用性差。Saga 模式通过将长事务拆分为一系列本地事务,以最终一致性来保证整体业务成功。

9.2 Saga 编排模式

编排式 Saga(Choreography):每个服务完成本地事务后发送事件,由事件触发下一个服务的操作。事件驱动,去中心化,但流程分散在各服务中,难以全局理解。

编排式 Saga(Orchestration):引入 Saga 编排器(Orchestrator)作为中央协调者,负责按顺序调用各服务,并处理失败回滚。流程集中定义,易于理解和调试,但编排器本身成为单点。

编排式 Saga 示例:订单创建流程

[编排器] ──1. 库存扣减──> [库存服务] ──OK──> [编排器]
              失败 → 库存充足失败 → [编排器] 中止

[编排器] ──2. 支付授权──> [支付服务] ──OK──> [编排器]
              失败 → 调用库存回滚 ──> [库存服务]

[编排器] ──3. 创建订单──> [订单服务] ──OK──> [编排器]
              失败 → 调用支付退款 + 库存回滚

[编排器] ──4. 发送通知──> [通知服务]
              失败 → 不必回滚(影响面小,可补偿重试)

9.3 补偿事务

Saga 的关键是补偿操作(Compensating Transaction)。每个正向操作都有一个对应的撤销操作。补偿不是回滚(undo),而是释放正向操作占用的资源。

// NestJS Saga 编排器示例
@Injectable()
export class CreateOrderSaga {
  constructor(
    private inventoryClient: InventoryGrpcClient,
    private paymentClient: PaymentGrpcClient,
    private orderClient: OrderGrpcClient,
    private sagaLog: SagaLogRepository,
  ) {}

  async execute(dto: CreateOrderDto): Promise<string> {
    const sagaId = generateUUID();
    const compensationStack: Compensation[] = [];

    try {
      // Step 1: 扣减库存
      const reservation = await this.inventoryClient.reserveStock({
        productId: dto.productId,
        quantity: dto.quantity,
        sagaId,
      });
      compensationStack.push({
        action: () => this.inventoryClient.releaseReservation(reservation.id),
      });
      await this.sagaLog.recordStep(sagaId, 'inventory_reserved', reservation);

      // Step 2: 支付授权
      const payment = await this.paymentClient.authorize({
        userId: dto.userId,
        amount: dto.amount,
        sagaId,
      });
      compensationStack.push({
        action: () => this.paymentClient.voidAuthorization(payment.id),
      });
      await this.sagaLog.recordStep(sagaId, 'payment_authorized', payment);

      // Step 3: 创建订单
      const order = await this.orderClient.createOrder({
        ...dto,
        reservationId: reservation.id,
        paymentId: payment.id,
        sagaId,
      });
      await this.sagaLog.recordStep(sagaId, 'order_created', order);

      return order.id;

    } catch (err) {
      // 执行补偿(LIFO 顺序)
      while (compensationStack.length > 0) {
        const comp = compensationStack.pop();
        try {
          await comp.action();
        } catch (compErr) {
          // 补偿失败 → 记录待人工介入
          await this.sagaLog.recordCompensationFailure(sagaId, compErr);
        }
      }
      throw new OrderCreationFailedException(sagaId, err.message);
    }
  }
}

Saga 的设计原则:

  1. 幂等性:每个步骤必须支持幂等执行。网络重试可能导致同一 sagaid 的重复调用
  2. 状态机持久化:Saga 状态变更必须持久化到数据库,防止编排器崩溃后丢失进度
  3. 补偿顺序:按正向操作反向顺序执行补偿(LIFO)
  4. 不可补偿的步骤:发邮件、发短信等操作无需补偿,采用"最多一次"语义即可
  5. 监控与告警:补偿失败是需要立即人工介入的严重事件,必须有独立的告警通道

10. NestJS 微服务实战:gRPC 完整示例

NestJS 对微服务的支持非常完善,原生支持 gRPC、RabbitMQ、Kafka、Redis、MQTT 等多种传输层。以下是一个基于 gRPC 的双服务通信示例:订单服务作为客户端调用库存服务。

10.1 项目结构

microservices-demo/
├── proto/
│   ├── order.proto
│   └── inventory.proto
├── order-service/
│   ├── src/
│   │   ├── main.ts
│   │   ├── order.controller.ts
│   │   ├── order.service.ts
│   │   └── app.module.ts
│   └── nest-cli.json
├── inventory-service/
│   ├── src/
│   │   ├── main.ts
│   │   ├── inventory.controller.ts
│   │   ├── inventory.service.ts
│   │   └── app.module.ts
│   └── nest-cli.json
└── docker-compose.yml

10.2 Protocol Buffers 定义

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

package inventory;

service InventoryService {
  rpc CheckStock (CheckStockRequest) returns (StockStatus);
  rpc ReserveStock (ReserveRequest) returns (Reservation);
}

message CheckStockRequest {
  string product_id = 1;
  int32 quantity = 2;
}

message StockStatus {
  string product_id = 1;
  int32 available = 2;
  bool sufficient = 3;
}

message ReserveRequest {
  string product_id = 1;
  int32 quantity = 2;
  string order_id = 3;
}

message Reservation {
  string reservation_id = 1;
  string product_id = 2;
  int32 reserved_quantity = 3;
  int64 expires_at = 4;
}

10.3 库存服务(gRPC 服务端)

// inventory-service/src/inventory.controller.ts
import { Controller } from '@nestjs/common';
import { GrpcMethod } from '@nestjs/microservices';
import { InventoryService } from './inventory.service';

interface CheckStockRequest {
  productId: string;
  quantity: number;
}

interface ReserveRequest {
  productId: string;
  quantity: number;
  orderId: string;
}

@Controller()
export class InventoryController {
  constructor(private readonly inventoryService: InventoryService) {}

  @GrpcMethod('InventoryService', 'CheckStock')
  async checkStock(data: CheckStockRequest) {
    return this.inventoryService.checkAvailability(
      data.productId,
      data.quantity,
    );
  }

  @GrpcMethod('InventoryService', 'ReserveStock')
  async reserveStock(data: ReserveRequest) {
    return this.inventoryService.reserve(
      data.productId,
      data.quantity,
      data.orderId,
    );
  }
}
// inventory-service/src/main.ts
import { NestFactory } from '@nestjs/core';
import { MicroserviceOptions, Transport } from '@nestjs/microservices';
import { join } from 'path';
import { InventoryModule } from './app.module';

async function bootstrap() {
  const app = await NestFactory.createMicroservice<MicroserviceOptions>(
    InventoryModule,
    {
      transport: Transport.GRPC,
      options: {
        package: 'inventory',
        protoPath: join(__dirname, '../proto/inventory.proto'),
        url: '0.0.0.0:50051',
      },
    },
  );

  await app.listen();
  console.log('Inventory gRPC Microservice running on port 50051');
}

bootstrap();

10.4 订单服务(gRPC 客户端 + HTTP 入口)

// order-service/src/app.module.ts
import { Module } from '@nestjs/common';
import { ClientsModule, Transport } from '@nestjs/microservices';
import { join } from 'path';
import { OrderController } from './order.controller';
import { OrderService } from './order.service';

@Module({
  imports: [
    ClientsModule.register([
      {
        name: 'INVENTORY_PACKAGE',
        transport: Transport.GRPC,
        options: {
          package: 'inventory',
          protoPath: join(__dirname, '../proto/inventory.proto'),
          url: process.env.INVENTORY_SERVICE_URL || 'localhost:50051',
        },
      },
    ]),
  ],
  controllers: [OrderController],
  providers: [OrderService],
})
export class OrderModule {}
// order-service/src/order.service.ts
import { Inject, Injectable, OnModuleInit } from '@nestjs/common';
import { ClientGrpc } from '@nestjs/microservices';
import { lastValueFrom } from 'rxjs';

interface InventoryService {
  checkStock(data: { productId: string; quantity: number }): Promise<any>;
  reserveStock(data: {
    productId: string;
    quantity: number;
    orderId: string;
  }): Promise<any>;
}

@Injectable()
export class OrderService implements OnModuleInit {
  private inventoryService: InventoryService;

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

  onModuleInit() {
    this.inventoryService = this.client.getService<InventoryService>('InventoryService');
  }

  async createOrder(productId: string, quantity: number, userId: string) {
    // Step 1: 检查库存
    const stock = await lastValueFrom(
      this.inventoryService.checkStock({ productId, quantity }),
    );

    if (!stock.sufficient) {
      throw new Error(`Insufficient stock for product ${productId}`);
    }

    // Step 2: 预留库存
    const reservation = await lastValueFrom(
      this.inventoryService.reserveStock({
        productId,
        quantity,
        orderId: `order_${Date.now()}`,
      }),
    );

    // Step 3: 创建订单记录(省略实际实现)
    return {
      orderId: reservation.orderId,
      status: 'created',
      reservationId: reservation.reservationId,
    };
  }
}
// order-service/src/order.controller.ts
import { Body, Controller, Post } from '@nestjs/common';
import { OrderService } from './order.service';

class CreateOrderDto {
  productId: string;
  quantity: number;
  userId: string;
}

@Controller('orders')
export class OrderController {
  constructor(private readonly orderService: OrderService) {}

  @Post()
  async createOrder(@Body() dto: CreateOrderDto) {
    return this.orderService.createOrder(
      dto.productId,
      dto.quantity,
      dto.userId,
    );
  }
}

10.5 Docker Compose 部署

version: '3.8'

services:
  inventory-service:
    build: ./inventory-service
    ports:
      - "50051:50051"
    environment:
      - NODE_ENV=production
    healthcheck:
      test: ["CMD", "grpc-health-probe", "-addr=:50051"]
      interval: 10s
      timeout: 5s
      retries: 3

  order-service:
    build: ./order-service
    ports:
      - "3000:3000"
    environment:
      - INVENTORY_SERVICE_URL=inventory-service:50051
      - NODE_ENV=production
    depends_on:
      inventory-service:
        condition: service_healthy

总结

Node.js 微服务架构的落地不是选择一个框架就能解决的问题,而是一系列设计决策的综合结果。从单体到微服务的拆分应以 DDD 限界上下文为依据,避免技术驱动的伪拆分。服务间通信优先选择 REST 对外、gRPC 对内、GraphQL 用于聚合层的组合策略。服务发现优先使用 Kubernetes 内建 DNS,仅在多云混编场景引入 Consul。API Gateway 承担横向关注点但不过度侵入业务,BFF 层为前端定制数据视图但不包含核心业务逻辑。

在弹性设计方面,熔断器保护下游不被过载压垮,带退避和抖动的重试处理瞬时故障,级联超时防止资源无谓等待。可观测性通过 OpenTelemetry + Jaeger 的分布式追踪、结构化日志和指标监控三位一体实现。分布式事务放弃 2PC,采用 Saga 模式以编排或协调方式达成最终一致。

NestJS 对 gRPC、消息队列和微服务架构的一等公民支持,使其成为 Node.js 微服务落地的理想框架。结合本文提供的 proto 定义、服务注册和调用示例,可以在此基础上快速构建生产级的微服务集群。

微服务架构不是目的,而是手段。只有当团队规模、系统复杂度和组织并行度达到了单体无法满足的程度时,微服务才成为自然的选择。保持对复杂度的敬畏,在正确的时间做出正确的架构决策,才是高级工程师的核心能力。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「nodejs」更多文章

  1. Node.js ORM 深度对比:Prisma、TypeORM、Sequelize 与 Drizzle
  2. Node.js 设计模式与最佳实践:从 SOLID 到六边形架构
  3. Node.js 高级测试策略:从单元测试到混沌工程的完整实践