トランザクション管理@マイクロサービス¶
はじめに¶
本サイトにつきまして、以下をご認識のほど宜しくお願いいたします。
01. Shared DB パターンのトランザクションパターン¶
- パターン不要 (各マイクロサービスの従来のトランザクション)
02. DB per service パターンのトランザクションパターン¶
ローカルトランザクション¶
▼ ローカルトランザクションとは¶
1 個のトランザクション処理によって、1 個のマイクロサービスの DB やスキーマ (MySQL の文脈ではスキーマが DB に相当) を操作する。
この方法が推奨である。
マイクロサービスアーキテクチャでローカルトランザクション処理を使用する場合、これを連続的に実行する仕組みが必要になる。
また、これらの各 DB に対する各トランザクション処理を紐付けられるように、トランザクションに ID (例:UUID) を割り当てる必要がある。
▼ ローカルトランザクション処理を実装できるパターン¶
ローカルトランザクション処理を連鎖させる必要がある場合、Saga パターン (オーケストレーションベース、コレグラフィベース、並列パイプラインベース) も採用する必要がある。
一方で、連鎖させる必要がない場合は特に名前がないが、便宜上『連鎖不要のローカルトランザクション』とする。
グローバルトランザクション (分散トランザクション)¶
▼ グローバルトランザクションとは¶
分散トランザクションとも言う。
1 個のトランザクションが含む各 CRUD 処理を異なる DB に対して実行する。
トランザクション中の各 CRUD で、宛先の DB を動的に切り替える必要がある。
非推奨である。
▼ グローバルトランザクション処理を実装できるパターン¶
グローバルトランザクション処理を実装できるパターンとして、二相コミット (2フェーズコミット) がある。
03. 2フェーズコミットパターン (二相コミットパターン)¶
2フェーズコミットパターンとは¶
『二相コミットパターン』ともいう。
非推奨である。
実装パターン¶
▼ 自前で実装する場合¶
マイクロサービスの処理のなかで実装する。
▼ OSS を使用する場合¶
二相コミットの OSS はなさそう。
▼ クラウドプロバイダーのマネージドサービスを使用する場合¶
- Scalar DB
- Google Spanner
04. 連鎖不要なローカルトランザクション¶
各マイクロサービスの永続化の間に依存関係がない場合、これらのマイクロサービスの永続化を調整する必要はない。
05. Saga パターン(連鎖が必要なローカルトランザクションの場合)¶
Saga パターンとは¶
各マイクロサービスの永続化の間に依存関係がある場合 (例:受注データの永続化には、配送データや決済データの永続化の結果が必要) に、これらのマイクロサービスの永続化を調整する必要がある。
Saga オーケストレーターにコールされるマイクロサービスへ、永続化とロールバックに関する API を実装する。
Saga オーケストレーターは、これらのマイクロサービスをコールし、ローカルトランザクション処理を連続的に実行する。

- https://iorilan.medium.com/i-asked-this-system-design-question-to-3-guys-during-a-developer-interview-and-none-of-them-gave-9c23abe45687
- https://thinkit.co.jp/article/14639?page=0%2C1
- https://qiita.com/nk2/items/d9e9a220190549107282
- https://qiita.com/yasuabe2613/items/b0c92ab8c45d80318420
- https://github.com/yongk/orderdemo?tab=readme-ov-file#bounded-context-mappings
Saga パターンと ACID¶
デザインパターン¶
▼ State パターン¶
Saga オーケストレーターをステートマシン図や State パターンでモデリングし、ステートマシンを実装する。
05-02. オーケストレーションベースの Saga パターン¶
オーケストレーションベースの Saga パターンとは¶
一連のローカルトランザクションの実行をまとめて制御する責務を持った Saga オーケストレーター (コーディネーター) と、これをコールする別のマイクロサービスを配置する。
各マイクロサービス間の通信方式は、リクエスト/レスポンス方式またはパブリッシュ/サブスクライブ方式のどちらでもよい。

