1.시스템&인프라/redis

8편. Redis Stream 실습: 이벤트 로그와 비동기 작업 큐

쿼드큐브 2026. 7. 27. 17:37
반응형
반응형

 

8편. Redis Stream 실습: 이벤트 로그와 비동기 작업 큐

 

📚 목차
1. 주문 이벤트 저장하기
2. 알림 이벤트 큐 만들기
3. 이메일 작업 큐 구현하기
4. 감사 로그 저장하기

 

 

Redis Stream은 Redis에서 이벤트 로그와 비동기 작업 큐를 구현할 때 사용할 수 있는 자료구조입니다.

List와 Pub/Sub도 메시지 처리에 사용할 수 있지만, Stream은 두 자료구조와 성격이 다릅니다.

 

📂 [GitHub 코드 보러가기] : https://github.com/cericube/nodejs-practice-lab/tree/main/redis-examples

 

1. 주문 이벤트 저장하기

주문 서비스에서는 주문 생성, 결제 완료, 배송 시작, 주문 취소 같은 이벤트가 계속 발생합니다.

예를 들어 주문 생성 시 다음과 같은 후속 처리가 필요할 수 있습니다.

주문 생성
  ↓
주문 생성 이벤트 저장
  ↓
재고 차감
  ↓
결제 요청
  ↓
알림 발송
  ↓
감사 로그 저장

이때 모든 작업을 주문 생성 API 안에서 동기적으로 처리하면 API 응답이 느려지고, 중간 작업 실패 시 흐름이 복잡해집니다.
Redis Stream을 사용하면 주문 생성 자체는 DB에 저장하고, 후속 처리를 위한 이벤트는 Stream에 기록할 수 있습니다.

OrderService
  ↓
DB에 주문 저장
  ↓
stream:orders에 주문 이벤트 저장
  ↓
worker가 stream:orders를 읽어서 후속 처리

 

✔️ 주문 생성

// src/services/order-stream.service.ts

/**
 * 주문을 생성하고 주문 생성 이벤트를 기록합니다.
 *
 * 1. Order 테이블에 주문을 저장합니다.
 * 2. 주문 생성 결과를 기준으로 Redis Stream에 order.created 이벤트를 기록합니다.
 * 3. 생성된 주문 정보를 반환합니다.
 *
 * 실습 포인트:
 * DB는 현재 주문 상태의 원본 저장소입니다.
 * Redis Stream은 주문 생성 사실을 다른 worker나 서비스가 나중에 처리할 수 있도록 남기는 이벤트 로그입니다.
 *
 * 참고:
 * DB 저장 후 Stream 기록에 실패할 수 있으므로 실무에서는 Outbox Pattern 등으로 두 저장소의 정합성을 보완할 수 있습니다.
 */
async createOrder(input: CreateOrderInput) {
  const order = await prisma.order.create({
    data: {
      userId: input.userId,
      totalPrice: input.totalPrice,
      status: 'CREATED',
    },
  });

  await this.addOrderEvent({
    eventType: 'order.created',
    orderId: order.id,
    userId: order.userId,
    status: order.status,
    totalPrice: order.totalPrice,
  });

  return order;
}

/**
 * 주문 이벤트를 Redis Stream에 추가합니다.
 *
 * 1. Redis Stream key를 가져옵니다.
 * 2. XADD 명령으로 주문 이벤트를 추가합니다.
 * 3. '*'를 사용하면 Redis가 Stream 메시지 ID를 자동 생성합니다.
 *
 * 실습 포인트:
 * Stream 메시지는 명시적으로 삭제하거나 보존 길이를 제한하지 않는 한 Redis에 로그처럼 남습니다.
 * 따라서 나중에 XRANGE로 다시 조회하거나 Consumer Group으로 처리할 수 있습니다.
 */
async addOrderEvent(input: {
  eventType: OrderEventType;
  orderId: number;
  userId: number;
  status: string;
  totalPrice: number;
}): Promise<string> {
  const key = RedisKey.stream.orders();

  // Redis 명령: XADD stream:orders * eventType value orderId value userId value status value totalPrice value createdAt value
  // 주문 이벤트를 Stream 끝에 추가하고 Redis가 자동 생성한 메시지 ID를 반환합니다.
  const messageId = await redis.xAdd(key, '*', {
    eventType: input.eventType,
    orderId: String(input.orderId),
    userId: String(input.userId),
    status: input.status,
    totalPrice: String(input.totalPrice),
    createdAt: new Date().toISOString(),
  });

  return messageId;
}

 

 

