soobook
DATABASE

MongoDB Change Stream으로 구현하는 Outbox 패턴

Outbox 패턴은 DB 쓰기와 이벤트 발행의 원자성을 보장하는 핵심 패턴이다. MongoDB를 사용한다면 외부 인프라 없이 Change Stream만으로 구현할 수 있다.

왜 Change Stream인가

CDC 글에서 Outbox 이벤트를 외부로 전달하는 방식을 두 가지로 분류했다.

  1. 폴링: outbox 테이블을 주기적으로 조회. 단순하지만 지연과 DB 부하가 있다
  2. CDC: WAL을 감시하여 즉시 전달. Debezium이 대표적

MongoDB 환경에서는 Change Stream이 세 번째 선택지가 된다. MongoDB가 자체적으로 제공하는 CDC 메커니즘이며, oplog(Operation Log)을 tailing하여 컬렉션의 변경을 실시간으로 구독한다.

Debezium처럼 별도 인프라를 운영할 필요 없이, Watch() 한 줄로 시작할 수 있다.

flowchart LR
    subgraph "MongoDB Transaction"
        W1["orders에<br/>주문 저장"]
        W2["outbox에<br/>이벤트 저장"]
    end
    W1 --> W2
    W2 --> Oplog["oplog.rs"]
    Oplog --> CS["Change Stream<br/>(Relay)"]
    CS --> K["Kafka / gRPC /<br/>외부 시스템"]

Change Stream은 Replica Set 또는 Sharded Cluster에서만 사용 가능하다.

Standalone 모드에서는 oplog이 존재하지 않기 때문이다.

로컬 개발 환경에서도 단일 노드 Replica Set으로 구성해야 한다.

Outbox 구조 설계

Outbox document

Outbox 컬렉션에 저장할 이벤트 도큐먼트의 구조부터 정의한다.

type OutboxEvent struct {
    ID            primitive.ObjectID `bson:"_id,omitempty"`
    AggregateType string             `bson:"aggregateType"` // e.g. "Order", "Payment"
    AggregateID   string             `bson:"aggregateId"`
    EventType     string             `bson:"eventType"`     // e.g. "OrderCreated"
    Payload       bson.Raw           `bson:"payload"`
    CreatedAt     time.Time          `bson:"createdAt"`
}
필드역할
aggregateType이벤트의 출처가 되는 도메인 객체 종류
aggregateId해당 도메인 객체의 고유 식별자
eventType발생한 이벤트의 종류. Consumer가 라우팅에 사용
payload이벤트 본문. bson.Raw로 두면 직렬화 형태를 유연하게 유지할 수 있다

Outbox 이벤트에 processed 같은 상태 필드를 두지 않는다.

Change Stream이 oplog 기반으로 이벤트를 감지하므로, 폴링처럼 “아직 처리 안 된 것”을 조회할 필요가 없다.

도큐먼트가 insert되는 순간 Change Stream이 바로 수신한다.

비즈니스 도큐먼트

예시로 대충 만들어보자.

type Order struct {
    ID        primitive.ObjectID `bson:"_id,omitempty"`
    UserID    string             `bson:"userId"`
    Items     []OrderItem        `bson:"items"`
    Amount    int64              `bson:"amount"`
    Status    string             `bson:"status"`
    CreatedAt time.Time          `bson:"createdAt"`
}

type OrderItem struct {
    ProductID string `bson:"productId"`
    Quantity  int    `bson:"quantity"`
    Price     int64  `bson:"price"`
}

Write: 트랜잭션으로 원자성 보장

핵심은 비즈니스 데이터와 outbox 이벤트를 같은 트랜잭션으로 저장하는 것이다. MongoDB 4.0부터 지원하는 multi-document 트랜잭션을 사용한다.

func CreateOrder(ctx context.Context, client *mongo.Client, order Order) error {
    session, err := client.StartSession()
    if err != nil {
        return fmt.Errorf("start session: %w", err)
    }
    defer session.EndSession(ctx)

    _, err = session.WithTransaction(ctx, func(sessCtx context.Context) (interface{}, error) {
        db := client.Database("shop")

        order.ID = primitive.NewObjectID()
        order.Status = "created"
        order.CreatedAt = time.Now()

        // 1. 주문 저장
        if _, err := db.Collection("orders").InsertOne(sessCtx, order); err != nil {
            return nil, fmt.Errorf("insert order: %w", err)
        }

        // 2. outbox에 이벤트 저장
        payload, err := bson.Marshal(order)
        if err != nil {
            return nil, fmt.Errorf("marshal payload: %w", err)
        }

        outboxEvent := OutboxEvent{
            AggregateType: "Order",
            AggregateID:   order.ID.Hex(),
            EventType:     "OrderCreated",
            Payload:       payload,
            CreatedAt:     time.Now(),
        }
        if _, err := db.Collection("outbox").InsertOne(sessCtx, outboxEvent); err != nil {
            return nil, fmt.Errorf("insert outbox: %w", err)
        }

        return nil, nil
    })

    return err
}