- https://learn.microsoft.com/ja-jp/azure/architecture/reference-architectures/saga/saga
- https://blogs.itmedia.co.jp/itsolutionjuku/2019/08/post_729.html
- https://medium.com/google-cloud-jp/gcp-saga-microservice-7c03a16a7f9d
- https://www.fiorano.com/jp/blog/integration/integration-architecture/%E3%82%B3%E3%83%AC%E3%82%AA%E3%82%B0%E3%83%A9%E3%83%95%E3%82%A3-vs-%E3%82%AA%E3%83%BC%E3%82%B1%E3%82%B9%E3%83%88%E3%83%AC%E3%83%BC%E3%82%B7%E3%83%A7%E3%83%B3/
実装パターン (リクエスト/レスポンス方式、パブリッシュ/サブスクライブ方式、ストリーミング方式を区別していない)¶
▼ 自前で実装する場合¶
マイクロサービスの処理のなかで実装する。
*実装例*
Saga オーケストレーターと、これをコールする別のマイクロサービスを配置する。
Saga オーケストレーターは、各ローカルトランザクションの成否を表す Saga ログを DB で管理する。
Saga オーケストレーターは、Order サービス (T1) 、Inventory サービス (T2) 、Payment サービス (T3) のローカルトランザクション処理を連続して実行する。
例えば、Payment サービスのローカルトランザクション (T3) が失敗した場合、Order サービスと Payment サービスのローカルトランザクション処理をロールバックする補償トランザクション (C1、C2) を実行する。

- https://docs.aws.amazon.com/prescriptive-guidance/latest/cloud-design-patterns/saga-orchestration.html#saga-orchestration-implementation
- https://dzone.com/articles/modelling-saga-as-a-state-machine
- https://www.baeldung.com/cs/saga-pattern-microservices
- https://medium.com/@vinciabhinav7/saga-design-pattern-569ec942079
- https://blog.knoldus.com/distributed-transactions-and-saga-patterns/
▼ OSS を使用する場合¶
Saga オーケストレーターの OSS (Temporal、Netflix Conductor、Uber Cadence など) を使用する。
Saga オーケストレーターのドメインモデリングにステートソーシングパターンを採用する必要がある。
▼ クラウドプロバイダーのマネージドサービスを使用する場合¶
クラウドプロバイダー (例:AWS、Google Cloud) が提供する Saga オーケストレーター (例:AWS Step Functions、Google Workflows など) を使用する。
各マイクロサービスは、Saga オーケストレーターをメッセージブローカー (これもクラウドプロバイダーが提供しているものでよい) を介して取得する。
Saga オーケストレーター¶
Saga オーケストレーターは、Saga オーケストレーターのワーカー(マイクロサービス)を調整し、ローカルトランザクション処理を連鎖的に実行する。
また、ローカルトランザクションの進捗度 (Saga ログ) を DB に都度記録し、Saga ログからいずれのローカルトランザクション処理を実行し終えたかを判断する。
途中で失敗した場合、失敗以前のローカルトランザクションの結果を元に戻す補償トランザクション処理を実行する。
Saga オーケストレーターのストレージ¶
▼ インメモリ¶
Saga オーケストレーターのインメモリでローカルトランザクションの進捗度 (Saga ログ) をリアルタイムに管理する。
Saga オーケストレーター自体で障害が起こった場合は進捗度が消えてしまうデメリットがあるが、Saga オーケストレーターの実装が簡単になる。
代わりに、すべて成功したか、どこかで失敗したか、キャンセルがあったかだけを DB に記録する。
▼ DB¶
DB でローカルトランザクションの進捗度 (Saga ログ) をリアルタイムに管理する。
通常のオーケストレーションベースの Saga パターンでは、DB に Saga ログテーブルを作成する。
Saga オーケストレーターは、ローカルトランザクションの進捗度 (Saga ログ) を DB に永続化する。
Saga オーケストレーターごとに DB を分割するとよい。
AWS StepFunctions のステートも設計例として、参考になる。
id |
order_saga_execution_id |
order_saga_current_step |
order_id |
order_saga_payload |
order_saga_status |
order_saga_state |
order_saga_version |
start_data |
end_data |
|---|---|---|---|---|---|---|---|---|---|
| 1 | 9db5b6da-daba-4633-b3cf-9c79f2bcf6f5 |
CreditApproval | 1 | "\"order-id\": 1, \"customer-id\": 456, \"payment-due\": 4999, \"credit-card-no\": \"xxxx-yyyy-dddd-9999\"}" |
SUCCEEDED | "{\"creditApproval\":\"SUCCEEDED\"}" |
楽観的ロックに使用するバージョン値 (例:最終更新日など) | 開始時刻 | 終了時刻 |
| 2 | b1f14b72-393d-432b-8ec2-782974a6ed60 |
Payment | 1 | "{ \"order-id\": 1, \"customer-id\": 456, ... }" |
STARTED | "{\"creditApproval\":\"SUCCEEDED\",\"payment\":\"STARTED\"}" |
〃 | 〃 | 〃 |
| 3 | b38229c6-30df-4166-a725-8b2c578e5ed5 |
CreditApproval | 2 | "{ \"order-id\": 2, \"customer-id\": 456, ... }" |
STARTED | "{\"creditApproval\":\"STARTED\"}" |
〃 | 〃 | 〃 |
| ... | ... | ... | ... | ... | ... | ... | ... | ... | ... |
Saga オーケストレーターのステータスチェッカー¶
▼ Saga オーケストレーターのステータスチェッカーとは¶
オーケストレーションベースの Saga パターンにて、Saga オーケストレーターにリクエストを送信するクライアントは、Saga オーケストレーターの処理結果を知る必要がある。
▼ ポーリングパターンの場合¶
ポーリングパターンの場合、Saga ステータスチェッカーを採用する。
Saga ステータスチェッカーは、トランザクション ID を使用してワークフローの処理結果を DB から読み込む。

