PostgreSQLを使用したGoでのトランザクションアウトレックスパターンの実装
イベントとデータを合わせて記載してください。分割しないでください。
同時に成功すべき2つの書き込みが、最終的にはそれぞれ個別に失敗する可能性があります。
あなたの注文サービスは、データベースに注文を保存し、次にメッセージブローカーに order.created イベントを公開します。
これらの2つの操作は、順番に実行されます。
その間、何かが失敗します:ブローカーがダウンしている、ネットワークタイムアウトが発生する、プロセスが再起動する、またはコンテナが追い出される。データベースへの書き込みは成功しました。公開は失敗しました。新しい注文を知る必要がある後続サービスは、そのことを知り得ません。顧客から連絡が入るまで誰も気づきませんでした。
これがデュアルライト問題であり、分散システムにおけるサイレントデータ損失の最も一般的な原因の一つです。トランザクショナルアウトボックスパターンがその標準的な解決策です。

デュアルライトの問題
この障害モードは、一度見れば推論しやすいものです:
BEGIN;
INSERT INTO orders ... -- 成功する
COMMIT;
PUBLISH order.created ... -- 失敗、クラッシュ、または到達しない
データベースとメッセージブローカーはトランザクション境界を共有していません。両方を巻き戻すロールバックはありません。順序で save -> publish を行うすべてのサービスにこのギャップが存在します。このパターンは多くの形で現れます:
db.Save(order)に続くevents.Publish(OrderCreated{...})- トランザクションをコミットした後、外部ウェブフックを呼び出すHTTPハンドラ
- 1つのキューからレコードを処理し、結果を別のキューに書き込むワーカー
どの場合でも結果は同じです。一方の側が成功し、もう一方が失敗し、両方の個別操作がいつか成功を返したため、システムはモニタリングでは見えない状態になります。
リトライループはこの問題を解決しません。データベースコミット後の公開をリトライするのは、リトライ自体が信頼可能である場合のみ有効です。それは、あなたが持っていない正確に同じ耐久性の保証が必要です。
トランザクショナルアウトボックスパターンが行うこと
アウトボックスパターンは、直接公開を完全に排除することで、このギャップを解消します。ビジネスロジック内でブローカーを呼び出すのではなく、ビジネスデータと同じデータベーストランザクション内で outbox テーブルにイベントレコードを書き込みます。別のバックグラウンドプロセス(リレー)が、アウトボックステーブルから読み取ってブローカーに公開します。
BEGIN;
INSERT INTO orders ... -- ビジネスデータ
INSERT INTO outbox_events ... -- イベントレコード
COMMIT;
-- リレープロセス(別個):
SELECT ... FROM outbox_events FOR UPDATE SKIP LOCKED;
PUBLISH order.created ...
UPDATE outbox_events SET processed_at = NOW() WHERE id = $1;
両方の書き込みが成功するか、両方が失敗します。PostgreSQLからすでに得ているトランザクション保証が、イベントレコードにも適用されます。リレーは、イベントが永続ストレージにあるため、必要に応じて何度でも公開をリトライできます。リレーが途中でクラッシュした場合、再起動してリトライします。最悪の結果は、イベントが複数回公開されることです。これは、コンシューマーを冪等性にする(分散システムにおける冪等性)によって処理されます。
アウトボックステーブルのPostgreSQLスキーマ
スキーマは意図的にシンプルです:
CREATE TABLE outbox_events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id VARCHAR(100) NOT NULL,
event_type VARCHAR(100) NOT NULL,
payload JSONB NOT NULL,
attempts INT NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
processed_at TIMESTAMPTZ
);
-- パーシャルインデックス:処理されていない行のみをインデックス化し、行が完了としてマークされるにつれて小さく保たれます
CREATE INDEX idx_outbox_unprocessed
ON outbox_events (created_at)
WHERE processed_at IS NULL;
created_at WHERE processed_at IS NULL 上のパーシャルインデックスが重要です。これがないと、インデックスは過去に書き込まれたすべてのイベントとともに成長し、リレーのポーリングクエリは時間とともに遅くなります。これがあることで、インデックスは保留中の行のみをカバーし、ステady状態では、公開されたイベント数に関係なく、小さく制限されたセットになります。
主要なフィールド選択:
aggregate_typeとaggregate_idは、イベントがどのエンティティに属しているかを示します。順序保証とルーティングに役立ちます。event_typeは、コンシューマーが期待するイベント名です。payload JSONBはイベントボディを保存します。必要に応じてクエリできるように、TEXTではなくJSONBを使用します。attemptsは、リレーがこの行の公開を試みた回数を追跡します。リトライ制限とデッドレター処理に使用されます。processed_atは保留中の行ではNULLであり、リレーが公開に成功すると設定されます。
ビジネスデータとアウトボックスイベントを1つのトランザクションで書き込む
ビジネスロジックは、単一の BeginTx / Commit 呼び出し内で両方のレコードを書き込みます。ここには公開呼び出しはありません。データベース書き込みのみです。
type OrderService struct {
db *sql.DB
}
func (s *OrderService) CreateOrder(ctx context.Context, order Order) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
if _, err := tx.ExecContext(ctx, `
INSERT INTO orders (id, customer_id, total, created_at)
VALUES ($1, $2, $3, NOW())
`, order.ID, order.CustomerID, order.Total); err != nil {
return fmt.Errorf("insert order: %w", err)
}
payload, err := json.Marshal(map[string]any{
"order_id": order.ID,
"customer_id": order.CustomerID,
"total": order.Total,
})
if err != nil {
return fmt.Errorf("marshal payload: %w", err)
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO outbox_events
(aggregate_type, aggregate_id, event_type, payload)
VALUES ($1, $2, $3, $4)
`, "order", order.ID, "order.created", payload); err != nil {
return fmt.Errorf("insert outbox event: %w", err)
}
return tx.Commit()
}
tx.Commit() が失敗した場合、注文行もアウトボックス行も永続化されません。成功した場合、両方がデータベース内に確実に存在することが保証されます。それ以降の任意の時点でリレーがイベントを公開できます。直ちに、1秒後、またはクラッシュ後のリレー再起動後でも構いません。
これはビジネスレイヤーで必要な唯一の変更です。パターンの残りの部分はリレーにあります。
Goリレーの実装
リレーは、タイマーでアウトボックステーブルをポーリングするバックグラウンドワーカーです。処理されていない行のバッチを取得し、それぞれを公開し、完了としてマークします。アプリケーションと同じバイナリ内に保持するか、別のプロセスとして実行します。どちらでも機能しますが、同じバイナリの方が運用が簡単です。
type OutboxRelay struct {
db *sql.DB
publisher Publisher
logger *slog.Logger
batchSize int
pollInterval time.Duration
maxAttempts int
}
func (r *OutboxRelay) Run(ctx context.Context) error {
ticker := time.NewTicker(r.pollInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
if err := r.processBatch(ctx); err != nil {
r.logger.Error("outbox relay batch failed", "err", err)
}
}
}
}
リレーはコンテキストの取消を尊重するため、グラシューラスシャットダウンとの統合が容易です。コンテキストのライフタイムと取消パターンの詳細については、Go context.Done()の正しい使い方を参照してください。
FOR UPDATE SKIP LOCKED:同時実行ワーカーパターン
processBatch 関数は、同時実行されるリレーワーカーを安全に処理するために FOR UPDATE SKIP LOCKED を使用します:
func (r *OutboxRelay) processBatch(ctx context.Context) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
rows, err := tx.QueryContext(ctx, `
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox_events
WHERE processed_at IS NULL
AND attempts < $1
ORDER BY created_at
LIMIT $2
FOR UPDATE SKIP LOCKED
`, r.maxAttempts, r.batchSize)
if err != nil {
return fmt.Errorf("query outbox: %w", err)
}
defer rows.Close()
type row struct {
id string
aggregateType string
aggregateID string
eventType string
payload json.RawMessage
}
var batch []row
for rows.Next() {
var e row
if err := rows.Scan(
&e.id, &e.aggregateType, &e.aggregateID, &e.eventType, &e.payload,
); err != nil {
return fmt.Errorf("scan row: %w", err)
}
batch = append(batch, e)
}
if err := rows.Err(); err != nil {
return err
}
for _, e := range batch {
if err := r.publisher.Publish(ctx, e.eventType, e.aggregateID, e.payload); err != nil {
r.logger.Error("publish failed", "event_id", e.id, "err", err)
if _, err := tx.ExecContext(ctx,
`UPDATE outbox_events SET attempts = attempts + 1 WHERE id = $1`, e.id,
); err != nil {
r.logger.Error("increment attempts failed", "event_id", e.id, "err", err)
}
continue
}
if _, err := tx.ExecContext(ctx,
`UPDATE outbox_events SET processed_at = NOW() WHERE id = $1`, e.id,
); err != nil {
return fmt.Errorf("mark processed: %w", err)
}
}
return tx.Commit()
}
FOR UPDATE SKIP LOCKED は2つのことを実行します。まず、FOR UPDATE はトランザクションの持続中に選択された行をロックし、他のトランザクションがそれらを選択するのを防ぎます。次に、SKIP LOCKED は、行がすでに他のトランザクションによってロックされている場合、クエリがそれを待たずにスキップすることを意味します。その結果、複数のリレーワーカーが並列で実行でき、それぞれが重複しない行のサブセットを取得します。
SKIP LOCKED がない場合、2番目のワーカーは、同じ行を見る前に最初のトランザクションがコミットするまでブロックされます。その時点では、それらはすでに完了としてマークされています。SKIP LOCKED があることで、2番目のワーカーは待たずに直ちに異なる行を取得し、安全な水平スケーリングを実現します。
上記のコードにあるスキャンと公開の分離に注意してください。すべての行は、公開ループが始まる前にスライスにスキャンされます。これにより、ブローカーへのネットワーク呼び出し全体にわたって開いた *sql.Rows カーソルを保持する必要がなくなり、トランザクションが不必要に長く開いたままになるのを防ぎます。
冪等性と重複排除
リレーは少なくとも1回の公開を行います。イベントを公開し、processed_at の更新をコミットする前にクラッシュした場合、再起動時に同じイベントを再度公開します。これは避けられません。データベースとメッセージブローカー間での正確に1回の実行は、分散トランザクションコーディネーターがない限り、このトレードオフが必要です。
コンシューマーは冪等性である必要があります。最も単純なアプローチは、processed_events テーブルで処理済みのイベントIDを追跡することです:
CREATE TABLE processed_events (
event_id UUID PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
func (h *OrderHandler) HandleOrderCreated(ctx context.Context, eventID string, payload []byte) error {
// イベントIDを自然キーとして使用して重複排除
_, err := h.db.ExecContext(ctx, `
INSERT INTO processed_events (event_id) VALUES ($1)
ON CONFLICT (event_id) DO NOTHING
`, eventID)
if err != nil {
return fmt.Errorf("dedup check: %w", err)
}
// 実際の挿入が行われたか(1行)、または何もしなかったか(0行)を確認します
// より簡単なアプローチ:RETURNINGを使用するか、影響を受けた行をチェックします
// 影響を受けた行が0の場合、これは重複です。スキップします
...
}
実際には、多くのチームはブローカー独自の重複排除ヘッダー(Kafkaのログ圧縮トピック用の key フィールドや、RabbitMQの message-id ヘッダーなど)に依存し、データベースレベルの重複排除をフォールバックとして扱います。どちらも適用可能な有効なレイヤーです。
公開メッセージにアウトボックスイベントの id(UUID)を重複排除キーとして含めます。コンシューマーは、どの重複排除メカニズムを好むかに関わらず、それを使用できます。
リトライポリシーとポイズンメッセージ
attempts カラムがリトライポリシーを駆動します。リレーは attempts >= maxAttempts の行をスキップし、それらをデッドレターとして扱います。それらは、別のプロセスまたはオペレーターのアラートによって処理されます。
単純なデッドレタービュー:
CREATE VIEW outbox_dead_letters AS
SELECT *
FROM outbox_events
WHERE attempts >= 5
AND processed_at IS NULL
ORDER BY created_at;
優れた本番環境のリトライポリシー:
maxAttemptsを5〜10に設定します(リトライのコストに応じて調整)。- 指数バックオフを検討します:
retry_afterカラムを含め、retry_after > NOW()の行をスキップします。 COUNT(*) FROM outbox_dead_lettersが閾値を超えた場合にアラートを設定します。- 手動リトライパスを提供します:特定の行に対して
attempts = 0とretry_after = NULLをリセットする管理エンドポイントまたはスクリプト。
ポイズンメッセージ(コンシューマーのバグやスキーマの不整合によって一貫して失敗する行)は、健全なメッセージをブロックしてはいけません。リレーはティックごとにバッチを処理し、失敗を削除するのではなく試行回数の増分でマークするため、健全な行は通常どおり進行し、ポイズンされた行はデッドレターの閾値に達するまで試行回数が累積されます。この outbox_dead_letters ビューは、同じパターンのデータベースサイドバージョンです。ブローカーネイティブなデッドレターキューが実装しています—閾値後に隔離し、量でアラートし、再生前に意図的な決定が必要です。
イベント順序付けとパーティショニング
ポーリングクエリは created_at で順序付けられ、バッチ内でFIFO(先入れ先出し)順序を提供します。ほとんどのユースケースではそれで十分です。エンティティごとの厳格な順序付けが重要な場合—例えば、同じ注文に対して order.updated が order.created より前に公開されないようにすること—アグリゲートごとの順序付けが必要です。
ORDER BY 句に aggregate_id を追加し、Apache Kafkaのようなパーティショニングされたトピックに公開する際にメッセージキーとして使用します。Kafkaは同じキーを持つすべてのメッセージを同じパーティションにルーティングし、パーティションは順序通りに消費されます。これにより、単一のリレーインスタンスを必要とするグローバル順序付けではなく、アグリゲートごとの順序付け保証が得られます。
ORDER BY aggregate_id, created_at
パーティショニング順序付けをサポートしていないブローカー(基本AMQPキューなど)の場合、単一インスタンスのリレーまたはコンシューマーでのアプリケーションレベルの順序付けチェックが実用的な代替手段です。
LISTEN/NOTIFYでポーリングレイテンシーを削減する
1秒のポーリング間隔は、平均イベントレイテンシーが500ミリ秒であることを意味します。ほとんどのワークロードではこれで問題ありません。ニアゼロレイテンシーが必要な場合、PostgreSQLの LISTEN/NOTIFY メカニズムにより、新しいアウトボックス行が挿入されるとすぐにリレーを起動できます。
アウトボックステーブルにトリガーを追加します:
CREATE OR REPLACE FUNCTION notify_outbox_insert() RETURNS trigger AS $$
BEGIN
PERFORM pg_notify('outbox_event', NEW.id::text);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER outbox_insert_notify
AFTER INSERT ON outbox_events
FOR EACH ROW EXECUTE FUNCTION notify_outbox_insert();
リレーでは、チャンネルをリッスンし、通知で起動しますが、定期的なポーリングへのフォールバックも維持します:
func (r *OutboxRelay) Run(ctx context.Context) error {
listener := pq.NewListener(r.dsn, 10*time.Second, time.Minute, nil)
defer listener.Close()
if err := listener.Listen("outbox_event"); err != nil {
return fmt.Errorf("listen: %w", err)
}
ticker := time.NewTicker(5 * time.Second) // フォールバックポーリング
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-listener.Notify:
if err := r.processBatch(ctx); err != nil {
r.logger.Error("outbox batch failed (notify)", "err", err)
}
case <-ticker.C:
if err := r.processBatch(ctx); err != nil {
r.logger.Error("outbox batch failed (poll)", "err", err)
}
}
}
}
フォールバックタイカーは、リレーの再起動中またはネットワークの一時的な障害中に逃された通知を処理します。フォールバック間隔はミリ秒ではなく数秒に保つ必要があります。その役割は低レイテンシーではなく回復です。
観測性:メトリクス、ログ、およびアラート
アウトボックスはインフラです。インフラとして扱い、それに合わせて計測します。
主要なメトリクス:
var (
outboxPublished = prometheus.NewCounter(prometheus.CounterOpts{
Name: "outbox_events_published_total",
Help: "正常に公開されたアウトボックスイベントの総数。",
})
outboxFailed = prometheus.NewCounterVec(prometheus.CounterOpts{
Name: "outbox_events_failed_total",
Help: "イベントタイプ別のアウトボックス公開失敗の総数。",
}, []string{"event_type"})
outboxPending = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "outbox_events_pending",
Help: "現在処理されていないアウトボックスイベントの数。",
})
outboxBatchDuration = prometheus.NewHistogram(prometheus.HistogramOpts{
Name: "outbox_batch_duration_seconds",
Help: "各アウトボックス処理バッチの所要時間。",
Buckets: prometheus.DefBuckets,
})
)
ゲージ更新: outbox_events_pending を正確に保つために定期的なクエリを実行します:
SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;
検討すべきアラート閾値:
- 2分以上
outbox_events_pending > 1000:リレーが遅れているか、停止している。 outbox_events_pendingが単調増加している:ブローカーがダウンしているか、リレーがクラッシュしている。- デッドレターカウントがゼロでない:スキーマまたはコンシューマーのバグの調査が必要。
outbox_batch_duration_seconds p95 > 5s:データベースが遅いか、バッチサイズが大きすぎる。
構造化ログフィールド: リレーからのすべてのログ行に event_id、event_type、aggregate_id、および attempt を含めます。これらのフィールドにより、失敗した公開を特定のアウトボックス行と後続のコンシューマートレースに関連付けることができます。
アウトボックス vs 直接キュー vs サガ
アウトボックスパターンは、すべての調整問題に適したツールではありません。以下に比較を示します:
| アプローチ | 原子性 | 複雑さ | 使用時 |
|---|---|---|---|
| 直接公開 | なし | 低 | イベントの稀な損失が許容可能 |
| トランザクショナルアウトボックス | 強い | 中 | 単一サービスからの信頼できるイベント配信 |
| サガパターン | 最終的な整合性 | 高 | 複数のデータベースにまたがるマルチサービストランザクション |
| 2フェーズコミット | 強い | 非常に高い | 実際にはほとんど実用的ではありません。ほとんどの分散システムで回避されます |
アウトボックスパターンは、単一サービスが自身の状態変化を反映するイベントを信頼できる方法でエミットすることを保証します。それは複数のサービス間で状態変化を調整するものではありません。それはサガパターンの役割です。ブローカーの選択—RabbitMQ、SQS, または Kafka—はアウトボックスパターン自体とは独立しています。リレーは、システムが使用するいずれのブローカーにも公開します。
サガを構築している場合、アウトボックスパターンは依然として有用です。サガの各参加者は、アウトボックスを使用してローカル状態変化とサガイベントを1つのトランザクションで書き込み、その後サガオーケストレーターまたはコリオグラフィがそれらのイベントを信頼できる方法で読み取ります。
WALベースのCDCを代替リレーとして
ポーリングの代わりに、PostgreSQLのWrite-Ahead Log (WAL) をテールし、レプリケーションストリームから直接アウトボックス挿入を読み取ることができます。Debeziumなどのツールがこれを行います。利点は、レイテンシーが低く、アウトボックステーブルへのロック圧力が発生しないことです。欠点は、運用の複雑さ、専用のPostgreSQLレプリケーションスロット、および実行および監視する外部サービスが必要となることです。
ほとんどのチームにとって、上記で説明したポーリングリレーが適切な出発点です。高いアウトボックス挿入率(秒あたり数万)、サブ100ミリ秒のイベントレイテンシーが必要、または他の変更キャプチャニーズのために既にDebeziumを実行している場合に、WALテールが意味を持ちます。
sqlc統合
型安全なGoデータベースコードにsqlcを使用している場合、アウトボックスクエリは自然に適合します:
-- name: InsertOutboxEvent :exec
INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload)
VALUES (@aggregate_type, @aggregate_id, @event_type, @payload);
-- name: FetchOutboxBatch :many
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox_events
WHERE processed_at IS NULL
AND attempts < @max_attempts
ORDER BY created_at
LIMIT @batch_size
FOR UPDATE SKIP LOCKED;
-- name: MarkOutboxProcessed :exec
UPDATE outbox_events SET processed_at = NOW() WHERE id = @id;
-- name: IncrementOutboxAttempts :exec
UPDATE outbox_events SET attempts = attempts + 1 WHERE id = @id;
-- name: OutboxPendingCount :one
SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;
sqlcは各クエリに対して型安全な関数を生成するため、文字列補間エラーを回避し、アウトボックスクエリロジックをデータベースアクセスレイヤーの残りと一緒に保つことができます。
本番環境チェックリスト
アウトボックス実装をリリースする前にこれを使用してください:
データベース
- アウトボックステーブルに
created_at WHERE processed_at IS NULLのパーシャルインデックスがある -
attemptsカラムが存在し、デフォルトが0 - デッドレタービューまたはクエリが定義されている
- 古い処理済み行は定期的にアーカイブまたは削除される(毎晩のクリーンアップジョーンで十分)
リレー
- ポーリングクエリで
FOR UPDATE SKIP LOCKEDが使用されている - リレーがトランザクション内で実行される(クエリ前に開始し、すべての更新後にコミット)
- バッチサイズは制限されている(50〜200行が一般的)
- リレーはグラシューラスシャットダウンのためにコンテキスト取消を尊重する
- 失敗した公開はバッチを中止させるのではなく、
attemptsを増分する
冪等性
- 公開メッセージに重複排除キーとしてアウトボックス
idが含まれている - コンシューマーが冪等性であるか、ブローカーが重複排除を提供している
- 重複排除パターンについては分散システムにおける冪等性を参照
観測性
-
outbox_events_pendingゲージが監視され、アラート設定されている - デッドレターカウントがアラート設定されている
- リレーバッチの所要時間が追跡されている
- 構造化ログに
event_id、event_type、およびaggregate_idが含まれている
オペレーション
- デッドレター行に対する手動リトライパスが存在する
- リレーの再起動動作がテストされている(正しく再公開するか?)
- ブローカー停止時の動作がテストされている(アウトボックスが正しく成長し、排水するか?)
最後の考え
デュアルライト問題は、インシデントを引き起こすまでエッジケースとして無視しやすいものです。トランザクショナルアウトボックスパターンは、すでに持っているツール—PostgreSQLトランザクション、バックグラウンドgoroutine、および1つの追加テーブル—でそれを解決します。リレーは構築が簡単で、運用が簡単で、推論も簡単です。
コストは、コンシューマーが少なくとも1回の実行を目的として設計されている必要があることです。それは妥当なトレードオフです。分散トランザクションなしでデータベースとブローカー間で正確に1回の実行を実現することは実際には達成できません。それ以外であることを pretending することは、障害条件下でイベントをサイレントにドロップまたは二重処理するシステムにつながります。
データと一緒にイベントを書き込みます。信頼できる方法で中継します。コンシューマーを冪等性にする。これこそがパターンの全体です。
この記事は、本番環境のアプリアーキテクチャクラスターの一部です。