WithTransaction이 커밋까지 처리한다. 주문 저장이나 outbox 저장 중 하나라도 실패하면 둘 다 롤백된다.

이것이 Outbox 패턴의 핵심이다. Dual write 문제를 방지한다.

DB 커밋과 메시지 발행이라는 dual write 문제를 DB 커밋 하나로 축소시킨다.

Read: Change Stream Relay

Outbox 컬렉션의 변경을 감지하여 외부 시스템에 전달하는 Relay 프로세스를 구현한다.

Change Stream 이벤트 구조

Change Stream이 전달하는 이벤트는 이런 형태다.

{
  "_id": { "_data": "826..." },
  "operationType": "insert",
  "clusterTime": { "$timestamp": { "t": 1712345678, "i": 1 } },
  "ns": { "db": "shop", "coll": "outbox" },
  "documentKey": { "_id": { "$oid": "..." } },
  "fullDocument": {
    "_id": { "$oid": "..." },
    "aggregateType": "Order",
    "aggregateId": "680...",
    "eventType": "OrderCreated",
    "payload": { ... },
    "createdAt": "2026-04-06T..."
  }
}

_id 필드가 Resume Token이다. 이 값을 저장해두면 장애 후 해당 지점부터 다시 수신할 수 있다.

Relay 구현

// ChangeEvent는 Change Stream에서 수신하는 이벤트를 디코딩하기 위한 구조체이다.
type ChangeEvent struct {
    FullDocument OutboxEvent `bson:"fullDocument"`
}

type OutboxRelay struct {
    client     *mongo.Client
    publisher  EventPublisher
    tokenStore ResumeTokenStore
}

// EventPublisher는 외부 시스템에 이벤트를 발행하는 인터페이스이다.
type EventPublisher interface {
    Publish(ctx context.Context, event OutboxEvent) error
}

// ResumeTokenStore는 Resume Token을 영속 저장하고 로드하는 인터페이스이다.
type ResumeTokenStore interface {
    Save(ctx context.Context, token bson.Raw) error
    Load(ctx context.Context) (bson.Raw, error)
}

func (r *OutboxRelay) Run(ctx context.Context) error {
    outbox := r.client.Database("shop").Collection("outbox")

    // insert 이벤트만 수신
    pipeline := mongo.Pipeline{
        {{"$match", bson.M{"operationType": "insert"}}},
    }

    opts := options.ChangeStream().
        SetFullDocument(options.Default)

    // 저장된 resume token이 있으면 해당 지점부터 재개
    if token, err := r.tokenStore.Load(ctx); err == nil && token != nil {
        opts.SetResumeAfter(token)
    }

    stream, err := outbox.Watch(ctx, pipeline, opts)
    if err != nil {
        return fmt.Errorf("watch outbox: %w", err)
    }
    defer stream.Close(ctx)

    for stream.Next(ctx) {
        var event ChangeEvent
        if err := stream.Decode(&event); err != nil {
            return fmt.Errorf("decode event: %w", err)
        }

        // 외부 시스템에 발행
        if err := r.publisher.Publish(ctx, event.FullDocument); err != nil {
            return fmt.Errorf("publish event: %w", err)
        }

        // 발행 성공 후 resume token 저장
        if err := r.tokenStore.Save(ctx, stream.ResumeToken()); err != nil {
            return fmt.Errorf("save resume token: %w", err)
        }
    }

    return stream.Err()
}

흐름을 정리하면 이렇다.

  1. 저장된 Resume Token이 있으면 해당 지점부터, 없으면 현재 시점부터 Change Stream을 연다
  2. stream.Next()가 블로킹으로 다음 이벤트를 기다린다
  3. 이벤트를 디코딩하고 외부 시스템에 발행한다
  4. 발행 성공 후 Resume Token을 저장한다
  5. 프로세스가 재시작되면 1번으로 돌아가 중단 지점부터 재개한다