Saga オーケストレーターのワーカー¶
▼ メッセージキューを経由しないポイントツーポイントの場合¶
メッセージブローカー (例:Apache Kafka、RabbitMQ など) を経由しないオーケストレーションベースの Saga パターンを実装する。
メッセージブローカーを経由するよりも、Saga オーケストレーターと各ワーカーの間の結合度が高まってしまうが、Saga オーケストレーターの実装が簡単になる。
ローカルトランザクションの進捗度に応じて、次のローカルトラザクションや補償トランザクション処理を実行する。
▼ メッセージキューを経由する場合¶
メッセージキュー (例:Amazon SQS など) を経由して、Saga オーケストレーターとワーカーの間で通信する。
▼ メッセージブローカーを経由する場合¶
メッセージブローカー (例:Apache Kafka、RabbitMQ など) を使い、オーケストレーションベースの Saga パターンを実装する。
Saga オーケストレーターは、メッセージブローカーに対してパブリッシュとサブスクライブを実行する。
サブスクライブしたメッセージに応じて、次のメッセージをパブリッシュする。

補償トランザクション¶
▼ 補償トランザクションとは¶
ローカルトランザクション処理を元に戻すトランザクション処理を逆順に実行し、Saga パターンによるトランザクションの結果を元に戻す仕組みのこと。
マイクロサービスアーキテクチャでは、トランザクションの通常のロールバック機能を使用した場合、処理に失敗したマイクロサービスのみロールバックする。その結果、それ以前のマイクロサービスではロールバックが起こらない可能性もある。
いずれかのマイクロサービスのローカルトランザクションが失敗した場合、まずそのマイクロサービスは自身のトランザクション処理をロールバックする。
その後、それまでのローカルトランザクション処理を擬似的にロールバックするトランザクション処理を逆順で実行する。
▼ 設計例¶
受注に関するトランザクションが異なるマイクロサービスにまたがる例。

補償トランザクションによって、各ローカルトランザクション処理を元に戻す逆順のトランザクション処理を実行する。

