deep-dive2025-02-23·10 min·177/348

Redis Streams 데이터 활용: 이벤트 소싱 및 실시간 처리

Redis Streams를 활용한 이벤트 소싱과 실시간 데이터 처리 아키텍처 구현 방법을 설명합니다.

Redis Streams 데이터 활용

Introduction

Redis Streams는 Redis 5.0부터 도입된 복잡한 메시징 데이터 구조로, Pub/Sub의 메시지 유실 문제를 해결하면서도 높은 처리 성능을 제공합니다. Append-only log 구조로 설계되어 이벤트 소싱 패턴과 실시간 데이터 처리에 적합합니다. 이 글에서는 Redis Streams의 활용 방법과 실제 구현 사례를 다루겠습니다.

Environment

# Redis Streams 테스트 환경
redis-cli

# Redis 버전 확인 (5.0 이상 필요)
redis-cli INFO server | grep redis_version
// Node.js Redis 클라이언트
const Redis = require('ioredis');
const redis = new Redis();

Problem

기존 메시징 시스템의 한계:

// Pub/Sub의 메시지 유실 문제
// 구독자가 오프라인이면 메시지 유실
await redis.publish('orders', JSON.stringify({ id: 1 }));
// 구독자가 연결되지 않으면 메시지는 사라짐

// List 기반 메시징의 한계
await redis.lpush('queue', JSON.stringify({ id: 1 }));
// 블로킹 polling 필요
const message = await redis.brpop('queue', 0);
// Consumer Group 기능 없음
# Redis Streams 기본 연산
XADD orders * orderId 1 amount 100.50
XLEN orders
XRANGE orders - +

Analysis

Redis Streams의 내부 구조를 분석했습니다:

// Redis Streams 구조 확인
await redis.xadd('test-stream', '*', 'key1', 'value1');
await redis.xadd('test-stream', '*', 'key2', 'value2');

// 스트림 정보 확인
const info = await redis.xinfo('STREAM', 'test-stream');
console.log('Stream Info:', info);

// 결과:
// 1) "length"
// 2) (integer) 2
// 3) "radix-tree-keys"
// 4) (integer) 1
// 5) "radix-tree-nodes"
// 6) (integer) 2
// 7) "groups"
// 8) (integer) 0
// 9) "last-generated-id"
// 10) "2-0"
# 스트림 데이터 구조
XRANGE orders - +
# 1) 1) "1-0"                    # Message ID (timestamp-sequence)
#    2) 1) "orderId"
#       2) "1"
#       3) "amount"
#       4) "100.50"
# 2) 1) "2-0"
#    2) 1) "orderId"
#       2) "2"
#       3) "amount"
#       4) "250.75"

Solution

1단계: 기본 Streams 구현

// 이벤트 발행
async function publishEvent(streamName, eventData) {
  const id = await redis.xadd(
    streamName,
    '*',
    'data', JSON.stringify(eventData),
    'timestamp', Date.now().toString(),
    'type', eventData.type
  );
  
  console.log(`Event published with ID: ${id}`);
  return id;
}

// 주문 이벤트 발행
await publishEvent('orders', {
  type: 'order_created',
  orderId: 12345,
  userId: 1001,
  amount: 99.99,
  items: ['item1', 'item2']
});

// 이벤트 조회
const events = await redis.xrange('orders', '-', '+');
console.log('Events:', events);

2단계: Consumer Group 구현

// Consumer Group 생성
async function createConsumerGroup(streamName, groupName) {
  try {
    await redis.xgroup('CREATE', streamName, groupName, '0', 'MKSTREAM');
    console.log(`Consumer group '${groupName}' created`);
  } catch (err) {
    if (err.message.includes('BUSYGROUP')) {
      console.log(`Consumer group '${groupName}' already exists`);
    } else {
      throw err;
    }
  }
}

// 컨슈머 그룹 생성
await createConsumerGroup('orders', 'order-processors');

