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
- Message ID 관리: Redis Streams의 Message ID는 타임스탬프 기반이므로 시간 동기화가 중요합니다
- Consumer Group 설계: 각 컨슈머는 고유한 이름을 가져야 하며, 실패 시 재처리 메커니즘이 필요합니다
- MAXLEN 설정: 스트림 크기를 제한하여 메모리 사용을 관리해야 합니다
- Pending 메시지 처리: 장시간 미처리된 메시지는 XCLAIM으로 다른 컨슈머에게 전달해야 합니다
- 이벤트 소싱 활용: Redis Streams는 이벤트 소싱 패턴의 구현에 최적화되어 있습니다
이 블로그는 외부 스폰서십, 제휴 마케팅 또는 광고 수익을 받지 않습니다.