通信模式概览
微服务通信模式的选择直接影响系统的可扩展性、可靠性和复杂度。
通信模式分类:
┌─────────────────────────────────────────────────┐
│ 同步通信(请求-响应) │
│ - REST API:简单、通用、HTTP生态 │
│ - gRPC:高性能、强类型、流式支持 │
│ - GraphQL:灵活查询、减少网络请求 │
│ │
│ 异步通信(消息驱动) │
│ - 消息队列:解耦、削峰、可靠投递 │
│ - 事件驱动:松耦合、可扩展、审计追踪 │
│ - 发布订阅:一对多、广播通知 │
│ │
│ 混合模式 │
│ - CQRS:读写分离、优化性能 │
│ - Saga:分布式事务、最终一致性 │
└─────────────────────────────────────────────────┘
同步通信:REST vs gRPC
REST API设计
// services/order-service/routes.js
const express = require('express');
const router = express.Router();
const orderController = require('../controllers/orderController');
const { authenticate, authorize } = require('../middleware/auth');
// 创建订单 - 调用库存服务和支付服务
router.post('/orders', authenticate, async (req, res) => {
try {
const { items, shippingAddress, paymentMethod } = req.body;
// 1. 调用库存服务验证库存
const inventoryCheck = await inventoryClient.checkStock(items);
if (!inventoryCheck.available) {
return res.status(400).json({
error: 'Insufficient stock',
details: inventoryCheck.details
});
}
// 2. 调用价格服务计算总价
const pricing = await pricingClient.calculateTotal(items);
// 3. 创建订单
const order = await orderService.createOrder({
userId: req.user.id,
items,
shippingAddress,
paymentMethod,
totalAmount: pricing.total,
status: 'PENDING'
});
// 4. 调用支付服务
const paymentResult = await paymentClient.processPayment({
orderId: order.id,
amount: pricing.total,
method: paymentMethod
});
if (paymentResult.success) {
// 5. 确认库存扣减
await inventoryClient.reserveStock(order.id, items);
order.status = 'CONFIRMED';
await orderService.updateOrder(order);
} else {
order.status = 'PAYMENT_FAILED';
await orderService.updateOrder(order);
}
res.status(201).json(order);
} catch (error) {
logger.error('Order creation failed', { error: error.message });
res.status(500).json({ error: 'Internal server error' });
}
});
// 服务客户端封装
class InventoryClient {
constructor() {
this.baseUrl = process.env.INVENTORY_SERVICE_URL;
this.circuitBreaker = new CircuitBreaker(this._checkStock.bind(this), {
timeout: 3000,
errorThresholdPercentage: 50,
resetTimeout: 30000
});
}
async checkStock(items) {
return this.circuitBreaker.fire(items);
}
async _checkStock(items) {
const response = await axios.post(`${this.baseUrl}/api/stock/check`, {
items
}, {
timeout: 2000,
headers: {
'X-Request-ID': generateRequestId()
}
});
return response.data;
}
}
gRPC高性能通信
// proto/order.proto
syntax = "proto3";
package order;
service OrderService {
// 一元RPC:创建订单
rpc CreateOrder(CreateOrderRequest) returns (OrderResponse);
// 服务端流:获取订单状态更新
rpc StreamOrderStatus(OrderStatusRequest) returns (stream OrderStatusUpdate);
// 双向流:实时订单处理
rpc ProcessOrderStream(stream OrderRequest) returns (stream OrderResponse);
}
message CreateOrderRequest {
string user_id = 1;
repeated OrderItem items = 2;
Address shipping_address = 3;
PaymentInfo payment_info = 4;
}
message OrderItem {
string product_id = 1;
int32 quantity = 2;
double price = 3;
}
message Address {
string street = 1;
string city = 2;
string country = 3;
string postal_code = 4;
}
message PaymentInfo {
string method = 1;
string token = 2;
}
message OrderResponse {
string order_id = 1;
string status = 2;
double total_amount = 3;
int64 created_at = 4;
}
message OrderStatusRequest {
string order_id = 1;
}
message OrderStatusUpdate {
string order_id = 1;
string status = 2;
string message = 3;
int64 timestamp = 4;
}
// services/order-service/grpc/server.go
package main
import (
"context"
"log"
"net"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
pb "order-service/proto"
)
type OrderServer struct {
pb.UnimplementedOrderServiceServer
orderRepo OrderRepository
inventoryCli *InventoryClient
paymentCli *PaymentClient
}
func (s *OrderServer) CreateOrder(ctx context.Context, req *pb.CreateOrderRequest) (*pb.OrderResponse, error) {
// 1. 验证库存(gRPC调用)
stockResp, err := s.inventoryCli.CheckStock(ctx, &pb.StockCheckRequest{
Items: req.Items,
})
if err != nil {
return nil, status.Errorf(codes.FailedPrecondition,
"库存检查失败: %v", err)
}
if !stockResp.Available {
return nil, status.Error(codes.FailedPrecondition, "库存不足")
}
// 2. 计算价格
totalAmount := calculateTotal(req.Items)
// 3. 创建订单
order := &Order{
UserID: req.UserId,
Items: convertItems(req.Items),
ShippingAddress: convertAddress(req.ShippingAddress),
TotalAmount: totalAmount,
Status: "PENDING",
}
if err := s.orderRepo.Create(ctx, order); err != nil {
return nil, status.Errorf(codes.Internal, "创建订单失败: %v", err)
}
// 4. 处理支付(带超时控制)
paymentCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
paymentResp, err := s.paymentCli.ProcessPayment(paymentCtx, &pb.PaymentRequest{
OrderId: order.ID,
Amount: totalAmount,
Method: req.PaymentInfo.Method,
Token: req.PaymentInfo.Token,
})
if err != nil || !paymentResp.Success {
order.Status = "PAYMENT_FAILED"
s.orderRepo.Update(ctx, order)
return nil, status.Error(codes.FailedPrecondition, "支付失败")
}
// 5. 确认订单
order.Status = "CONFIRMED"
s.orderRepo.Update(ctx, order)
return &pb.OrderResponse{
OrderId: order.ID,
Status: order.Status,
TotalAmount: totalAmount,
CreatedAt: order.CreatedAt.Unix(),
}, nil
}
// 服务端流:实时推送订单状态
func (s *OrderServer) StreamOrderStatus(req *pb.OrderStatusRequest, stream pb.OrderService_StreamOrderStatusServer) error {
orderID := req.OrderId
// 订阅订单状态变更
statusChan := s.subscribeOrderStatus(orderID)
for {
select {
case update := <-statusChan:
if err := stream.Send(&pb.OrderStatusUpdate{
OrderId: orderID,
Status: update.Status,
Message: update.Message,
Timestamp: update.Timestamp.Unix(),
}); err != nil {
return err
}
if update.Status == "COMPLETED" || update.Status == "CANCELLED" {
return nil
}
case <-stream.Context().Done():
return stream.Context().Err()
}
}
}
func main() {
lis, err := net.Listen("tcp", ":50051")
if err != nil {
log.Fatalf("failed to listen: %v", err)
}
s := grpc.NewServer(
grpc.UnaryInterceptor(loggingInterceptor),
grpc.StreamInterceptor(streamLoggingInterceptor),
)
pb.RegisterOrderServiceServer(s, &OrderServer{
orderRepo: NewOrderRepository(),
inventoryCli: NewInventoryClient("inventory-service:50052"),
paymentCli: NewPaymentClient("payment-service:50053"),
})
log.Printf("gRPC server listening on :50051")
if err := s.Serve(lis); err != nil {
log.Fatalf("failed to serve: %v", err)
}
}
异步通信:消息队列与事件驱动
RabbitMQ消息队列
// services/order-service/events/publisher.js
const amqp = require('amqplib');
class OrderEventPublisher {
constructor() {
this.connection = null;
this.channel = null;
this.exchangeName = 'order_events';
}
async connect() {
this.connection = await amqp.connect(process.env.RABBITMQ_URL);
this.channel = await this.connection.createChannel();
// 声明交换机
await this.channel.assertExchange(this.exchangeName, 'topic', {
durable: true
});
console.log('Connected to RabbitMQ');
}
async publishOrderCreated(order) {
const event = {
eventType: 'ORDER_CREATED',
orderId: order.id,
userId: order.userId,
items: order.items,
totalAmount: order.totalAmount,
timestamp: new Date().toISOString()
};
await this.channel.publish(
this.exchangeName,
'order.created',
Buffer.from(JSON.stringify(event)),
{
persistent: true,
contentType: 'application/json',
messageId: generateMessageId(),
timestamp: Date.now()
}
);
console.log(`Published ORDER_CREATED event for order ${order.id}`);
}
async publishOrderPaid(orderId, paymentId) {
const event = {
eventType: 'ORDER_PAID',
orderId,
paymentId,
timestamp: new Date().toISOString()
};
await this.channel.publish(
this.exchangeName,
'order.paid',
Buffer.from(JSON.stringify(event)),
{ persistent: true }
);
}
async publishOrderShipped(orderId, trackingNumber) {
const event = {
eventType: 'ORDER_SHIPPED',
orderId,
trackingNumber,
timestamp: new Date().toISOString()
};
await this.channel.publish(
this.exchangeName,
'order.shipped',
Buffer.from(JSON.stringify(event)),
{ persistent: true }
);
}
}
// services/inventory-service/events/consumer.js
class InventoryEventConsumer {
constructor() {
this.connection = null;
this.channel = null;
this.queueName = 'inventory.order_events';
}
async connect() {
this.connection = await amqp.connect(process.env.RABBITMQ_URL);
this.channel = await this.connection.createChannel();
// 声明队列
await this.channel.assertQueue(this.queueName, {
durable: true,
arguments: {
'x-message-ttl': 86400000, // 24小时TTL
'x-dead-letter-exchange': 'dead_letters'
}
});
// 绑定到交换机
await this.channel.bindQueue(this.queueName, 'order_events', 'order.*');
// 设置预取数量(流量控制)
await this.channel.prefetch(10);
}
async startConsuming() {
console.log('Inventory service started consuming events');
await this.channel.consume(this.queueName, async (msg) => {
if (!msg) return;
try {
const event = JSON.parse(msg.content.toString());
switch (event.eventType) {
case 'ORDER_CREATED':
await this.handleOrderCreated(event);
break;
case 'ORDER_PAID':
await this.handleOrderPaid(event);
break;
case 'ORDER_CANCELLED':
await this.handleOrderCancelled(event);
break;
}
// 确认消息
this.channel.ack(msg);
} catch (error) {
console.error('Error processing message:', error);
// 重试计数
const retryCount = msg.properties.headers['x-retry-count'] || 0;
if (retryCount < 3) {
// 重新入队,增加重试计数
this.channel.nack(msg, false, false);
this.channel.publish('', this.queueName, msg.content, {
headers: { 'x-retry-count': retryCount + 1 }
});
} else {
// 超过重试次数,发送到死信队列
this.channel.nack(msg, false, false);
}
}
});
}
async handleOrderCreated(event) {
console.log(`Reserving stock for order ${event.orderId}`);
// 预留库存
await inventoryService.reserveStock(
event.orderId,
event.items
);
}
async handleOrderPaid(event) {
console.log(`Confirming stock deduction for order ${event.orderId}`);
// 确认扣减库存
await inventoryService.confirmDeduction(event.orderId);
}
async handleOrderCancelled(event) {
console.log(`Releasing reserved stock for order ${event.orderId}`);
// 释放预留库存
await inventoryService.releaseStock(event.orderId);
}
}
Kafka事件流
// services/notification-service/kafka/consumer.js
const { Kafka } = require('kafkajs');
const kafka = new Kafka({
clientId: 'notification-service',
brokers: ['kafka1:9092', 'kafka2:9092', 'kafka3:9092']
});
const consumer = kafka.consumer({
groupId: 'notification-group',
maxWaitTimeInMs: 100,
minBytes: 1,
maxBytes: 10485760 // 10MB
});
async function startConsumer() {
await consumer.connect();
// 订阅多个topic
await consumer.subscribe({
topic: 'order-events',
fromBeginning: false
});
await consumer.subscribe({
topic: 'payment-events',
fromBeginning: false
});
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const event = JSON.parse(message.value.toString());
try {
switch (topic) {
case 'order-events':
await handleOrderEvent(event);
break;
case 'payment-events':
await handlePaymentEvent(event);
break;
}
// 提交offset
await consumer.commitOffsets([{
topic,
partition,
offset: (parseInt(message.offset) + 1).toString()
}]);
} catch (error) {
console.error(`Error processing message from ${topic}:`, error);
// 不提交offset,消息会被重新消费
}
}
});
}
async function handleOrderEvent(event) {
switch (event.eventType) {
case 'ORDER_CREATED':
await sendOrderConfirmationEmail(event);
break;
case 'ORDER_SHIPPED':
await sendShippingNotification(event);
break;
case 'ORDER_DELIVERED':
await sendDeliveryConfirmation(event);
break;
}
}
async function handlePaymentEvent(event) {
switch (event.eventType) {
case 'PAYMENT_SUCCESS':
await sendPaymentReceipt(event);
break;
case 'PAYMENT_FAILED':
await sendPaymentFailureAlert(event);
break;
}
}
startConsumer().catch(console.error);
分布式事务:Saga模式
编排式Saga(Orchestration)
// services/orchestrator/saga/orderSaga.js
class OrderSaga {
constructor() {
this.steps = [];
this.compensations = [];
}
async execute(context) {
try {
// Step 1: 创建订单
const order = await this.createOrder(context);
this.compensations.push(() => this.cancelOrder(order.id));
context.orderId = order.id;
// Step 2: 预留库存
const reservation = await this.reserveInventory(context);
this.compensations.push(() => this.releaseInventory(reservation.id));
context.reservationId = reservation.id;
// Step 3: 处理支付
const payment = await this.processPayment(context);
this.compensations.push(() => this.refundPayment(payment.id));
context.paymentId = payment.id;
// Step 4: 确认订单
await this.confirmOrder(context);
// 所有步骤成功,清理补偿操作
this.compensations = [];
return { success: true, orderId: order.id };
} catch (error) {
console.error('Saga execution failed, starting compensation:', error);
await this.compensate();
throw error;
}
}
async compensate() {
// 逆序执行补偿操作
for (let i = this.compensations.length - 1; i >= 0; i--) {
try {
await this.compensations[i]();
} catch (error) {
console.error(`Compensation step ${i} failed:`, error);
// 记录到死信队列,人工处理
await this.logCompensationFailure(i, error);
}
}
}
async createOrder(context) {
const response = await axios.post(`${ORDER_SERVICE_URL}/orders`, {
userId: context.userId,
items: context.items,
shippingAddress: context.shippingAddress
});
return response.data;
}
async cancelOrder(orderId) {
await axios.patch(`${ORDER_SERVICE_URL}/orders/${orderId}/cancel`);
}
async reserveInventory(context) {
const response = await axios.post(`${INVENTORY_SERVICE_URL}/reservations`, {
orderId: context.orderId,
items: context.items
});
return response.data;
}
async releaseInventory(reservationId) {
await axios.delete(`${INVENTORY_SERVICE_URL}/reservations/${reservationId}`);
}
async processPayment(context) {
const response = await axios.post(`${PAYMENT_SERVICE_URL}/payments`, {
orderId: context.orderId,
amount: context.totalAmount,
method: context.paymentMethod
});
return response.data;
}
async refundPayment(paymentId) {
await axios.post(`${PAYMENT_SERVICE_URL}/payments/${paymentId}/refund`);
}
async confirmOrder(context) {
await axios.patch(`${ORDER_SERVICE_URL}/orders/${context.orderId}/confirm`);
}
}
// 使用示例
const saga = new OrderSaga();
const context = {
userId: 'user123',
items: [{ productId: 'p1', quantity: 2 }],
shippingAddress: { /* ... */ },
paymentMethod: 'credit_card',
totalAmount: 100.00
};
saga.execute(context)
.then(result => console.log('Order completed:', result))
.catch(error => console.error('Order failed:', error));
协同式Saga(Choreography)
// services/order-service/events/handlers.js
class OrderEventHandler {
async handleInventoryReserved(event) {
const { orderId, reservationId } = event;
// 更新订单状态
await orderService.updateStatus(orderId, 'INVENTORY_RESERVED');
// 发布事件触发下一步:处理支付
await eventBus.publish('payment.initiate', {
orderId,
amount: await orderService.getTotalAmount(orderId)
});
}
async handleInventoryReservationFailed(event) {
const { orderId, reason } = event;
// 取消订单
await orderService.cancelOrder(orderId, reason);
// 通知用户
await notificationService.sendOrderFailedNotification(orderId, reason);
}
}
// services/payment-service/events/handlers.js
class PaymentEventHandler {
async handlePaymentInitiate(event) {
const { orderId, amount } = event;
try {
const payment = await paymentService.processPayment({
orderId,
amount
});
if (payment.success) {
await eventBus.publish('payment.completed', {
orderId,
paymentId: payment.id
});
} else {
await eventBus.publish('payment.failed', {
orderId,
reason: payment.error
});
}
} catch (error) {
await eventBus.publish('payment.failed', {
orderId,
reason: error.message
});
}
}
}
// services/inventory-service/events/handlers.js
class InventoryEventHandler {
async handlePaymentCompleted(event) {
const { orderId, paymentId } = event;
// 确认库存扣减
await inventoryService.confirmDeduction(orderId);
// 发布事件
await eventBus.publish('inventory.confirmed', {
orderId,
paymentId
});
}
async handlePaymentFailed(event) {
const { orderId, reason } = event;
// 释放预留库存
await inventoryService.releaseReservation(orderId);
// 发布事件
await eventBus.publish('inventory.released', {
orderId,
reason
});
}
}
通信模式选择指南
决策树:
┌─────────────────────────────────────────────────┐
│ 需要立即响应? │
│ ├─ 是 → 同步通信 │
│ │ ├─ 简单CRUD → REST │
│ │ ├─ 高性能/流式 → gRPC │
│ │ └─ 灵活查询 → GraphQL │
│ │ │
│ └─ 否 → 异步通信 │
│ ├─ 需要可靠投递 → 消息队列(RabbitMQ) │
│ ├─ 高吞吐/持久化 → Kafka │
│ └─ 松耦合/审计 → 事件驱动 │
│ │
│ 需要跨服务事务? │
│ ├─ 强一致性 → 避免分布式事务,重构设计 │
│ └─ 最终一致性 → Saga模式 │
│ ├─ 集中控制 → 编排式Saga │
│ └─ 松耦合 → 协同式Saga │
└─────────────────────────────────────────────────┘
总结
微服务通信模式的选择应基于业务需求和技术约束:
- 同步通信:适合需要实时响应的场景,但要考虑超时、重试、熔断
- 异步通信:适合解耦、削峰、提高可靠性,但要处理最终一致性
- Saga模式:处理分布式事务,权衡编排式和协同式的优缺点
- 混合使用:根据场景选择最合适的模式,不必拘泥于单一方案
关键原则:
- 优先使用异步通信,减少服务间耦合
- 同步调用必须有超时、重试、熔断机制
- 使用幂等性设计,处理消息重复消费
- 实现分布式追踪,便于问题排查
- 监控通信延迟和失败率,及时发现问题
延伸阅读
- Microservices Patterns - Chris Richardson
- Building Microservices - Sam Newman
- gRPC官方文档
- Apache Kafka文档
- Saga Pattern - Microsoft
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。