// 메시지 처리
async function processMessages(consumerName) {
  while (true) {
    try {
      const results = await redis.xreadgroup(
        'GROUP', 'order-processors', consumerName,
        'COUNT', 10,
        'BLOCK', 5000,
        'STREAMS', 'orders', '>'
      );
      
      if (results) {
        for (const [stream, messages] of results) {
          for (const [id, fields] of messages) {
            await processMessage(id, fields);
            await redis.xack('orders', 'order-processors', id);
          }
        }
      }
    } catch (err) {
      console.error('Consumer error:', err);
      await new Promise(resolve => setTimeout(resolve, 1000));
    }
  }
}

async function processMessage(id, fields) {
  const message = {};
  for (let i = 0; i < fields.length; i += 2) {
    message[fields[i]] = fields[i + 1];
  }
  
  console.log(`Processing message ${id}:`, message);
  
  // 비즈니스 로직 처리
  const data = JSON.parse(message.data);
  if (data.type === 'order_created') {
    await handleOrderCreated(data);
  }
}

3단계: 스트림 관리 및 최적화

// 스트림 길이 제한 (MAXLEN)
await redis.xadd('orders', 'MAXLEN', '~', '10000', '*',
  'data', JSON.stringify({ type: 'test' })
);

// 스트림 트리밍 (오래된 메시지 삭제)
await redis.xtrim('orders', 'MAXLEN', 5000);

// 스트림 삭제
await redis.xdel('orders', '1-0', '2-0');

// Pending 메시지 확인
const pending = await redis.xpending('orders', 'order-processors');
console.log('Pending messages:', pending);

// 미처리 메시지 복구
const claimed = await redis.xclaim(
  'orders',
  'order-processors',
  'backup-consumer',
  3600000,  // 1시간 이상 미처리
  '0-0'  // 모든 ID
);
console.log('Claimed messages:', claimed);

4단계: 이벤트 소싱 구현

// 이벤트 소싱 아키텍처
class EventSourcedOrder {
  constructor(orderId) {
    this.orderId = orderId;
    this.events = [];
    this.state = {};
  }
  
  async loadEvents() {
    const events = await redis.xrange(
      `order:${this.orderId}`,
      '-',
      '+'
    );
    
    for (const [id, fields] of events) {
      const event = {};
      for (let i = 0; i < fields.length; i += 2) {
        event[fields[i]] = fields[i + 1];
      }
      this.events.push(event);
      this.applyEvent(event);
    }
  }
  
  applyEvent(event) {
    switch (event.type) {
      case 'order_created':
        this.state = {
          orderId: event.orderId,
          userId: event.userId,
          status: 'created',
          amount: parseFloat(event.amount)
        };
        break;
      case 'order_paid':
        this.state.status = 'paid';
        this.state.paidAt = event.timestamp;
        break;
      case 'order_shipped':
        this.state.status = 'shipped';
        this.state.shippedAt = event.timestamp;
        break;
    }
  }
  
  async addEvent(event) {
    const id = await redis.xadd(
      `order:${this.orderId}`,
      '*',
      'type', event.type,
      'data', JSON.stringify(event.data),
      'timestamp', Date.now().toString()
    );
    
    this.events.push({ ...event, id });
    this.applyEvent({ ...event, id });
    
    return id;
  }
  
  getState() {
    return { ...this.state };
  }
}

// 사용 예
const order = new EventSourcedOrder(12345);
await order.addEvent({
  type: 'order_created',
  data: { userId: 1001, amount: 99.99 }
});
await order.addEvent({
  type: 'order_paid',
  data: { paymentMethod: 'credit_card' }
});
console.log('Current state:', order.getState());

Lessons Learned

  1. Message ID 관리: Redis Streams의 Message ID는 타임스탬프 기반이므로 시간 동기화가 중요합니다
  2. Consumer Group 설계: 각 컨슈머는 고유한 이름을 가져야 하며, 실패 시 재처리 메커니즘이 필요합니다
  3. MAXLEN 설정: 스트림 크기를 제한하여 메모리 사용을 관리해야 합니다
  4. Pending 메시지 처리: 장시간 미처리된 메시지는 XCLAIM으로 다른 컨슈머에게 전달해야 합니다
  5. 이벤트 소싱 활용: Redis Streams는 이벤트 소싱 패턴의 구현에 최적화되어 있습니다

이 블로그는 외부 스폰서십, 제휴 마케팅 또는 광고 수익을 받지 않습니다.