✔️ Stream 이벤트 처리 예

실제로 후속 처리를 하려면 별도의 worker가 필요합니다.
예를 들어 이런 worker가 있을 수 있습니다.

// 예시: order-worker.ts

async function orderWorker() {
  const orderStreamService = new OrderStreamService();

  const events = await orderStreamService.getOrderEvents(10);

  for (const event of events) {
    if (event.eventType === 'order.created') {
      console.log('재고 차감 처리:', event.orderId);
      console.log('알림 발송 처리:', event.userId);
      console.log('이메일 발송 처리:', event.userId);
      console.log('감사 로그 저장:', event.orderId);
    }
  }
}

흐름은 다음과 같습니다.

 

✔️ 주문 상태를 변경

Redis Stream은 이벤트 흐름을 저장하는 곳이고, DB는 현재 주문 상태를 저장하는 곳입니다.
역할을 분리해서 봐야 합니다.

/**
 * 주문 상태를 변경하고 상태 변경 이벤트를 기록합니다.
 *
 * 1. DB에서 지정한 주문의 상태를 수정합니다.
 * 2. 변경된 상태에 대응하는 이벤트 종류를 결정합니다.
 * 3. 수정된 주문 정보를 Redis Stream에 기록합니다.
 *
 * 실습 포인트:
 * DB에는 현재 상태를 저장하고 Stream에는 상태가 변경된 이력을 순서대로 남깁니다.
 *
 * 참고:
 * 정의되지 않은 상태는 order.created로 기록되므로 실무에서는 허용 상태를 검증하거나 상태 타입을 제한해야 합니다.
 */
async changeOrderStatus(orderId: number, status: string): Promise<void> {
  const order = await prisma.order.update({
    where: {
      id: orderId,
    },
    data: {
      status,
    },
  });

  const eventType: OrderEventType =
    status === 'PAID'
      ? 'order.paid'
      : status === 'CANCELLED'
        ? 'order.cancelled'
        : status === 'SHIPPED'
          ? 'order.shipped'
          : 'order.created';

  await this.addOrderEvent({
    eventType,
    orderId: order.id,
    userId: order.userId,
    status: order.status,
    totalPrice: order.totalPrice,
  });
}

 

예를 들어 결제 완료 처리를 하면 다음과 같습니다.

await orderStreamService.changeOrderStatus(1, 'PAID');

그러면 DB는 이렇게 바뀝니다.

Order Table

id | status
1  | PAID

그리고 Stream에는 이벤트가 추가됩니다.

stream:orders

eventType order.paid
orderId   1
status    PAID

 

이 이벤트를 읽은 worker는 다음 후속 작업을 할 수 있습니다.

order.paid 이벤트 수신
  ↓
결제 완료 알림 발송
  ↓
주문 완료 이메일 발송
  ↓
배송 준비 작업 생성
  ↓
감사 로그 저장

따라서 상태는 DB에 저장하고, 상태 변경 사실은 Stream에 기록하는 구조가 됩니다.

 

2. 알림 이벤트 큐 만들기

알림은 사용자가 직접 기다릴 필요가 없는 작업입니다.
예를 들어 다음과 같은 알림은 API 응답과 분리해서 처리해도 됩니다.

- 주문 생성 알림
- 댓글 작성 알림
- 좋아요 알림
- 관리자 공지 알림

Redis Stream을 사용하면 알림 이벤트를 Queue처럼 저장하고, worker가 하나씩 읽어서 처리할 수 있습니다.

 

 

✔️ Consumer Group

단순히 xRange()로 Stream을 읽으면 모든 worker가 같은 메시지를 읽을 수 있습니다.
하지만 Consumer Group을 사용하면 Redis가 메시지 분배 상태를 관리합니다.

따라서 알림 발송, 이메일 발송, 주문 후처리처럼 여러 worker가 나눠 처리해야 하는 작업은 Consumer Group을 사용하는 것이 좋습니다.

 

✔️ 알림 이벤트 흐름


✔️ 알림 및 ConsumerGroup 생성

// src/services/notification-stream.service.ts

/**
 * 알림 작업을 Redis Stream에 추가합니다.
 *
 * 1. 알림 대상 사용자와 알림 내용을 Stream 메시지 필드로 구성합니다.
 * 2. 사용자 ID를 Redis에 저장할 문자열로 변환합니다.
 * 3. Redis가 생성한 메시지 ID를 반환합니다.
 *
 * 실습 포인트:
 * Stream에 기록된 작업은 worker가 즉시 실행 중이지 않아도 나중에 Consumer Group으로 읽을 수 있습니다.
 */
