상품 구매 예시 이벤트 중재자 구현

공부용·2025년 10월 5일

1. What (무엇을 구현할 것인지)

쇼핑 예제를 통해 주문하는 과정을 이벤트 중재자를 사용해 풀어내려한다.

워크 플로우

  1. 주문을 한다
    • 주문을 생성한다.
  2. 주문을 처리
    • 주문이 접수됐다고 고객에게 알림을 전송한다.
    • 결제를 한다.
    • 재고를 차감한다.
  3. 주문을 이행
    • 주문한 물품을 포장한다.
  4. 주문을 배송
    • 배송 처리 한다.
    • 물품의 배송 준비가 되었다고 고객에게 알린다.
  5. 배송을 고객에게 알린다.
    • 주문한 물품을 배송했다고 고객에게 알린다.

각 단계별 세부 과정은 비동기적으로 수행된다.

2. How (어떻게 구현할까)

1. 이벤트 기반 아키텍처 (EDA) 관점

  • 처음 고민: 주문 프로세스는 단계별로 많은 하위 작업(결제, 재고 차감, 알림 발송 등)이 발생한다.
  • 이를 동기적으로 처리하면 특정 단계가 오래 걸릴 때 전체 요청이 블로킹될 수 있다.
  • 따라서 각 단계를 이벤트로 정의하고, 이벤트를 발행하면 다른 모듈들이 구독해서 비동기적으로 처리하도록 만들 필요가 있다.
  • 예:
    • OrderCreated 이벤트 -> 결제 모듈/재고 모듈/알림 모듈이 각각 구독
    • InventoryReserved 이벤트 -> 포장 모듈이 구독
      이렇게 하면 모듈 간 직접 호출이 아닌 이벤트를 매개로 한 느슨한 결합이 만들어진다.

2. 이벤트 중재자 (Mediator) 관점

  • 단순 Pub/Sub 구조라면 "누가 언제 어떤 이벤트를 처리해야 하는지" 통제가 어렵다.
  • 예를 들어, 주문이 생성되면 결제와 재고 차감은 동시에 가능하지만, 배송은 반드시 포장 이후에만 가능해야 한다.
  • 따라서 이벤트 중재자(Mediator)를 두어 워크플로우를 조율할 필요가 있다.
  • Mediator는 특정 이벤트가 발생하면 다음 단계의 커맨드를 비동기로 발행해야한다.
    • OrderCreated → PayCommand, ReserveStockCommand, NotifyAcceptedCommand
    • InventoryReserved → PackCommand
    • OrderPacked → ShipCommand, NotifyShippingReady
    • OrderShipped → NotifyShipped
      Mediator는 업무 플로우의 순서와 병렬성을 관리하는 조정자 역할을 한다.

3. 이벤트 저장소 (Event Store) 관점

  • 실제 운영에서는 네트워크 오류, 모듈 장애 등으로 이벤트가 실패하거나 놓칠 수 있다.
  • 이를 대비해 모든 이벤트를 Event Store에 기록하고, 상태(READY, PROCESSED, FAILED)를 관리했다.
  • 실패한 이벤트는 저장소에 남아 있으므로 Replay를 통해 다시 처리할 수 있다.
  • 예: 재고 차감 이벤트가 실패하면, Event Store에서 다시 꺼내 브로커에 발행하여 재처리한다.
    이를 통해 발행 보장성과 장애 복구력을 확보한다.

4. 멱등성 (Idempotency) 관점

  • 재처리를 가능하게 만들면, 동일 이벤트가 중복 실행될 위험이 생긴다.
  • 예: IventoryReserved 이벤트가 두 번 처리되면 재고가 두 번 차감될 수 있다.
  • 이를 막기 위해 각 모듈에서 이미 처리한 이벤트 ID는 무시하는 멱등성 로직을 넣었다.
    이로써 재처리와 중복 처리 모두 대응할 수 있다.

