Watermill 라이브러리로 구현하는 Golang CQRS

CQRS 개요

CQRS는 명령-질의 책임 분리(Command-Query Responsibility Segregation)의 약자입니다. 이 패턴의 핵심은 명령(쓰기 요청)과 질의(읽기 요청)를 서로 다른 객체가 처리하도록 분리하는 데 있습니다. 이것이 기본 개념입니다.

이를 더 확장하여 데이터 저장소를 분리하고, 읽기 전용 저장소와 쓰기 전용 저장소를 별도로 사용할 수 있습니다. 이렇게 하면 다양한 질의 유형을 처리하거나 여러 바운디드 컨텍스트(Bounded Context)를 지원하기 위해 최적화된 여러 읽기 저장소를 가질 수 있습니다. CQRS와 관련하여 읽기/쓰기 저장소 분리가 자주 논의되지만, 이것은 CQRS 자체의 필수 조건은 아닙니다. CQRS는 기본적으로 명령과 질의를 분리하는 것에서 시작합니다.

주요 용어

명령(Command)

명령은 어떤 작업을 수행하라는 요청을 나타내는 단순한 데이터 구조체입니다.

명령 버스(Command Bus)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/command_bus.go

// CommandBus는 명령(commands)을 명령 핸들러(command handlers)로 전송합니다.
type CommandBus struct {
    // ...
}

명령 프로세서(Command Processor)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/command_processor.go

// CommandProcessor는 명령 버스로부터 수신된 명령을 처리할 
// 적절한 CommandHandler를 결정합니다.
type CommandProcessor struct {
    // ...
}

명령 핸들러(Command Handler)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/command_processor.go

// CommandHandler는 NewCommand에서 정의한 명령을 수신하고 
// Handle 메서드를 사용하여 처리합니다.
// DDD를 사용하는 경우 CommandHandler는 집계(Aggregate)를 수정하고 영속화할 수 있습니다.
//
// EventHandler와 달리 각 명령에는 하나의 CommandHandler만 있어야 합니다.
//
// 메시지 처리 중에는 CommandHandler의 단일 인스턴스가 사용됩니다.
// 여러 명령이 동시에 전송되면 Handle 메서드가 동시에 여러 번 실행될 수 있습니다.
// 따라서 Handle 메서드는 스레드 안전(thread-safe)해야 합니다!
type CommandHandler interface {
    // ...
}

이벤트(Event)

이벤트는 이미 발생한 사실을 나타냅니다. 이벤트는 불변(immutable)합니다.

이벤트 버스(Event Bus)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/event_bus.go

// EventBus는 이벤트를 이벤트 핸들러로 전송합니다.
type EventBus struct {
    // ...
}

이벤트 프로세서(Event Processor)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/event_processor.go

// EventProcessor는 이벤트 버스로부터 수신된 이벤트를 처리할 
// EventHandler를 결정합니다.
type EventProcessor struct {
    // ...
}

이벤트 핸들러(Event Handler)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/event_processor.go

// EventHandler는 NewEvent에서 정의한 이벤트를 수신하고 Handle 메서드로 처리합니다.
// DDD를 사용하는 경우 CommandHandler가 집계를 수정하고 유지할 수 있습니다.
// 또한 프로세스 관리자, Saga를 호출하거나 읽기 모델(Read Model)을 구축할 수도 있습니다.
// CommandHandler와 달리 각 이벤트에는 여러 EventHandler가 있을 수 있습니다.
//
// 메시지 처리 중에는 EventHandler의 단일 인스턴스가 사용됩니다.
// 여러 이벤트가 동시에 전달되면 Handle 메서드가 동시에 여러 번 실행될 수 있습니다.
// 따라서 Handle 메서드는 스레드 안전해야 합니다!
type EventHandler interface {
    // ...
}

CQRS 퍼사드(Facade)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/cqrs.go

// Facade는 Command 및 Event 버스와 프로세서를 생성하기 위한 퍼사드입니다.
// 표준 방식으로 CQRS를 사용할 때 보일러플레이트(boilerplate) 코드를 피하기 위해 생성되었습니다.
// 버스와 프로세서를 수동으로 생성할 수도 있으며, NewFacade에서 영감을 얻을 수 있습니다.
type Facade struct {
    // ...
}

명령 및 이벤트 마샬러(Marshaler)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/components/cqrs/marshaler.go