▼ 実装例 (Go の defer() 関数)¶
この例では、Go の defer() 関数で補償トランザクションの仕組みを実装している。
ローカルトランザクションで失敗した場合は、まずそのマイクロサービスが自身のトランザクション処理をロールバックする。
その後、それまでにコールされた defer() 関数を実行し補償トランザクション処理を実行する。
package saga
import (
"time"
"go.uber.org/multierr"
"go.temporal.io/sdk/temporal"
"go.temporal.io/sdk/workflow"
)
func TransferMoney(ctx workflow.Context, transferDetails TransferDetails) (err error) {
retryPolicy := &temporal.RetryPolicy{
InitialInterval: time.Second,
BackoffCoefficient: 2.0,
MaximumInterval: time.Minute,
MaximumAttempts: 3,
}
options := workflow.ActivityOptions{
StartToCloseTimeout: time.Minute,
RetryPolicy: retryPolicy,
}
ctx = workflow.WithActivityOptions(ctx, options)
err = workflow.ExecuteActivity(ctx, Withdraw, transferDetails).Get(ctx, nil)
if err != nil {
return err
}
// 補償トランザクション
defer func() {
if err != nil {
errCompensation := workflow.ExecuteActivity(ctx, WithdrawCompensation, transferDetails).Get(ctx, nil)
err = multierr.Append(err, errCompensation)
}
}()
// ローカルトランザクション
// 失敗した場合、まずは自身のトランザクション処理をロールバックする
// その後、前のdefer関数を実行し、前のローカルトランザクション処理を元に戻す補償トランザクション処理を実行する
err = workflow.ExecuteActivity(ctx, Deposit, transferDetails).Get(ctx, nil)
if err != nil {
return err
}
// 補償トランザクション
defer func() {
if err != nil {
errCompensation := workflow.ExecuteActivity(ctx, DepositCompensation, transferDetails).Get(ctx, nil)
err = multierr.Append(err, errCompensation)
}
// uncomment to have time to shut down worker to simulate worker rolling update and ensure that compensation sequence preserves after restart
// workflow.Sleep(ctx, 10*time.Second)
}()
// ローカルトランザクション
// 失敗した場合、まずは自身のトランザクション処理をロールバックする
// その後、前のdefer関数を実行し、前のローカルトランザクション処理を元に戻す補償トランザクション処理を実行する
err = workflow.ExecuteActivity(ctx, StepWithError, transferDetails).Get(ctx, nil)
if err != nil {
return err
}
return nil
}
▼ 実装例 (Go の slice)¶
この例では、Go の slice で補償トランザクションの仕組みを実装している。
スライス内のローカルトランザクション処理を順番に実行し、どこかで失敗した場合は逆順に補償トランザクション処理を実行する。
package main
import (
"fmt"
"errors"
)
// ローカルトランザクション処理を表Cookieす関数型
type LocalTransaction func() error
// 補償トランザクション処理を表す関数型
type CompensatingAction func() error
// Sagaオーケストレーターの各ステップを表す構造体
type SagaStep struct {
Transaction LocalTransaction // ローカルトランザクション
Compensate CompensatingAction // 補償トランザクション
}
// Sagaを表す構造体
type Saga struct {
Steps []SagaStep // 複数のSagaStepから成る
}
// ローカルトランザクションと補償トランザクション処理を実行する関数
func (s *Saga) Execute() error {
for _, step := range s.Steps {
// ローカルトランザクション処理を順番に実行する
if err := step.Transaction(); err != nil {
// 失敗した場合は、補償トランザクション処理を逆順で実行する
for i := len(s.Steps) - 1; i >= 0; i-- {
if err := s.Steps[i].Compensate(); err != nil {
// 補償トランザクションが失敗した場合はエラーメッセージを返す
return errors.New(fmt.Sprintf("failed to compensate for step %d: %v", i, err))
}
}
// 最初のエラーを返す
return err
}
}
// 全てのトランザクションが成功した場合、nilを返す
return nil
}
// ローカルトランザクション
// 資金移動
func transferFunds() error {
return nil
}
// 補償トランザクション
// 資金移動の取り消し
func reverseTransfer() error {
return nil
}
func main() {
// Sagaオーケストレーター
saga := Saga{
Steps: []SagaStep{
SagaStep{
Transaction: transferFunds, // 1つ目のローカルトランザクション
Compensate: reverseTransfer, // 1つ目の補償トランザクション
},
SagaStep{
Transaction: transferFunds, // 2つ目のローカルトランザクション
Compensate: reverseTransfer, // 2つ目の補償トランザクション
},
},
}
// Sagaの実行
if err := saga.Execute(); err != nil {
fmt.Println("saga failed:", err) // エラーが発生した場合
} else {
fmt.Println("saga succeeded") // 正常に完了した場合
}
}
- https://dsysd-dev.medium.com/writing-temporal-workflows-in-golang-part-1-9f50f6ef23d5
- https://qiita.com/somen440/items/a6c323695627235128e9#%E3%82%AA%E3%83%BC%E3%82%B1%E3%82%B9%E3%83%88%E3%83%AC%E3%83%BC%E3%82%B7%E3%83%A7%E3%83%B3%E3%83%99%E3%83%BC%E3%82%B9%E3%81%AE%E3%82%B5%E3%83%BC%E3%82%AC%E5%AE%9F%E8%A3%85
▼ 実装例 (TypeScript の配列)¶
この例では、Azure の Durable Function にて、TypeScript の配列で補償トランザクションの仕組みを実装している。
スライス内のローカルトランザクション処理を順番に実行し、どこかで失敗した場合は逆順に補償トランザクション処理を実行する。
import df from "durable-functions";
import {Task} from "durable-functions/lib/src/classes";
// APIError型の定義。ステータスコードとボディを持つ
type APIError = {
status: 200 | 400 | 500;
body: object | string;
};
// APIErrorかどうかをチェックする関数
const isAPIError = (arg: any): arg is APIError => {
// 引数がオブジェクトでない場合は、トランザクション不可とする
if (typeof arg !== "object") return false;
// ステータスコードが200, 400, 500のいずれかでない場合は、トランザクション不可とする
if (!(arg.status && [200, 400, 500].includes(arg.status))) return false;
// メッセージが文字列でない場合は、トランザクション不可とする
if (typeof arg.message !== "string") return false;
// 全ての条件を満たす場合はトランザクション可とする
return true;
};
// Sagaオーケストレーター
export const saga = df.orchestrator(function* (context) {
// 補償トランザクション処理を格納する配列
const compensatingTransactions: Task[] = [];
try {
// Sagaオーケストレーターの入力を取得
const {input} = context.df.getInput();
// ローカルトランザクションとして、doActivityAを実行する
const a = yield context.df.callActivity("doActivityA", input.body);
// 補償トランザクションとして、rejectActivityAを追加する
compensatingTransactions.push(
context.df.callActivity("rejectActivityA", input.body),
);
// ローカルトランザクションとして、doActivityBを実行する
const b = yield context.df.callActivity("doActivityB", a);
// 補償トランザクションとして、rejectActivityBを追加する
compensatingTransactions.push(
context.df.callActivity("rejectActivityB", b),
);
// ローカルトランザクションとして、doActivityCを実行する
const c = yield context.df.callActivity("doActivityC", b);
// 補償トランザクションとして、rejectActivityCを追加する
compensatingTransactions.push(
context.df.callActivity("rejectActivityC", c),
);
// Sagaオーケストレーターのクライアントに正常終了のレスポンスを返す
return {
status: 200,
body: "The process has succeeded.",
};
} catch (e) {
// 例外発生時に補償トランザクション処理をまとめて実行する
yield context.df.Task.all(compensatingTransactions);
// 例外がAPIError型の場合、そのまま返す
if (isAPIError(e)) return e;
// その他の例外は500エラーとして、返す
return {
status: 500,
body: (e as Error).message,
};
}
});
▼ 実装例 (Go の slice)¶
この例では、アウトボックスパターンで Saga オーケストレーションを実装している。
ちょっと難しいかな...
package saga
import (
"context"
"database/sql"
"fmt"
"github.com/google/uuid"
"go.example/saga/pkg/jsonmap"
)
type SagaState struct {
ID uuid.UUID
Version int8
Type string
Payload jsonmap.JSONMap
CurrentStep SagaStep
StepStatus jsonmap.JSONMap
SagaStatus SagaStatus
}
// Repository
type Repository interface {
Persist(ctx context.Context, tx *sql.Tx, ss SagaState) error
Update(ctx context.Context, tx *sql.Tx, ss SagaState) error
QueryByID(ctx context.Context, tx *sql.Tx, ID string) (*SagaState, error)
}
func NewSaga(sagaType string, payload jsonmap.JSONMap, currentStep SagaStep) SagaState {
// ステートマシン
return SagaState{
ID: uuid.New(),
Version: 1,
Type: sagaType,
Payload: payload,
CurrentStep: currentStep,
StepStatus: jsonmap.JSONMap{string(currentStep): SagaStepStatusStarted},
SagaStatus: SagaStatusStarted,
}
}
// NextSagaStatus evaluate current SagaStepStatuses and set SagaStatus
func (s *SagaState) NextSagaStatus() {
ss := map[string]bool{}
for _, v := range s.StepStatus {
ss[fmt.Sprintf("%v", v)] = true
}
if ss[SagaStepStatusSucceeded] && len(ss) == 1 {
s.SagaStatus = SagaStatusCompleted
} else if (ss[SagaStepStatusStarted] && len(ss) == 1) || (ss[SagaStepStatusSucceeded] && ss[SagaStepStatusStarted] && len(ss) == 2) {
s.SagaStatus = SagaStatusStarted
} else if !ss[SagaStepStatusCompensating] {
s.SagaStatus = SagaStatusAborted
} else {
s.SagaStatus = SagaStatusAborting
}
}
// IncrementVersion
func (s *SagaState) IncrementVersion() {
s.Version++
}
// SagaStatus represents the saga status based on steps status
type SagaStatus string
// SagaStatus type
const (
SagaStatusStarted = "STARTED"
SagaStatusAborting = "ABORTING"
SagaStatusAborted = "ABORTED"
SagaStatusCompleted = "COMPLETED"
)
// SagaStepStatus represent current saga step status
type SagaStepStatus string
// SagaStepStatus type
const (
SagaStepStatusStarted = "STARTED"
SagaStepStatusFailed = "FAILED"
SagaStepStatusSucceeded = "SUCCEEDED"
SagaStepStatusCompensating = "COMPENSATING"
SagaStepStatusCompensated = "COMPENSATED"
)
// SagaStep define saga service step in order to follow
type SagaStep string
// NextSagaStep find saga next step from provided steps and current saga step
func NextSagaStep(steps []SagaStep, currentStep SagaStep) SagaStep {
if currentStep == "" {
return steps[0]
}
curr := -1
for i := 0; i < len(steps); i++ {
if steps[i] == currentStep {
curr = i
break
}
}
if curr == -1 || curr+1 == len(steps) {
return ""
}
return steps[curr+1]
}
// PrevSagaStep find saga previous step from provided steps and current saga step
func PrevSagaStep(steps []SagaStep, currentStep SagaStep) SagaStep {
curr := -1
for i := 0; i < len(steps); i++ {
if steps[i] == currentStep {
curr = i
break
}
}
if curr == -1 || curr-1 == -1 {
return ""
}
return steps[curr-1]
}
package reservation
import (
"context"
"database/sql"
"fmt"
"go.example/saga/pkg/saga"
"go.example/saga/pkg/store/postgres"
"go.example/saga/reservation/pkg/model"
"log"
)
...
func (c *Controller) PostReservation(ctx context.Context, cmd model.ReservationCmd) (*model.Reservation, error) {
r := model.NewReservation(cmd.HotelID, cmd.RoomID, cmd.GuestID, cmd.PaymentDue, cmd.StartDate, cmd.EndDate, cmd.CreditCardNO)
if _, err := c.store.Transact(ctx, func(tx *sql.Tx) (interface{}, error) {
// persist reservation
if err := c.repository.Add(ctx, tx, r); err != nil {
return nil, err
}
payload := r.ToJSONMap()
currStep := saga.NextSagaStep(sagaSteps, "")
sagaState := saga.NewSaga(roomReservationSaga, payload, currStep)
if err := c.sagaRepository.Persist(ctx, tx, sagaState); err != nil {
return nil, err
}
outboxEvent := postgres.NewEvent(sagaState.ID.String(), string(currStep), postgres.RequestEventType, payload)
if err := outboxEvent.Persist(ctx, tx); err != nil {
return nil, err
}
log.Printf("Started Saga for reservationID %s sagaID %s", r.ID, sagaState.ID)
return r, nil
}); err != nil {
return nil, err
}
return r, nil
}
...
05-03. Choreography (コレオグラフィ) ベースの Saga パターン¶
コレオグラフィベースの Saga パターンとは¶
マイクロサービスは、自身のローカルトランザクション処理を完了させた後に、次のマイクロサービスをコールする。
各マイクロサービス間の通信方式は、パブリッシュ/サブスクライブ方式にする必要がある。
そのために、マイクロサービス間にメッセージ中継システムを配置する。