async addNotificationEvent(input: NotificationEventInput): Promise<string> {
  const key = RedisKey.stream.notifications();

  // Redis 명령: XADD stream:notifications 
  // 알림 작업을 Stream 끝에 추가하고 Redis가 자동 생성한 메시지 ID를 반환합니다.
  return redis.xAdd(key, '*', {
    userId: String(input.userId),
    type: input.type,
    title: input.title,
    message: input.message,
    createdAt: new Date().toISOString(),
  });
}

/**
 * 알림 worker가 공유할 Consumer Group을 생성합니다.
 *
 * 1. 알림 Stream과 Consumer Group이 없으면 함께 생성합니다.
 * 2. `$`를 시작 ID로 사용해 그룹 생성 이후에 추가되는 메시지부터 처리합니다.
 * 3. 이미 그룹이 존재해서 발생한 BUSYGROUP 오류만 무시합니다.
 *
 * 실습 포인트:
 * Consumer Group을 사용하면 여러 worker가 같은 Stream의 새 메시지를 나누어 처리할 수 있습니다.
 *
 * 참고:
 * BUSYGROUP 이외의 오류는 연결 장애나 잘못된 명령일 수 있으므로 호출자에게 다시 전달합니다.
 */
async createConsumerGroup(): Promise<void> {
  const key = RedisKey.stream.notifications();

  try {
    // Redis 명령: XGROUP CREATE stream:notifications notification-workers $ MKSTREAM
    // Stream이 없으면 생성하고, 현재 마지막 메시지 다음부터 읽는 Consumer Group을 만듭니다.
    await redis.xGroupCreate(key, this.groupName, '$', {
      MKSTREAM: true,
    });
  } catch (error) {
    if (error instanceof Error && error.message.includes('BUSYGROUP')) {
      return;
    }

    throw error;
  }
}

 

✔️ 이벤트 처리

/**
 * Consumer Group에 전달되지 않은 새 알림 작업을 읽습니다.
 *
 * 1. 호출한 worker를 Consumer Group의 consumer 이름으로 사용합니다.
 * 2. `>` ID로 아직 다른 consumer에게 전달되지 않은 메시지만 요청합니다.
 * 3. 최대 count개의 메시지를 1초 동안 기다려 읽고 알림 작업 데이터로 변환합니다.
 *
 * 실습 포인트:
 * 읽은 메시지는 ACK 전까지 Consumer Group의 pending 목록에 남습니다.
 *
 * 참고:
 * 대기 시간 안에 새 메시지가 없으면 Redis가 null을 반환하므로 빈 배열로 변환합니다.
 */
async readNotificationJobs(consumerName: string, count = 10): Promise<NotificationJob[]> {
  const key = RedisKey.stream.notifications();

  // 새 메시지를 최대 count개까지 1초 동안 기다려 읽으며, 메시지가 없으면 null을 반환합니다.
  const result = await redis.xReadGroup(
    this.groupName,
    consumerName,
    [
      {
        key,
        id: '>',
      },
    ],
    {
      COUNT: count,
      BLOCK: 1000,
    },
  );

  if (!result) {
    return [];
  }

  const stream = result[0];

  if (!stream) {
    return [];
  }

  return stream.messages.map(parseNotificationJob);
}

/**
 * 처리 완료한 알림 작업을 Consumer Group에 확인 처리합니다.
 *
 * 1. 처리한 Stream 메시지 ID를 전달받습니다.
 * 2. Consumer Group에 ACK를 보내 해당 메시지를 pending 목록에서 제거합니다.
 *
 * 실습 포인트:
 * 작업이 성공한 뒤 ACK해야 worker 장애 시 미완료 작업을 pending 목록에서 확인하거나 재처리할 수 있습니다.
 */
async ackNotificationJob(messageId: string): Promise<void> {
  const key = RedisKey.stream.notifications();

  // Redis 명령: XACK stream:notifications notification-workers messageId
  // 지정한 메시지를 처리 완료로 표시하고 pending 목록에서 제거한 메시지 수를 반환합니다.
  await redis.xAck(key, this.groupName, messageId);
}

 

worker는 Consumer Group을 기준으로 메시지를 읽습니다.

여기서 중요한 값은 id: '>'입니다.

