编排与编排:设计微服务中的分布式工作流程
在整体架构中,执行复杂的业务交易(例如履行电子商务订单)非常简单。所有数据都驻留在单个关系数据库中,允许开发人员将跨库存、支付和运输的多个数据库写入包装在单个 ACID 事务中。如果任何时候发生错误,SQL ROLLBACK 会立即恢复系统一致性。
然而,现代云原生系统采用微服务架构,其中每个服务都拥有自己的数据并公开不同的 API 边界。在这种分布式范例中,单个端到端业务操作跨越多个独立的微服务和数据库引擎。
由于两阶段提交 (2PC) 协议在云网络中速度缓慢、阻塞且脆弱,因此分布式系统必须异步协调工作流程,同时保持最终一致性。
这使软件架构师做出了一个基本的设计决策:您应该使用编排还是编排来管理分布式微服务工作流程?
在这份综合指南中,我们将分解这两种架构模式,探索现实世界的类比,分析架构权衡,详细介绍 Go 和 Java 中的生产代码实现,并建立一个框架来为您的基础设施选择正确的模式。
现实世界的类比:快闪族与交响乐团
要为这两种模式建立直观的心理模型,请考虑人类表演者群体如何协调他们的行动:
编舞类比:快闪舞者网络
想象一下一群专业街舞者表演快闪表演。没有教练站在舞台上指着个别舞者指挥他们的下一步动作。相反,每个舞者都会聆听中心音乐,并对旁边舞者的动作做出动态反应。
- 当舞者 A 完成翻转时,舞者 B 识别出该视觉提示并开始旋转。
- 当舞者B完成旋转后,舞者C向前迈出一步。
- 关键特征:去中心化、反应性和自主性。每个参与者都明白自己的责任,无需中央指导。
编排类比:交响乐团
现在想象一下一个由 70 名乐手组成的古典交响乐团。小提琴手、打击乐手和大提琴手不会通过观察舞台上彼此的手来获得暗示。相反,所有人都直视指挥。
- 指挥向小提琴发出何时演奏软弦的信号。
- 指挥指向鼓以发出打击乐的信号。
- 如果音乐家错过节奏,指挥会协调节奏调整或发出暂停信号。
- 关键特征:集中、明确、命令驱动。一名领导者指导所有参与者。
1.编排架构:去中心化和事件驱动
在编排中,分布式微服务在没有中央主协调器的情况下进行反应式通信。每当服务的内部状态发生变化时,服务都会将域事件发布到异步消息代理(例如 Apache Kafka、RabbitMQ 或 AWS EventBridge)。下游微服务订阅相关事件主题并独立决定下一步采取什么行动。
编排下的电子商务流程
考虑使用编排的电子商务结账流程:
- 订单服务:接收 HTTP POST 结帐请求,将挂单写入其数据库,并向 Kafka 事件总线发出
OrderCreated域事件。 - 支付服务:订阅
OrderCreated主题。收到事件后,它会向客户的信用卡收费并发出PaymentProcessed事件。 - 库存服务:订阅
PaymentProcessed主题。它保留仓库物品并发出InventoryReserved事件。 - 运送服务:订阅
InventoryReserved主题。它生成运输标签并发出OrderShipped事件。 - 通知服务:订阅
OrderShipped并向客户发送跟踪电子邮件。
编排的优点
- 高度自治和松耦合:服务不知道下游处理程序的存在。订单服务只知道订单已创建;它并不关心谁使用该信息。
- 独立的可扩展性和速度:团队可以独立构建、部署和扩展微服务。添加新功能(例如,跟踪销售的分析服务)需要订阅现有事件,而无需修改上游代码。
- 无单点故障(SPOF):因为没有中央工作流协调器,所以不相关服务的故障不会导致整个执行引擎瘫痪。
- 高性能和吞吐量:事件驱动的发布/订阅流异步处理大量事件,没有同步 HTTP/gRPC 阻塞延迟。
编排的缺点
- 隐式工作流逻辑:没有单个代码位置定义端到端业务流程。了解整个工作流程需要将多个代码库中的事件处理程序拼凑在一起。
- 循环依赖风险:如果微服务在没有仔细主题设计的情况下发布和订阅重叠主题,则无限事件循环可能会导致系统消息队列崩溃。
- 复杂的可观测性和分布式跟踪:跨 10 个事件主题跟踪单个订单交易需要强大的分布式跟踪基础设施(例如 OpenTelemetry、Jaeger、W3C Trace Context)。
- 困难的错误处理和补偿:如果付款成功后库存服务失败,则库存服务必须发出
InventoryFailed事件。支付服务必须侦听此事件并手动触发退款补偿。
2.编排架构:集中式和命令驱动
在编排中,专用协调器服务(Saga Orchestrator)显式指导执行顺序。 Orchestrator 持有工作流状态机,向工作微服务发送命令请求(通过 gRPC、HTTP REST 或专用命令队列),等待响应并确定下一个执行步骤。
编排下的电子商务流程
- 订单服务/Saga Orchestrator:接收结帐请求并实例化处于状态
ORDER_PENDING的OrderSagaCoordinator工作流实例。 - 第 1 步(支付命令):Orchestrator 调用
PaymentService.ExecutePayment()。付款服务处理付款并返回SUCCESS。 - 第 2 步(库存命令):Orchestrator 接收
SUCCESS并调用InventoryService.ReserveStock()。库存服务保留库存并退货SUCCESS。 - 第 3 步(运输命令):Orchestrator 调用
ShippingService.CreateShipment()。运输服务返回跟踪详细信息。 - 第 4 步(完成):Orchestrator 将其状态存储中的订单状态更新为
ORDER_COMPLETED。
如果 InventoryService.ReserveStock() 在步骤 2 期间失败,Orchestrator 将按顺序执行回滚命令:
- 调用
PaymentService.RefundPayment()撤消步骤 1。 - 将 Saga 状态更新为
ORDER_CANCELLED。
编排的优点
- 显式且集中的工作流可见性:整个业务流程在单个状态机定义或工作流 DSL(例如,临时工作流定义)中清晰可见。
- 简化的故障管理:如果某个步骤失败,Orchestrator 会直接调用所有先前完成的步骤的补偿事务,而不依赖于间接事件链。
- 防止循环依赖:工作人员服务与 Orchestrator 来回通信,而不是直接相互调用。
- 更轻松的测试和审核:您可以通过在单元测试中模拟服务响应来确定性地测试工作流状态转换。
编排的缺点
- 过度中心化的风险(“上帝服务”):如果开发人员将域业务逻辑推入 Orchestrator,工作微服务可能会变成“愚蠢的 CRUD 服务”,从而重新创建一个整体核心。
- 更紧密的 API 耦合:Orchestrator 必须明确了解所有参与者微服务的 API 契约和端点。
- 潜在的可扩展性瓶颈:中央 Orchestrator 处理每个活动事务的状态持久性。高吞吐量系统需要水平可扩展的状态引擎后端。
3. 综合架构比较
要并行评估编排与编排,请考虑它们的关键操作特征:
| 尺寸 | 编排(事件驱动) | 编排(命令驱动) |
|---|---|---|
| 沟通方式 | 异步发布/订阅(Event 广播) |
点对点/RPC(Command + 响应) |
| 服务耦合 | 非常低(服务只知道域事件) | 中等(Orchestrator 了解工作 API) |
| 状态管理 | 跨服务数据库分布 | 集中在 Orchestrator 状态引擎内部 |
| 工作流程可见性 | 隐式(跨处理程序传播) | 显式(集中状态机代码) |
| 故障恢复 | 复杂(补偿事件级联) | 简单明了(Orchestrator 管理回滚) |
| 分布式追踪 | 需要所有主题的关联 ID | 通过 Orchestrator 日志简化跟踪 |
| 理想的团队规模 | 拥有自治团队的大型工程组织 | 管理复杂企业流程的中型/大型团队 |
| 最适合 | 高通量、简单的线性工作流程 | 具有繁重业务规则的复杂多分支工作流程 |
4. 混合方法:宏观编排+微观编排
现代企业架构很少强制做出全有或全无的选择。相反,领先的工程团队采用混合架构:
- 宏观级别(编排):高级限界上下文(例如,销售限界上下文、供应链限界上下文、客户支持)通过 Kafka 或 NATS 使用 事件驱动编排 进行通信。
- 微观级别(编排):在特定的有界上下文中(例如,在处理多网关重试、欺诈验证和分类帐条目的支付有界上下文内),本地 编排器 协调细粒度的服务执行。
[EVENT BROKER: KAFKA]
/ | \
(OrderCreated) (PaymentSuccess) (StockReserved)
/ | \
[Order Domain] [Payment Domain] [Inventory Domain]
| | |
(Local Saga (Local Saga (Local Saga
Orchestrator) Orchestrator) Orchestrator)
这种混合模式产生跨域边界的编排松散耦合,同时保持各个服务团队内编排的状态可见性。
5. 生产代码示例
让我们看看如何使用 Go 和 Java (Spring Boot) 在生产环境中实现这两种模式。
Go 实现:编排事件消费者与 Saga Orchestrator
1. Go 中的编排(Kafka 事件消费者)
在 Choreography 中,库存服务被动地监听来自 Kafka 的 PaymentProcessedEvent:
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/segmentio/kafka-go"
)
type PaymentProcessedEvent struct {
OrderID string `json:"order_id"`
Amount float64 `json:"amount"`
Status string `json:"status"`
}
type InventoryReservedEvent struct {
OrderID string `json:"order_id"`
Status string `json:"status"`
}
func main() {
reader := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"localhost:9092"},
Topic: "payment-events",
GroupID: "inventory-service-group",
})
defer reader.Close()
writer := kafka.NewWriter(kafka.WriterConfig{
Brokers: []string{"localhost:9092"},
Topic: "inventory-events",
})
defer writer.Close()
fmt.Println("Inventory Service listening for payment events...")
for {
msg, err := reader.ReadMessage(context.Background())
if err != nil {
log.Fatalf("Error reading message: %v", err)
}
var event PaymentProcessedEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
log.Printf("Invalid message payload: %v", err)
continue
}
if event.Status == "SUCCESS" {
log.Printf("[Choreography] Reserved stock for Order: %s", event.OrderID)
// Publish downstream domain event reactively
resEvent := InventoryReservedEvent{
OrderID: event.OrderID,
Status: "RESERVED",
}
payload, _ := json.Marshal(resEvent)
err = writer.WriteMessages(context.Background(), kafka.Message{
Key: []byte(event.OrderID),
Value: payload,
})
if err != nil {
log.Printf("Failed to publish inventory event: %v", err)
}
}
}
}
2. Go 中的编排(中央状态机协调器)
在编排中,显式状态机执行步骤并处理补偿:
package main
import (
"context"
"errors"
"fmt"
"log"
)
type OrderSagaOrchestrator struct {
paymentClient *PaymentClient
stockClient *StockClient
}
func NewOrderSagaOrchestrator(p *PaymentClient, s *StockClient) *OrderSagaOrchestrator {
return &OrderSagaOrchestrator{paymentClient: p, stockClient: s}
}
func (o *OrderSagaOrchestrator) ExecuteSaga(ctx context.Context, orderID string, amount float64) error {
log.Printf("[Orchestrator] Starting Saga execution for Order ID: %s", orderID)
// Step 1: Charge Payment
if err := o.paymentClient.Charge(ctx, orderID, amount); err != nil {
log.Printf("[Orchestrator] Payment failed for Order %s: %v", orderID, err)
return err
}
log.Printf("[Orchestrator] Step 1 Complete: Payment Charged")
// Step 2: Reserve Inventory
if err := o.stockClient.Reserve(ctx, orderID); err != nil {
log.Printf("[Orchestrator] Inventory reservation failed: %v. Initiating Compensation...", err)
// Compensation Step: Refund Payment
if refundErr := o.paymentClient.Refund(ctx, orderID, amount); refundErr != nil {
log.Printf("[CRITICAL] Compensation failed! Manual intervention required for Order %s", orderID)
}
return errors.New("saga aborted: inventory unavailable")
}
log.Printf("[Orchestrator] Saga Completed Successfully for Order ID: %s", orderID)
return nil
}
type PaymentClient struct{}
func (p *PaymentClient) Charge(ctx context.Context, id string, amt float64) error { return nil }
func (p *PaymentClient) Refund(ctx context.Context, id string, amt float64) error { return nil }
type StockClient struct{}
func (s *StockClient) Reserve(ctx context.Context, id string) error { return errors.New("out of stock") }
func main() {
saga := NewOrderSagaOrchestrator(&PaymentClient{}, &StockClient{})
_ = saga.ExecuteSaga(context.Background(), "ORD-9982", 149.99)
}
Java(Spring Boot)实现
1. Java 编排(Spring Cloud Stream / Kafka Listener)
package com.ghaznix.microservices.choreography;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.function.Function;
public record PaymentProcessedEvent(String orderId, String status, double amount) {}
public record InventoryReservedEvent(String orderId, String status) {}
@Configuration
public class InventoryChoreographyProcessor {
@Bean
public Function<PaymentProcessedEvent, InventoryReservedEvent> processPaymentEvent() {
return paymentEvent -> {
System.out.println("[Choreography Java] Processing payment event for order: " + paymentEvent.orderId());
if ("SUCCESS".equals(paymentEvent.status())) {
// Reserve stock in database...
System.out.println("[Choreography Java] Reserved inventory for: " + paymentEvent.orderId());
return new InventoryReservedEvent(paymentEvent.orderId(), "SUCCESS");
} else {
return new InventoryReservedEvent(paymentEvent.orderId(), "FAILED");
}
};
}
}
2. Java 中的编排(声明式状态协调器)
package com.ghaznix.microservices.orchestration;
import org.springframework.stereotype.Service;
@Service
public class OrderSagaOrchestratorService {
private final PaymentServiceClient paymentClient;
private final InventoryServiceClient inventoryClient;
private final ShippingServiceClient shippingClient;
public OrderSagaOrchestratorService(PaymentServiceClient p, InventoryServiceClient i, ShippingServiceClient s) {
this.paymentClient = p;
this.inventoryClient = i;
this.shippingClient = s;
}
public boolean processCheckoutSaga(String orderId, double totalAmount) {
System.out.println("[Orchestrator Java] Initiating Saga Workflow for Order: " + orderId);
// Step 1: Execute Payment
boolean paymentSuccess = paymentClient.processPayment(orderId, totalAmount);
if (!paymentSuccess) {
System.err.println("[Orchestrator Java] Step 1 Failed: Aborting Saga.");
return false;
}
// Step 2: Reserve Inventory
boolean inventorySuccess = inventoryClient.reserveStock(orderId);
if (!inventorySuccess) {
System.err.println("[Orchestrator Java] Step 2 Failed: Triggering Compensation.");
paymentClient.refundPayment(orderId, totalAmount);
return false;
}
// Step 3: Trigger Shipping
boolean shippingSuccess = shippingClient.createShipment(orderId);
if (!shippingSuccess) {
System.err.println("[Orchestrator Java] Step 3 Failed: Compensating Step 2 & Step 1.");
inventoryClient.releaseStock(orderId);
paymentClient.refundPayment(orderId, totalAmount);
return false;
}
System.out.println("[Orchestrator Java] Saga Executed Successfully.");
return true;
}
}
6. 流行的行业工具格局
根据您选择的架构方向,开源和云生态系统提供专用的基础设施引擎:
编排生态系统
- 消息流:Apache Kafka、Apache Pulsar、RabbitMQ、NATS JetStream。
- 云事件路由器:AWS EventBridge、Azure Event Grid、Google Cloud Eventarc。
- 模式注册表:Confluence 模式注册表(用于 Avro/Protobuf 治理)。
编排生态系统
- 工作流代码引擎:Temporal.io(Go/Java/TypeScript 持久执行引擎)、Cadence。
- 云托管协调器:AWS Step Functions、Azure 逻辑应用程序、GCP 工作流程。
- BPMN 和企业引擎:Camunda 8 (Zeebe)、Netflix Conductor。
7.决策矩阵:如何选择?
在为微服务平台选择编排还是编排时,请使用以下决策规则矩阵:
[Start: System Architecture Assessment]
|
Is the workflow complex with >4 steps
or strict business auditing rules?
/ \
(YES) (NO)
/ \
[Choose: Saga Orchestration] Does the system require
(e.g., Temporal / Camunda) ultra-high event streaming velocity?
/ \
(YES) (NO)
/ \
[Choose: Event Choreography] [Choose: Simple Choreography]
(e.g., Apache Kafka / NATS) (e.g., RabbitMQ Pub/Sub)
如果满足以下条件,请选择编排:
- 您的工作流程由 2-4 个简单的线性步骤组成。
- 高事件流吞吐量和亚毫秒级传输延迟是首要任务。
- 您的工程团队被组织成自治域小组,独立构建和部署服务。
- 您已经拥有强大的分布式跟踪和 APM 工具(OpenTelemetry、Datadog)。
如果满足以下条件,请选择编排:
- 您的业务流程涉及复杂的状态转换、多分支条件逻辑或时间延迟(例如“等待 3 天以获得客户批准”)。
- 您的合规性和审计要求要求对每笔交易的准确状态进行集中日志。
- 您需要对失败的步骤进行稳健、自动的补偿回滚,而无需编写自定义事件链逻辑。
- 您正在管理企业财务交易(例如银行业务、保险索赔处理)。
## 结论
编排和编排都不是普遍优越的。 编排最大化松散耦合、事件吞吐量和服务自主性,但代价是隐式工作流可见性和复杂的分布式跟踪。 编排提供显式状态管理、集中式可审计性和确定性故障恢复,但代价是更紧密的 API 耦合和编排器基础设施管理。
通过了解这两种模式的优势,并在适当的情况下利用混合宏编排与微编排,您可以构建弹性、可扩展的微服务架构,以优雅地处理跨云环境的分布式事务。
Empower Your Digital Presence & Workflows
Explore top-tier tools built by Ghaznix to streamline your links, surveys, and brand growth.