- https://learn.microsoft.com/ja-jp/azure/architecture/reference-architectures/saga/saga
- https://blogs.itmedia.co.jp/itsolutionjuku/2019/08/post_729.html
- https://zenn.dev/yoshii0110/articles/74dfcf4132a805
- https://www.fiorano.com/jp/blog/integration/integration-architecture/%E3%82%B3%E3%83%AC%E3%82%AA%E3%82%B0%E3%83%A9%E3%83%95%E3%82%A3-vs-%E3%82%AA%E3%83%BC%E3%82%B1%E3%82%B9%E3%83%88%E3%83%AC%E3%83%BC%E3%82%B7%E3%83%A7%E3%83%B3/
実装パターン (パブリッシュ/サブスクライブ方式)¶
▼ 自前で実装する場合¶
マイクロサービスの処理のなかで実装する。
*実装例*
以下のリポジトリを参考にせよ。

▼ OSS を使用する場合¶
コレオグラフィの OSS (例:Debezium、Maxwell など) を使用する。
Saga オーケストレーターのドメインモデリングにイベントソーシングパターンを採用する必要がある。
▼ クラウドプロバイダーのマネージドサービスを使用する場合¶
各マイクロサービス間の通信方式は、パブリッシュ/サブスクライブ方式にする必要がある。
- AWS Lambda、マイクロサービス間のパブリッシュ/サブスクライブ方式の AWS リソース (例:Amazon EventBridge、Amazon SQS、Amazon SNS)
- Google Cloud Run Functions、マイクロサービス間のパブリッシュ/サブスクライブ方式の Google Cloud リソース (例:Google Eventarc)
宛先マイクロサービスとの通信方式¶
各マイクロサービスにパブリッシュとサブスクライブを処理する責務を持たせる。