id = '>' → 아직 어떤 consumer에게도 전달되지 않은 새 메시지만 읽음
id = '0' → 내 Consumer에게 이미 전달되었지만 ACK하지 않은 Pending 메시지

즉 여러 worker가 있어도 같은 메시지를 중복해서 가져가지 않고 나눠 처리할 수 있습니다.

stream:notifications
  ↓
notification-workers group
  ↓
worker-1 → message A
worker-2 → message B
worker-3 → message C

처리가 끝나면 ACK를 보냅니다.

await redis.xAck(key, this.groupName, messageId);

ACK는 “이 메시지는 정상 처리되었다”는 표시입니다.
ACK하지 않은 메시지는 pending 상태로 남습니다.

 

 

반응형

 

3. 이메일 작업 큐 구현하기

이메일 발송은 대표적인 비동기 작업입니다.
회원가입 직후 환영 메일을 보내거나, 주문 완료 후 주문 확인 메일을 보내는 작업은 API 응답을 막지 않고 뒤에서 처리하는 것이 좋습니다.

✔️ 이메일 작업 추가

// src/services/email-stream.service.ts

/**
 * 이메일 발송 작업을 Redis Stream에 추가합니다.
 *
 * 1. 수신자, 작업 종류, 제목과 본문을 Stream 메시지 필드로 구성합니다.
 * 2. 최초 재시도 횟수를 0으로 설정하고 생성 시각을 기록합니다.
 * 3. Redis가 생성한 메시지 ID를 반환합니다.
 *
 * 실습 포인트:
 * 이메일 생성 요청과 실제 발송을 분리하면 요청 처리 중 외부 메일 서버의 응답을 기다리지 않아도 됩니다.
 */
async addEmailJob(input: EmailJobInput): Promise<string> {
  const key = RedisKey.stream.emails();

  // 이메일 발송 작업을 Stream에 저장합니다.
  // 이벤트를 추가하고 생성된 메시지 ID를 반환합니다.
  return redis.xAdd(key, '*', {
    to: input.to,
    type: input.type,
    subject: input.subject,
    body: input.body,
    retryCount: '0',
    createdAt: new Date().toISOString(),
  });
}

/**
 * 회원가입 환영 이메일 작업을 생성합니다.
 *
 * 1. 수신자 이메일과 사용자 이름을 전달받습니다.
 * 2. 환영 이메일의 작업 종류, 제목과 본문을 구성합니다.
 * 3. 공통 이메일 작업 추가 메서드에 발송을 위임합니다.
 */
async addWelcomeEmailJob(email: string, name: string): Promise<string> {
  return this.addEmailJob({
    to: email,
    type: 'welcome',
    subject: '회원가입을 환영합니다.',
    body: `${name}님, 회원가입을 환영합니다.`,
  });
}

 

이메일 발송에 필요한 정보를 Stream에 저장합니다.

stream:emails
  1710000000000-0
    to         kim@example.com
    type       welcome
    subject    회원가입을 환영합니다.
    body       Kim님, 회원가입을 환영합니다.
    retryCount 0
    createdAt  2026-06-16T00:00:00.000Z

 

✔️ 이메일 전송

/**
 * Consumer Group에 전달되지 않은 새 이메일 작업을 읽습니다.
 *
 * 1. 호출한 worker를 Consumer Group의 consumer 이름으로 사용합니다.
 * 2. `>` ID로 아직 다른 consumer에게 전달되지 않은 메시지만 요청합니다.
 * 3. 최대 count개의 메시지를 1초 동안 기다려 읽고 이메일 작업 데이터로 변환합니다.
 *
 * 실습 포인트:
 * 읽은 메시지는 ACK 전까지 Consumer Group의 pending 목록에 남습니다.
 *
 * 참고:
 * 이 메서드를 호출하기 전에 createConsumerGroup으로 Consumer Group을 준비해야 합니다.
 * 대기 시간 안에 새 메시지가 없으면 Redis가 null을 반환하므로 빈 배열로 변환합니다.
 */
async readEmailJobs(consumerName: string, count = 5): Promise<EmailJob[]> {
  const key = RedisKey.stream.emails();

  // Consumer Group에서 아직 전달되지 않은 새 이메일 발송 작업을 읽습니다.
  // COUNT만큼 최대 1초 동안 기다리며, 새 작업이 없으면 null을 반환합니다.
  const result = await redis.xReadGroup(
    this.groupName,
    consumerName,
    [
      {
        key,
        id: '>',
      },
    ],
    {
      COUNT: count,
      BLOCK: 1000,
    },
  );

  if (!result) {
    return [];
  }

  const stream = result[0];

  if (!stream) {
    return [];
  }

  return stream.messages.map(parseEmailJob);
}