최종 흐름

  1. 주문 발송

    • OrderCreated 이벤트 발행 -> Event Store에 기록
  2. 주문 처리 (비동기)

    • Mediator가 이벤트를 받고 결제/재고 차감/알림 발송 커맨드를 병렬 실행
    • 각 모듈은 스레드로 독립 실행 (MSA를 흉내)
  3. 주문 이행

    • 재고 차감 성공 이벤트 -> Mediator가 포장 커맨드 발행
  4. 주문 배송

    • 포장 완료 이벤트 -> Mediator가 배송 커멘드 + 배송 준비 알림 발행
  5. 고객 알림

    • 배송 완료 이벤트 -> Mediator가 배송 완료 알림 발행

4. 구현 방식 요약

  • Thread 기반 MSA 가정
    각 모듈(결제, 재고, 포장, 배송, 알림)은 독립적인 스레드로 동작 -> 실제 MSA처럼 독립 서비스 느낌을 낸다.

  • Broker (토픽별 큐)
    BlockingQueue를 사용하여 토픽별 이벤트 큐를 만듦.
    broker.publish(event) → 구독자가 broker.take(topic)으로 수신한다.

  • Event Store (저장/재처리)
    이벤트를 모두 저장하고 상태(READY, PROCESSED, FAILED)를 기록.
    실패한 이벤트는 replayFailed()를 통해 재처리한다.

  • Mediator (중재자/Orchestrator)
    완료 이벤트를 구독하여 다음 커맨드를 비동기로 발행.
    -> OrderCreated -> Pay, ReserveStock, NotifyAccepted
    -> InventoryReserved -> Pack
    -> OrderPacked -> Ship, NotifyShippingReady
    -> OrderShipped->NotifyShipped

  • Command 패턴 (비동기 실행 로그)
    Mediatorsms AsyncCommand를 통해 커맨드를 발행하고,
    로그에 [COMMAND][ASYNC], [COMMAND][RUN], [COMMAND][DONE]이 찍혀서 비동기성이 드러남.


5. 주요 코드

Event Store (재처리)

static final class EventStore {
    private final Map<Long, Event> byId = new ConcurrentHashMap<>();
    private final Map<Status, Set<Long>> byStatus = new ConcurrentHashMap<>();

    public void append(Event e) { ... }  // 이벤트 기록
    public void markProcessed(long id) { ... }
    public void markFailed(long id, String error) { ... }

    public void replayFailed(java.util.function.Consumer<Event> sink) {
        System.out.println("[REPROCESS] scanning FAILED events…");
        for (Event e : findByStatus(Status.FAILED)) {
            System.out.printf("[REPROCESS] replay %s (reason=%s)%n", e, e.error);
            updateStatus(e.id, Status.READY, null);
            sink.accept(e); // 브로커에 다시 발행
        }
    }
}

실패한 이벤트를 다시 READY로 바꾸고 브로커에 발행 -> 재처리 가능.

Mediator (중재자)

public void onOrderCreated(Event e) {
    issue(new AsyncCommand("PayCommand",
        () -> emitCmd(Topic.CMD_PAY, e), commandsPool)).executeAsync();

    issue(new AsyncCommand("ReserveStockCommand",
        () -> emitCmd(Topic.CMD_RESERVE_STOCK, e), commandsPool)).executeAsync();

    issue(new AsyncCommand("NotifyAcceptedCommand",
        () -> emitCmdWith(Topic.CMD_NOTIFY, e, Map.of("notifyType", "ORDER_ACCEPTED")),
        commandsPool)).executeAsync();
}

주문 생성 이벤트가 들어오면 동시에 결제/재고 차감/알림 커맨드를 발행.

모듈 (구독자)

static final class InventoryModule extends Module {
    private final Random rnd = new Random();
    protected void handle(Event e) throws Exception {
        sleep(300);
        // 30% 확률로 실패 (재처리 데모용)
        if (rnd.nextInt(10) < 3) throw new RuntimeException("재고 차감 실패(일시적)");
        System.out.printf("[%s] 재고 차감 성공 (orderId=%s)%n", id, e.payload.get("orderId"));
    }
}