Resume Token 영속화

Resume Token을 MongoDB 자체에 저장하는 구현이다. 별도의 외부 저장소가 필요 없다.

type MongoTokenStore struct {
    coll *mongo.Collection
}

func NewMongoTokenStore(client *mongo.Client) *MongoTokenStore {
    return &MongoTokenStore{
        coll: client.Database("shop").Collection("resumeTokens"),
    }
}

func (s *MongoTokenStore) Save(ctx context.Context, token bson.Raw) error {
    filter := bson.M{"_id": "outbox-relay"}
    update := bson.M{
        "$set": bson.M{
            "token":     token,
            "updatedAt": time.Now(),
        },
    }
    opts := options.Update().SetUpsert(true)
    _, err := s.coll.UpdateOne(ctx, filter, update, opts)
    return err
}

func (s *MongoTokenStore) Load(ctx context.Context) (bson.Raw, error) {
    var result struct {
        Token bson.Raw `bson:"token"`
    }
    err := s.coll.FindOne(ctx, bson.M{"_id": "outbox-relay"}).Decode(&result)
    if err == mongo.ErrNoDocuments {
        return nil, nil
    }
    return result.Token, err
}

Resume Token은 oplog에 대응하므로, oplog에서 해당 엔트리가 삭제되면 재개할 수 없다. oplog은 capped collection이라 용량이 가득 차면 오래된 엔트리부터 삭제된다. Relay가 오래 중단되지 않도록 모니터링하고, oplog 크기를 충분히 확보해야 한다.

장애 복구 시나리오

Relay 프로세스의 장애 지점에 따라 어떤 일이 일어나는지 정리한다.

flowchart TD
    A["이벤트 수신"] --> B["외부 발행"]
    B --> C["Resume Token 저장"]

    B -->|"여기서 장애"| R1["재시작 시 같은 이벤트 재수신<br/>→ Consumer 멱등성 필요"]
    C -->|"여기서 장애"| R2["재시작 시 같은 이벤트 재수신<br/>→ Consumer 멱등성 필요"]
장애 지점결과대응
발행 전이벤트 미발행. 재시작 시 같은 이벤트를 다시 수신하므로 유실 없음없음
발행 후, Token 저장 전이벤트가 중복 발행될 수 있음Consumer 멱등성
Token 저장 후정상. 다음 이벤트부터 재개없음

어떤 시점에 장애가 나더라도 이벤트가 유실되지는 않는다.

다만 중복이 발생할 여지가 있으므로 Consumer 측에서 멱등성을 처리해야 한다.

Debezium을 포함한 모든 CDC 기반 시스템에 똑같이 적용되는 원칙이다.

Consumer 멱등성

Change Stream은 at-least-once 보장이다. Consumer에서 중복 이벤트를 걸러내야 한다.

Outbox 이벤트의 _id(ObjectID)를 멱등성 키로 활용하면 된다. 처리 완료된 이벤트 ID를 별도 컬렉션에 기록하고, 중복 수신 시 스킵한다.

type IdempotentConsumer struct {
    processed *mongo.Collection
    handler   func(ctx context.Context, event OutboxEvent) error
}

func (c *IdempotentConsumer) Handle(ctx context.Context, event OutboxEvent) error {
    // 이미 처리된 이벤트인지 확인
    eventID := event.ID.Hex()
    _, err := c.processed.InsertOne(ctx, bson.M{
        "_id":         eventID,
        "processedAt": time.Now(),
    })
    if mongo.IsDuplicateKeyError(err) {
        // 이미 처리됨, 스킵
        return nil
    }
    if err != nil {
        return fmt.Errorf("record processed: %w", err)
    }

    // 비즈니스 로직 실행
    if err := c.handler(ctx, event); err != nil {
        // 실패 시 처리 기록 삭제하여 재시도 가능하게
        c.processed.DeleteOne(ctx, bson.M{"_id": eventID})
        return err
    }

    return nil
}

_id에 unique index가 자동으로 걸리므로, InsertOne이 반환하는 duplicate key error로 중복을 원자적으로 감지한다.

checkProcessedmarkProcessed와 같은 두 단계 방식보다 race condition에 안전하다.

처리 기록이 무한히 쌓이지 않도록 TTL 인덱스를 설정한다. 이벤트 재발행 가능 기간(oplog 보존 기간)보다 길게 잡으면 된다.

db.processedEvents.createIndex(
  { "processedAt": 1 },
  { expireAfterSeconds: 604800 } // 7일
)