/**
 * 발송을 완료한 이메일 작업을 Consumer Group에 확인 처리합니다.
 *
 * 1. 발송을 완료한 Stream 메시지 ID를 전달받습니다.
 * 2. Consumer Group에 ACK를 보내 해당 메시지를 pending 목록에서 제거합니다.
 *
 * 실습 포인트:
 * 이메일 발송에 성공한 뒤 ACK해야 장애가 발생했을 때 미완료 작업을 확인하거나 재처리할 수 있습니다.
 */
async ackEmailJob(messageId: string): Promise<void> {
  const key = RedisKey.stream.emails();

  // 처리가 끝난 이메일 발송 작업을 완료 상태로 표시합니다.
  // Pending 목록에서 제거한 메시지 수를 반환하며, 대상 메시지가 없으면 0을 반환합니다.
  await redis.xAck(key, this.groupName, messageId);
}

 

이메일 worker는 readEmailJobs()로 작업을 가져옵니다.

const jobs = await emailStreamService.readEmailJobs('email-worker-1');

작업을 처리한 뒤 성공하면 ACK를 보냅니다.

await emailStreamService.ackEmailJob(job.id);

실패하면 다시 Stream에 retry 작업을 추가할 수 있습니다.

await emailStreamService.retryEmailJob(job);

 

✔️ Stream이 이메일 큐에 적합한 이유

이메일 발송은 실패할 수 있습니다.

- SMTP 서버 장애
- 네트워크 장애
- 외부 이메일 API 장애
- 일시적인 rate limit

따라서 이메일 작업은 단순히 “보내고 끝”이 아니라 다음 처리가 필요합니다.

- 작업 저장
- worker 분산 처리
- 성공 시 ACK
- 실패 시 재시도
- pending 작업 확인

Redis Stream은 이런 흐름을 학습하기에 적합합니다.
다만 실무에서 이메일 큐를 더 안정적으로 운영해야 한다면 다음 기능도 필요합니다.

- 지연 재시도
- 최대 재시도 횟수
- dead-letter queue
- 작업 우선순위
- 작업 처리 timeout

이 수준까지 필요하다면 Redis Stream을 직접 다루기보다 BullMQ 같은 큐 라이브러리를 사용하는 것도 좋은 선택입니다.

 

4. 감사 로그 저장하기

감사 로그는 사용자의 중요한 행위나 관리자의 작업 이력을 남기는 기능입니다.
예를 들어 다음과 같은 이벤트를 기록할 수 있습니다.

- 관리자 로그인
- 사용자 권한 변경
- 주문 상태 변경
- 게시글 삭제
- 결제 취소

감사 로그는 보통 DB에 최종 저장해야 합니다.
하지만 Redis Stream을 함께 사용하면 이벤트를 먼저 빠르게 기록하고, worker가 나중에 DB에 저장하는 구조를 만들 수 있습니다.

 

 

✔️ 감사 로그 이벤트 처리 흐름

1) 감사 로그 이벤트를 Stream에 추가하는 코드는 다음입니다.

return redis.xAdd(key, '*', {
  action: input.action,
  target: input.target,
  message: input.message,
  actorId: input.actorId !== undefined ? String(input.actorId) : '',
  createdAt: new Date().toISOString(),
});

 

예를 들어 관리자가 주문 상태를 변경했다면 다음과 같은 메시지를 남길 수 있습니다.

stream:audit-logs
  1710000000000-0
    action    order.status.changed
    target    order:1
    message   주문 상태가 PAID로 변경되었습니다.
    actorId   10
    createdAt 2026-06-16T00:00:00.000Z

 

2) worker는 Stream에서 감사 로그 작업을 읽고 DB에 저장합니다.

const auditLog = await prisma.auditLog.create({
  data: {
    action: job.action,
    target: job.target,
    message: job.message,
  },
});

 

3) DB 저장에 성공하면 ACK를 보냅니다.

await this.ackAuditLogJob(job.id);

실패하면 ACK하지 않습니다.
그러면 해당 메시지는 pending 상태로 남아 나중에 확인하거나 재처리할 수 있습니다.

 

 


※ 게시된 글 및 이미지 중 일부는 AI 도구의 도움을 받아 생성되거나 다듬어졌습니다.

반응형

 

반응형