// CommandEventMarshaler는 명령과 이벤트를 Watermill 메시지로 마샬링하거나
// 그 반대의 작업을 수행합니다. 명령의 페이로드는 []bytes로 마샬링되어야 합니다.
type CommandEventMarshaler interface {
    // Marshal은 명령 또는 이벤트를 Watermill 메시지로 마샬링합니다.
    Marshal(v interface{}) (*message.Message, error)

    // Unmarshal은 Watermill 메시지를 명령 또는 이벤트 v로 언마샬링합니다.
    Unmarshal(msg *message.Message, v interface{}) (err error)

    // Name은 명령 또는 이벤트의 이름을 반환합니다.
    // Name은 수신된 명령 또는 이벤트가 우리가 처리하려는 것인지 확인하는 데 사용됩니다.
    Name(v interface{}) string

    // NameFromMessage는 Watermill 메시지(마샬링에 의해 생성됨)로부터 
    // 명령 또는 이벤트의 이름을 반환합니다.
    //
    // 명령 또는 이벤트를 Watermill 메시지로 마샬링할 때 
    // 불필요한 언마샬링을 피하기 위해 Name 대신 NameFromMessage를 사용해야 합니다.
    NameFromMessage(msg *message.Message) string
}

사용 예제

예제 도메인

호텔 객실 예약을 처리하는 간단한 도메인을 예제로 사용하겠습니다.

도메인은 간단합니다:

  • 고객이 객실을 예약할 수 있습니다(book a room).
  • 객실이 예약될 때마다 고객을 위해 맥주를 주문합니다(Whenever a room is booked, we order a beer).
  • 예약 정보를 기반으로 재무 보고서(financial report)를 생성합니다.

명령 전송(Sending a command)

전체 소스 코드: github.com/ThreeDotsLabs/watermill/_examples/basic/5-cqrs-protobuf/main.go

func publishCommands(commandBus *cqrs.CommandBus) {
    for i := 0; ; i++ {
        startDate := time.Now().Add(time.Duration(i) * time.Hour * 24)
        endDate := startDate.Add(time.Hour * 48)

        bookRoomCmd := &BookRoom{
            RoomId:    fmt.Sprintf("%d", i),
            GuestName: "John",
            StartDate: timestamppb.New(startDate),
            EndDate:   timestamppb.New(endDate),
        }

        if err := commandBus.Send(context.Background(), bookRoomCmd); err != nil {
            panic(err)
        }

        time.Sleep(1 * time.Second)
    }
}

명령 핸들러(Command handler)

BookRoomHandlerBookRoom 명령을 처리하고 RoomBooked 이벤트를 발생시킵니다.

type BookRoomHandler struct {
    eventBus *cqrs.EventBus
}

func (b BookRoomHandler) HandlerName() string {
    return "BookRoomHandler"
}

func (b BookRoomHandler) NewCommand() interface{} {
    return &BookRoom{}
}

func (b BookRoomHandler) Handle(ctx context.Context, c interface{}) error {
    cmd := c.(*BookRoom)
    price := (rand.Int63n(40) + 1) * 10

    log.Printf("객실 %s이(가) %s님에 의해 %s부터 %s까지 예약되었습니다.",
        cmd.RoomId,
        cmd.GuestName,
        cmd.StartDate.AsTime(),
        cmd.EndDate.AsTime(),
    )

    return b.eventBus.Publish(ctx, &RoomBooked{
        ReservationId: watermill.NewUUID(),
        RoomId:        cmd.RoomId,
        GuestName:     cmd.GuestName,
        Price:         price,
        StartDate:     cmd.StartDate,
        EndDate:       cmd.EndDate,
    })
}

이벤트 핸들러(Event handler)

OrderBeerOnRoomBooked 이벤트 핸들러는 RoomBooked 이벤트를 수신하여 OrderBeer 명령을 발행합니다.

type OrderBeerOnRoomBooked struct {
    commandBus *cqrs.CommandBus
}

func (o OrderBeerOnRoomBooked) HandlerName() string {
    return "OrderBeerOnRoomBooked"
}

func (OrderBeerOnRoomBooked) NewEvent() interface{} {
    return &RoomBooked{}
}

func (o OrderBeerOnRoomBooked) Handle(ctx context.Context, e interface{}) error {
    event := e.(*RoomBooked)
    orderBeerCmd := &OrderBeer{
        RoomId: event.RoomId,
        Count:  rand.Int63n(10) + 1,
    }
    return o.commandBus.Send(ctx, orderBeerCmd)
}