재고 모듈은 일부러 실패 확률을 넣어두어, 재처리 기능이 동작하는 걸 확인할 수 있다.

실행 결과

전체 로그

[BROKER][PUB] Event{id=1, topic=EVT_ORDER_CREATED, status=PROCESSED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][ASYNC] PayCommand - submitted
[COMMAND][ASYNC] ReserveStockCommand - submitted
[COMMAND][RUN   ] PayCommand - running
[COMMAND][RUN   ] ReserveStockCommand - running
[BROKER][PUB] Event{id=2, topic=CMD_PAY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][DONE  ] PayCommand - completed
[payment#2] start: Event{id=2, topic=CMD_PAY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][ASYNC] NotifyAcceptedCommand - submitted
[BROKER][PUB] Event{id=3, topic=CMD_RESERVE_STOCK, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][DONE  ] ReserveStockCommand - completed
[inventory#2] start: Event{id=3, topic=CMD_RESERVE_STOCK, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][RUN   ] NotifyAcceptedCommand - running
[BROKER][PUB] Event{id=4, topic=CMD_NOTIFY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=ORDER_ACCEPTED}}
[COMMAND][DONE  ] NotifyAcceptedCommand - completed
[notify#1] start: Event{id=4, topic=CMD_NOTIFY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=ORDER_ACCEPTED}}
[notify#1] 고객 알림 발송 (type=ORDER_ACCEPTED, orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=5, topic=EVT_CUSTOMER_NOTIFIED, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=ORDER_ACCEPTED}}
[notify#1] done -> EVT_CUSTOMER_NOTIFIED
[MEDIATOR] 고객 알림 완료 이벤트 수신: Event{id=5, topic=EVT_CUSTOMER_NOTIFIED, status=PROCESSED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=ORDER_ACCEPTED}}
[payment#2] 결제 완료 (orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=6, topic=EVT_PAYMENT_COMPLETED, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[payment#2] done -> EVT_PAYMENT_COMPLETED
[inventory#2] FAIL: Event{id=3, topic=CMD_RESERVE_STOCK, status=FAILED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}} (재고 차감 실패(일시적))
[REPROCESS] scanning FAILED events…
[REPROCESS] replay Event{id=3, topic=CMD_RESERVE_STOCK, status=FAILED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}} (reason=재고 차감 실패(일시적))
[BROKER][PUB] Event{id=3, topic=CMD_RESERVE_STOCK, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[inventory#1] start: Event{id=3, topic=CMD_RESERVE_STOCK, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[inventory#1] 재고 차감 성공 (orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=7, topic=EVT_INVENTORY_RESERVED, status=PROCESSED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[inventory#1] done -> EVT_INVENTORY_RESERVED
[COMMAND][ASYNC] PackCommand - submitted
[COMMAND][RUN   ] PackCommand - running
[BROKER][PUB] Event{id=8, topic=CMD_PACK, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][DONE  ] PackCommand - completed
[packing#1] start: Event{id=8, topic=CMD_PACK, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[packing#1] 포장 완료 (orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=9, topic=EVT_ORDER_PACKED, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[packing#1] done -> EVT_ORDER_PACKED
[COMMAND][ASYNC] ShipCommand - submitted
[COMMAND][RUN   ] ShipCommand - running
[BROKER][PUB] Event{id=10, topic=CMD_SHIP, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][DONE  ] ShipCommand - completed
[COMMAND][ASYNC] NotifyShippingReady - submitted
[shipping#2] start: Event{id=10, topic=CMD_SHIP, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[COMMAND][RUN   ] NotifyShippingReady - running
[BROKER][PUB] Event{id=11, topic=CMD_NOTIFY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPING_READY}}
[COMMAND][DONE  ] NotifyShippingReady - completed
[notify#2] start: Event{id=11, topic=CMD_NOTIFY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPING_READY}}
[notify#2] 고객 알림 발송 (type=SHIPPING_READY, orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=12, topic=EVT_CUSTOMER_NOTIFIED, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPING_READY}}
[notify#2] done -> EVT_CUSTOMER_NOTIFIED
[MEDIATOR] 고객 알림 완료 이벤트 수신: Event{id=12, topic=EVT_CUSTOMER_NOTIFIED, status=PROCESSED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPING_READY}}
[shipping#2] 배송 처리 (orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=13, topic=EVT_ORDER_SHIPPED, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123}}
[shipping#2] done -> EVT_ORDER_SHIPPED
[COMMAND][ASYNC] NotifyShipped - submitted
[COMMAND][RUN   ] NotifyShipped - running
[BROKER][PUB] Event{id=14, topic=CMD_NOTIFY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPED}}
[COMMAND][DONE  ] NotifyShipped - completed
[notify#1] start: Event{id=14, topic=CMD_NOTIFY, status=READY, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPED}}
[notify#1] 고객 알림 발송 (type=SHIPPED, orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34)
[BROKER][PUB] Event{id=15, topic=EVT_CUSTOMER_NOTIFIED, status=PROCESSED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPED}}
[notify#1] done -> EVT_CUSTOMER_NOTIFIED
[MEDIATOR] 고객 알림 완료 이벤트 수신: Event{id=15, topic=EVT_CUSTOMER_NOTIFIED, status=PROCESSED, payload={orderId=782b4c7b-0f2f-42a1-a583-3d5e0160bb34, userId=u-123, notifyType=SHIPPED}}

0) 주문 생성(트리거)

  • [BROKER][PUB] Event{id=1, topic=EVT_ORDER_CREATED, status=PROCESSED, ...}
    • 주문 생성 이벤트가 발행됨. 곧바로 Mediator가 이 이벤트를 소비해 PROCESSED로 마킹(비동기이므로 출력 시점에 PROCESSED로 보일 수 있음).

1) 주문 처리(접수 알림·결제·재고 차감) — 병렬

  • [COMMAND][ASYNC] PayCommand - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=2, topic=CMD_PAY, ...}
  • [payment#2] start: ... → 결제 완료 → [BROKER][PUB]
    EVT_PAYMENT_COMPLETED(id=6) → done
    • Mediator가 결제 커맨드를 비동기로 제출 -> 결제 워커가 소비 -> 결제 완료 이벤트(id=6) 발생.
  • [COMMAND][ASYNC] ReserveStockCommand - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=3, topic=CMD_RESERVE_STOCK, ...}
  • [inventory#2] start: ... → FAIL ... (재고 차감 실패(일시적))
    • Mediator가 재고 차감 커맨드 제출 -> 첫 시도는 실패로 FAILED 상태 기록.
  • [COMMAND][ASYNC] NotifyAcceptedCommand - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=4, topic=CMD_NOTIFY, ... notifyType=ORDER_ACCEPTED}
  • [notify#1] start: ... → 고객 알림 발송(ORDER_ACCEPTED)
  • [BROKER][PUB] EVT_CUSTOMER_NOTIFIED(id=5)
  • [MEDIATOR] 고객 알림 완료 이벤트 수신: id=5
    • Mediator가 주문 접수 알림 커맨드 제출 -> 알림 워커가 처리 -> 고객 알림 완료 이벤트(id=5) 발생 및 Mediator가 수신(후속 플로우용 체크).

      요약: 주문 생성 후, 결제/재고.접수알림이 동시에 진행된다. 결제 성공, 알림 성공, 재고만 1차 실패.

2) 실패 재처리(Replay)로 재고 차감 재시도 -> 성공

  • [REPROCESS] scanning FAILED events…
  • [REPROCESS] replay Event{id=3, ...} (reason=재고 차감 실패(일시적))
  • [BROKER][PUB] Event{id=3, topic=CMD_RESERVE_STOCK, status=READY, ...}
  • [inventory#1] start: id=3 ... → 재고 차감 성공
  • [BROKER][PUB] EVT_INVENTORY_RESERVED(id=7) → done
    • Event Store가 FAILED 이벤트(id=3) 를 다시 READY로 되돌려 재발행 → 이번엔 다른 워커가 성공 → 재고 확보 완료 이벤트(id=7) 발생.

3) 주문 이행(포장)

  • [COMMAND][ASYNC] PackCommand - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=8, topic=CMD_PACK, ...}
  • [packing#1] start: id=8 ... → 포장 완료
  • [BROKER][PUB] EVT_ORDER_PACKED(id=9) → done
    • Mediator는 EVT_INVENTORY_RESERVED(id=7) 를 수신한 뒤에만 포장 커맨드를 발행한다(의존성 보장). 포장 완료 이벤트(id=9) 발생.

4) 주문 배송(배송 처리 + 배송 준비 알림)

  • [COMMAND][ASYNC] ShipCommand - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=10, topic=CMD_SHIP, ...}
  • [shipping#2] start: id=10 ... → 배송 처리
  • [BROKER][PUB] EVT_ORDER_SHIPPED(id=13) → done
    • Mediator는 포장 완료(id=9) 를 보고 배송 커맨드를 발행 → 배송 완료 이벤트(id=13) 발생.
  • [COMMAND][ASYNC] NotifyShippingReady - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=11, topic=CMD_NOTIFY, ... notifyType=SHIPPING_READY}
  • [notify#2] start: id=11 ... → 고객 알림 발송(SHIPPING_READY)
  • [BROKER][PUB] EVT_CUSTOMER_NOTIFIED(id=12) → done
  • [MEDIATOR] 고객 알림 완료 이벤트 수신: id=12
    • 동시에 배송 준비 알림도 발행/처리되어 고객에게 안내.

      요약: 포장 이후 Mediator가 배송과 배송준비 알림을 동시에 흘려보낸다.

5) 배송 알림

  • [COMMAND][ASYNC] NotifyShipped - submitted / RUN / DONE
  • [BROKER][PUB] Event{id=14, topic=CMD_NOTIFY, ... notifyType=SHIPPED}
  • [notify#1] start: id=14 ... → 고객 알림 발송(SHIPPED)
  • [BROKER][PUB] EVT_CUSTOMER_NOTIFIED(id=15) → done
  • [MEDIATOR] 고객 알림 완료 이벤트 수신: id=15
    • Mediator는 EVT_ORDER_SHIPPED(id=13) 를 보고 마지막 배송 완료 알림을 발행하고, 알림 완료 이벤트(id=15)를 수신합니다.

정리

  • 의존성/순서 보장: Mediator가 완료 이벤트(예: EVT_INVENTORY_RESERVED, EVT_ORDER_PACKED, EVT_ORDER_SHIPPED)를 수신한 뒤 다음 커맨드를 발행한다. 그래서 “재고 확보 전에는 포장 X, 포장 전에는 배송 X”가 지켜진다.
  • 병렬성: 가능한 구간(접수 알림/결제/재고, 포장 후 배송·알림)은 동시에 커맨드가 발행된다.
  • 재처리(Replay): 실패한 재고 차감(id=3)을 Event Store가 READY로 되감고 재발행해서 성공시킨다.
  • 멱등성: 모듈은 이미 처리한 이벤티 ID를 건너뛰도록 설계되어 중복 실행을 방지한다.

데모 아키텍처

의도

  • OrderService: 클라이언트 요청을 처리하는 API 서버
  • Mediator Service: 오케스트레이션 서비스
  • 도메인 서비스 모듈: 각각은 독립된 MSA 서버
    하나의 모듈은 2개의 Thread로 이루어져있다.
    이러한 구조의 의도는 모듈은 MSA 서버로 scale-out이 가능하게 구성했다.

레퍼런스

회원시스템 이벤트기반 아키텍처 구축하기
[소프트웨어 아키텍처 101] Ch. 14 이벤트 기반 아키텍처 스타일

profile
공부 내용을 가볍게 적어놓는 블로그.

0개의 댓글