Outbox 컬렉션 관리

Outbox 컬렉션의 도큐먼트는 이벤트 발행이 완료되면 더 이상 필요하지 않다.

방치하면 컬렉션이 무한히 커지므로 주기적으로 비워야 한다.

TTL 인덱스로 자동 삭제

가장 간단한 방법이다. createdAt 필드에 TTL 인덱스를 걸어, 일정 시간이 지나면 MongoDB가 자동으로 삭제한다.

db.outbox.createIndex(
  { "createdAt": 1 },
  { expireAfterSeconds: 86400 } // 24시간
)

TTL은 oplog 보존 기간보다 충분히 길게 설정해야 한다. Relay가 일시적으로 중단되었다가 재개할 때, 아직 처리하지 못한 이벤트가 삭제되면 유실이 발생한다.

Capped Collection 활용

Outbox 컬렉션을 capped collection으로 생성하면, 크기 제한을 초과할 때 오래된 도큐먼트가 자동으로 삭제된다.

db.createCollection("outbox", {
  capped: true,
  size: 104857600,  // 100MB
  max: 1000000      // 최대 100만 건
})

다만 capped collection에서는 도큐먼트 삭제(deleteOne)나 크기를 늘리는 업데이트가 불가능하다는 제약이 있다. Outbox처럼 insert-only 용도에는 적합하다.

Tips

Connection Pool 관리

각 Change Stream은 connection을 하나 점유한다. 여러 컬렉션을 감시해야 한다면 컬렉션마다 Watch()를 호출하는 대신 데이터베이스 수준으로 열고 파이프라인에서 필터링하는 편이 connection을 덜 쓴다.

// 컬렉션별 watch → connection 3개 소비
stream1, _ := ordersCol.Watch(ctx, pipeline)
stream2, _ := usersCol.Watch(ctx, pipeline)
stream3, _ := paymentsCol.Watch(ctx, pipeline)

// 데이터베이스 수준 watch → connection 1개
pipeline := mongo.Pipeline{
    {{"$match", bson.M{
        "ns.coll": bson.M{"$in": bson.A{"orders", "users", "payments"}},
    }}},
}
stream, _ := db.Watch(ctx, pipeline)

Aggregation Pipeline 필터링

Watch()에 aggregation pipeline을 전달하면 관심 있는 이벤트만 수신한다.

MongoDB 8.0에서는 $match가 oplog reader 수준으로 push down되어, 불필요한 이벤트가 아예 네트워크를 타지 않는다.

pipeline := mongo.Pipeline{
    // OrderCreated 이벤트만 수신
    {{"$match", bson.M{
        "operationType":              "insert",
        "fullDocument.aggregateType": "Order",
        "fullDocument.eventType":     "OrderCreated",
    }}},
    // 불필요한 필드 제거하여 페이로드 축소
    {{"$project", bson.M{
        "fullDocument": 1,
        "operationType": 1,
    }}},
}
stream, err := outbox.Watch(ctx, pipeline)

Pre/Post Image

MongoDB 6.0부터 변경 전후의 도큐먼트 상태를 함께 받을 수 있다.

Outbox 패턴에서는 주로 fullDocument 옵션을 사용하지만, 상태 변경 이벤트를 발행할 때는 변경 전 상태가 필요한 경우가 있다.

// 컬렉션에 pre/post image 활성화
db.runCommand({
  collMod: "orders",
  changeStreamPreAndPostImages: { enabled: true }
})
opts := options.ChangeStream().
    SetFullDocumentBeforeChange(options.WhenAvailable).
    SetFullDocument(options.WhenAvailable)

stream, err := collection.Watch(ctx, pipeline, opts)

Pre-image는 config.system.preimages 컬렉션에 저장되어 추가 스토리지를 소비한다.

필요한 컬렉션에만 선택적으로 활성화하고, expireAfterSeconds로 보존 기간을 설정해야 한다.

Sharded Cluster에서의 순서 보장

MongoDB는 global logical clock을 사용하여 샤드 간에도 변경 순서를 보장한다.

다만 특정 샤드에 활동이 없는 “cold shard”가 있으면, 해당 샤드 확인을 기다리느라 응답 지연이 발생할 수 있다.

Outbox 컬렉션을 단일 샤드에 두면 이 문제를 피할 수 있다.

이벤트 순서가 비즈니스 로직에 중요하다면 Outbox 컬렉션은 샤딩하지 않는 것이 안전하다.

References