읽기 모델 구축(Event handler로 읽기 모델 구성)

BookingsFinancialReportRoomBooked 이벤트를 수신하여 재무 보고서를 업데이트하는 읽기 모델입니다.

type BookingsFinancialReport struct {
    handledBookings map[string]struct{}
    totalCharge     int64
    lock            sync.Mutex
}

func NewBookingsFinancialReport() *BookingsFinancialReport {
    return &BookingsFinancialReport{handledBookings: map[string]struct{}{}}
}

func (b BookingsFinancialReport) HandlerName() string {
    return "BookingsFinancialReport"
}

func (BookingsFinancialReport) NewEvent() interface{} {
    return &RoomBooked{}
}

func (b *BookingsFinancialReport) Handle(ctx context.Context, e interface{}) error {
    b.lock.Lock()
    defer b.lock.Unlock()

    event := e.(*RoomBooked)

    if _, ok := b.handledBookings[event.ReservationId]; ok {
        return nil
    }
    b.handledBookings[event.ReservationId] = struct{}{}
    b.totalCharge += event.Price

    fmt.Printf(">>> 현재까지 예약된 객실 총액: $%d\n", b.totalCharge)
    return nil
}

CQRS 퍼사드로 연결하기

모든 구성 요소를 cqrs.Facade를 사용하여 연결합니다.

const amqpAddress = "amqp://guest:guest@rabbitmq:5672/"

func main() {
    logger := watermill.NewStdLogger(false, false)
    cqrsMarshaler := cqrs.ProtobufMarshaler{}

    // AMQP 구독자 및 발행자 설정
    commandsAMQPConfig := amqp.NewDurableQueueConfig(amqpAddress)
    commandsPublisher, err := amqp.NewPublisher(commandsAMQPConfig, logger)
    if err != nil {
        panic(err)
    }
    commandsSubscriber, err := amqp.NewSubscriber(commandsAMQPConfig, logger)
    if err != nil {
        panic(err)
    }

    eventsPublisher, err := amqp.NewPublisher(
        amqp.NewDurablePubSubConfig(amqpAddress, nil), logger)
    if err != nil {
        panic(err)
    }

    // 메시지 라우터 생성
    router, err := message.NewRouter(message.RouterConfig{}, logger)
    if err != nil {
        panic(err)
    }
    router.AddMiddleware(middleware.Recoverer)

    // Facade 생성
    cqrsFacade, err := cqrs.NewFacade(cqrs.FacadeConfig{
        GenerateCommandsTopic: func(commandName string) string {
            return commandName
        },
        CommandHandlers: func(cb *cqrs.CommandBus, eb *cqrs.EventBus) []cqrs.CommandHandler {
            return []cqrs.CommandHandler{
                BookRoomHandler{eventBus: eb},
                OrderBeerHandler{eventBus: eb},
            }
        },
        CommandsPublisher: commandsPublisher,
        CommandsSubscriberConstructor: func(handlerName string) (message.Subscriber, error) {
            return commandsSubscriber, nil
        },
        GenerateEventsTopic: func(eventName string) string {
            return "events"
        },
        EventHandlers: func(cb *cqrs.CommandBus, eb *cqrs.EventBus) []cqrs.EventHandler {
            return []cqrs.EventHandler{
                OrderBeerOnRoomBooked{commandBus: cb},
                NewBookingsFinancialReport(),
            }
        },
        EventsPublisher: eventsPublisher,
        EventsSubscriberConstructor: func(handlerName string) (message.Subscriber, error) {
            config := amqp.NewDurablePubSubConfig(
                amqpAddress,
                amqp.GenerateQueueNameTopicNameWithSuffix(handlerName),
            )
            return amqp.NewSubscriber(config, logger)
        },
        Router:                router,
        CommandEventMarshaler: cqrsMarshaler,
        Logger:                logger,
    })
    if err != nil {
        panic(err)
    }

    // 명령 전송 고루틴 시작
    go publishCommands(cqrsFacade.CommandBus())

    // 라우터 실행 (프로세서 자동 시작)
    if err := router.Run(context.Background()); err != nil {
        panic(err)
    }
}

이제 CQRS 애플리케이션이 실행됩니다.

태그: Watermill Golang CQRS CommandBus eventbus

7월 30일 02:18에 게시됨