05-04. 並列パイプラインベースの Saga パターン¶
並列パイプラインベースの Saga パターンとは¶
オーケストレーションベースとコレオグラフィベースのパターンを組み合わせる。
イベント駆動のマイクロサービスを連続的にコールするルーターサービスを配置する。
実装パターン¶
▼ 自前で実装する場合¶
マイクロサービスの処理のなかで実装する。
▼ OSS を使用する場合¶
並列パイプラインの OSS (例:Apache Camel など) を使用する。
06. TCC パターン¶
TCC パターンとは¶
各マイクロサービスに各処理フェーズ (Try、Confirm、Cancel) を実行する API を実装し、API を順に実行する。
Try フェーズでは、ローカルトランザクション処理を開始する。
Confirm フェーズでは、ローカルトランザクション処理をコミットする。
Cancel フェーズでは、以前のフェーズで問題があった場合、ロールバックする。
07. クエリパターン¶
API Composition¶
複数のマイクロサービスにまたがる read 処理が必要になることがある。
API Composition サービスは、クライアントからのリクエストを受信し、複数のマイクロサービスに連鎖的にリクエストをルーティングする。
メモリ上で取得結果を結合し、クライアントにレスポンスする。
08. Outbox パターン¶
Outbox パターンとは¶
マイクロサービスアーキテクチャで、永続化処理とパブリッシュ処理が必要な場合に使える設計パターンである。
通信元マイクロサービスが永続化処理の結果をメッセージブローカーにパブリッシュしたい場合に、パブリッシュ処理をメッセージリレージョブに切り分け、通信元には永続化処理だけをもたせる。
Outbox パターンでは、Saga ログテーブルに加えて、Outbox テーブルを作成する。
メッセージリレージョブは Outbox テーブルのイベントを定期的にクエリし、メッセージブローカーにパブリッシュする。
| Id | AggregateType | AggregateId | Type | Payload |
|---|---|---|---|---|
| ec6e | Order | 123 | OrderCreated | {"id": 123, ...} |
| 8af8 | Order | 456 | OrderDetailCanceled | {"id": 456, ...} |
| 890b | Customer | 789 | InvoiceCreated | {"id": 789, ...} |

