コレオグラフィーとオーケストレーション: マイクロサービスでの分散ワークフローの設計
モノリシック アーキテクチャでは、電子商取引の注文の処理など、複雑なビジネス トランザクションの実行が簡単になります。すべてのデータは単一のリレーショナル データベースに存在するため、開発者は在庫、支払い、配送にわたる複数のデータベースへの書き込みを単一の ACID トランザクション内でラップできます。どこかの時点でエラーが発生した場合、SQL ROLLBACK によってシステムの一貫性が即座に復元されます。
ただし、最新のクラウドネイティブ システムは マイクロサービス アーキテクチャ を採用しており、各サービスがデータを所有し、明確な API 境界を公開しています。この分散パラダイムでは、単一のエンドツーエンドのビジネス操作が複数の独立したマイクロサービスとデータベース エンジンにまたがります。
2 フェーズ コミット (2PC) プロトコルはクラウド ネットワーク全体で遅く、ブロックされ、脆弱であるため、分散システムは 最終整合性を維持しながらワークフローを非同期に調整する必要があります。
これにより、ソフトウェア アーキテクトは、次のような基本的な設計上の決定を迫られます。 分散マイクロサービス ワークフローを管理するには、コレオグラフィーとオーケストレーションを使用する必要がありますか?
この包括的なガイドでは、両方のアーキテクチャ パターンを分析し、現実世界の類似性を探り、アーキテクチャのトレードオフを分析し、Go と Java での実稼働コードの実装を詳しく説明し、インフラストラクチャに適切なパターンを選択するためのフレームワークを確立します。
現実世界の例え: フラッシュ モブ vs. 交響楽団
両方のパターンの直感的なメンタル モデルを構築するには、人間のパフォーマーのグループが自分たちの行動をどのように調整するかを検討してください。
振り付けのたとえ: フラッシュモブ ダンサー ネットワーク
プロのストリート ダンサーのグループがフラッシュ モブ ルーチンを実行しているところを想像してください。ステージに立って個々のダンサーを指差して次の動きを指示するインストラクターはいません。代わりに、各ダンサーは中央の音楽トラックを聴き、隣のダンサーの動きにダイナミックに反応します。
- ダンサー A がフリップを完了すると、ダンサー B はその視覚的な合図を認識し、回転を開始します。 ※ダンサーBが回転し終わるとダンサーCが前に出ます。
- 主な特徴: 分散型、事後対応型、自律型。すべての参加者は、中心的な指示がなくても、自分自身の責任を理解しています。
オーケストレーションのたとえ: 交響楽団
ここで、70 人編成のクラシック交響楽団を想像してください。ヴァイオリニスト、打楽器奏者、チェロ奏者は、ステージ上のお互いの手を見て合図することはありません。代わりに、全員が指揮者を直接見ます。
- 指揮者はヴァイオリンに弱い弦を演奏するタイミングを知らせます。
- 指揮者はドラムを指差し、パーカッションの打撃を合図します。
- 演奏者がテンポを間違えた場合、指揮者がテンポ調整を調整するか、一時停止の合図をします。
- 主な特徴: 集中管理、明示的、コマンド駆動型。一人のリーダーが参加者全員を指揮します。
1. コレオグラフィー アーキテクチャ: 分散型およびイベント駆動型
Choreography では、分散マイクロサービスは中央のマスター コーディネーターなしで反応的に通信します。サービスは、内部状態が変化するたびに、ドメイン イベントを非同期メッセージ ブローカー (Apache Kafka、RabbitMQ、AWS EventBridge など) に発行します。ダウンストリームのマイクロサービスは、関連するイベント トピックをサブスクライブし、次に実行するアクションを個別に決定します。
コレオグラフィーによる電子商取引フロー
コレオグラフィーを使用した電子商取引のチェックアウト プロセスを考えてみましょう。
- 注文サービス: HTTP POST チェックアウト リクエストを受信し、保留中の注文をデータベースに書き込み、
OrderCreatedドメイン イベントを Kafka イベント バスに発行します。 - 支払いサービス:
OrderCreatedトピックを購読します。イベントを受信すると、顧客のクレジット カードに請求し、PaymentProcessedイベントを生成します。 - インベントリ サービス:
PaymentProcessedトピックをサブスクライブします。倉庫アイテムを予約し、InventoryReservedイベントを発行します。 - 配送サービス:
InventoryReservedトピックを購読します。配送ラベルを生成し、OrderShippedイベントを発行します。 - 通知サービス:
OrderShippedを購読し、追跡メールを顧客に送信します。
コレオグラフィーの利点
- 高度な自律性と疎結合: サービスはダウンストリーム ハンドラーの存在を認識しません。 Order Service は、注文が作成されたことのみを認識します。誰がその情報を利用するかは気にしません。
- 独立したスケーラビリティと速度: チームはマイクロサービスを独立して構築、デプロイ、拡張できます。新しい機能 (売り上げを追跡する分析サービスなど) を追加するには、上流のコードを変更せずに既存のイベントをサブスクライブする必要があります。
- 単一障害点 (SPOF) がない: 中央のワークフロー コーディネーターが存在しないため、無関係なサービスの障害によって実行エンジン全体が停止することはありません。
- 高パフォーマンスとスループット: イベント ドリブンのパブリッシュ/サブスクライブ ストリーミングは、同期 HTTP/gRPC ブロッキング レイテンシを発生させることなく、大量のイベントを非同期的に処理します。
コレオグラフィーのデメリット
- 暗黙的なワークフロー ロジック: エンドツーエンドのビジネス プロセスを定義する単一のコードの場所はありません。ワークフロー全体を理解するには、複数のコードベースにわたるイベント ハンドラーをつなぎ合わせる必要があります。
- 循環依存関係のリスク: 慎重なトピック設計を行わずにマイクロサービスが重複するトピックをパブリッシュおよびサブスクライブすると、無限イベント ループによりシステム メッセージ キューがクラッシュする可能性があります。
- 複雑な可観測性と分散トレース: 10 個のイベント トピックにわたる単一の注文トランザクションを追跡するには、堅牢な分散トレース インフラストラクチャ (OpenTelemetry、Jaeger、W3C Trace Context など) が必要です。
- 難しいエラー処理と補償: 支払いが成功した後に Inventory Service が失敗した場合、Inventory Service は
InventoryFailedイベントを発行する必要があります。決済サービスはこのイベントをリッスンし、手動で返金補償をトリガーする必要があります。
2. オーケストレーション アーキテクチャ: 集中型およびコマンド駆動型
オーケストレーション では、専用のコーディネーター サービス (Saga Orchestrator) が実行シーケンスを明示的に指示します。オーケストレーターはワークフロー ステート マシンを保持し、コマンド リクエストを (gRPC、HTTP REST、または専用のコマンド キュー経由で) ワーカー マイクロサービスに送信し、応答を待って、次の実行ステップを決定します。
オーケストレーション下の電子商取引フロー
- Order Service / Saga Orchestrator: チェックアウト リクエストを受信し、
ORDER_PENDING状態のOrderSagaCoordinatorワークフロー インスタンスをインスタンス化します。 - ステップ 1 (支払いコマンド): オーケストレーターは
PaymentService.ExecutePayment()を呼び出します。 Payment Service は支払いを処理し、SUCCESSを返します。 - ステップ 2 (インベントリ コマンド): オーケストレーターは
SUCCESSを受信し、InventoryService.ReserveStock()を呼び出します。 Inventory Service は在庫を予約し、SUCCESSを返します。 - ステップ 3 (出荷コマンド): オーケストレーターは
ShippingService.CreateShipment()を呼び出します。配送サービスは追跡詳細を返します。 - ステップ 4 (完了): オーケストレーターは、状態ストア内の注文ステータスを
ORDER_COMPLETEDに更新します。
ステップ 2 で InventoryService.ReserveStock() が失敗した場合、Orchestrator はロールバック コマンドを順番に実行します。
PaymentService.RefundPayment()を呼び出してステップ 1 を元に戻します。- Saga の状態を
ORDER_CANCELLEDに更新します。
オーケストレーションの利点
- 明示的かつ一元化されたワークフローの可視性: ビジネス プロセス全体が、単一のステート マシン定義またはワークフロー DSL (例: 一時的なワークフロー定義) で明確に表示されます。
- 簡素化された障害管理: ステップが失敗した場合、Orchestrator は間接的なイベント チェーンに依存せずに、以前に完了したすべてのステップの補償トランザクションを直接呼び出します。
- 循環依存関係の防止: ワーカー サービスは、相互に直接呼び出すのではなく、Orchestrator とやり取りします。
- より簡単なテストと監査: 単体テストでサービスの応答を模擬することで、ワークフローの状態遷移を決定論的にテストできます。
オーケストレーションの欠点
- 過剰集中化 (「ゴッド サービス」) のリスク: 開発者がドメイン ビジネス ロジックを Orchestrator にプッシュすると、ワーカー マイクロサービスが「ダムな CRUD サービス」に変わり、モノリシック コアが再作成される危険があります。
- より緊密な API 結合: オーケストレーターは、すべての参加マイクロサービスの API コントラクトとエンドポイントを明示的に認識する必要があります。
- 潜在的なスケーラビリティのボトルネック: 中央のオーケストレーターは、アクティブなトランザクションごとに状態の永続性を処理します。高スループット システムには、水平方向にスケーラブルなステート エンジン バックエンドが必要です。
3. 包括的なアーキテクチャの比較
コレオグラフィーとオーケストレーションを並べて評価するには、主な操作上の特徴を考慮してください。
| 寸法 | コレオグラフィー (イベント駆動型) | オーケストレーション (コマンド駆動) |
|---|---|---|
| コミュニケーション スタイル | 非同期 Pub/Sub (Event ブロードキャスト) |
ポイントツーポイント / RPC (Command + 応答) |
| サービス カップリング | 非常に低い (サービスはドメイン イベントのみを認識します) | 中 (オーケストレーターはワーカー API を知っています) |
| 状態管理 | サービス データベース全体に分散 | 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 実装: コレオグラフィー イベント コンシューマー vs. Saga Orchestrator
1. Go でのコレオグラフィー (Kafka イベント コンシューマー)
Choreography では、Inventory Service は 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 リスナー)
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。
- スキーマ レジストリ: Confluent スキーマ レジストリ (Avro/Protobuf ガバナンス用)。
オーケストレーション エコシステム
- ワークフロー コード エンジン: Temporal.io (Go/Java/TypeScript 永続実行エンジン)、Cadence。
- クラウド マネージド コーディネーター: AWS Step Functions、Azure Logic Apps、GCP ワークフロー。
- BPMN とエンタープライズ エンジン: Camunda 8 (Zeebe)、Netflix の指揮者。
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 日間待つ」) が含まれます。
- コンプライアンスと監査の要件により、すべてのトランザクションの正確な状態を一元的に記録する必要があります。
- カスタム イベント チェーン ロジックを作成せずに、失敗したステップに対する堅牢な自動補正ロールバックが必要です。
- あなたは企業の金融取引 (銀行業務、保険請求処理など) を管理しています。
## 結論
振付もオーケストレーションも普遍的に優れているわけではありません。 Choreography は、暗黙的なワークフローの可視性と複雑な分散トレースを犠牲にして、疎結合、イベント スループット、サービスの自律性を最大化します。 オーケストレーション は、より緊密な API 結合とオーケストレーター インフラストラクチャ管理を犠牲にして、明示的な状態管理、一元的な監査機能、決定的な障害回復を実現します。
両方のパターンの長所を理解し、必要に応じて マイクロ オーケストレーションを使用したハイブリッド マクロ コレオグラフィを活用することで、クラウド環境全体で分散トランザクションを適切に処理する、回復力とスケーラブルなマイクロサービス アーキテクチャを構築できます。
Empower Your Digital Presence & Workflows
Explore top-tier tools built by Ghaznix to streamline your links, surveys, and brand growth.