- https://microservices.io/patterns/data/transactional-outbox.html
- https://qiita.com/jokoshi/items/5016c3226f3009ddee10#31-transactional-messaging%E4%B8%8D%E6%95%B4%E5%90%88%E7%99%BA%E7%94%9F%E3%82%B1%E3%83%BC%E3%82%B91%E3%81%B8%E3%81%AE%E5%87%A6%E6%96%B9%E7%AE%8B
- https://debezium.io/blog/2019/02/19/reliable-microservices-data-exchange-with-the-outbox-pattern/
*実装例*
カタログコンテキストサービスが永続化処理の結果をメッセージブローカーにパブリッシュしたい場合に、次の設計パターンがある。
- Outbox パターンを使わない
- カタログコンテキストは永続化処理とパブリッシュ処理の両方を持つ
- Outbox パターンを使う
- パブリッシュ処理をメッセージリレージョブに切り分け、カタログコンテキストには永続化処理だけをもたせる。

メッセージリレーの設計パターン¶
▼ Polling publisher パターン¶
DB のイベントチェッカー (例:Debezium) を使用して、Outbox テーブルのイベントを検知する。
また、検知したイベントをメッセージブローカー (例:Apache Kafka、RabbitMQ など) にパブリッシュする。
Saga オーケストレーターのクライアントやマイクロサービス側では、これをポーリングする。

▼ Transaction log tailing パターンとは¶
トランザクションログ (例:MySQL バイナリログ、PostgreSQL WAL など) を追跡する。