Skip to content

redis PubSub

qbsb147 edited this page Mar 21, 2026 · 7 revisions

⚡ Chat 메시지 읽음 처리 이벤트

📌 목적

메시지 조회 시 읽음 처리 작업을 비동기로 분리하여 빠른 응답 시간(Latency) 제공

🔹 메시지 조회 시 이벤트 발행

@Override
public Page<MessageDto.Response> getMessages(Long roomNo, Pageable pageable) {
    ChatRoom chatRoom = chatRoomRepository.findById(roomNo)
            .orElseThrow(() -> new NotFoundException("해당 채팅방을 찾을 수가 없습니다."));
    Page<MessageProjection> chatMessages = chatMessageRepository.findMessagesWithUnreadAndMember(roomNo, pageable);

    UUID publicUuid = jwtTokenProvider.getPublicUuidFromToken();

     // redis를 통한 서비스 구현
👉  chatEventProducer.sendReadEvent(roomNo, publicUuid);

    return MessageDto.Response(chatMessage);
}
  • chatEventProducer.sendReadEvent() 메서드를 호출

🏗️ Pub/Sub 패턴 구현

Pub/Sub = 발행자(Publisher)가 메시지를 발행하면, 구독자(Subscriber)가 이를 받아 처리하는 구조

@Service
@RequiredArgsConstructor
public class ChatEventProducer {

    private final RedisTemplate<String, String> redisTemplate;

    public void sendReadEvent(Long roomNo, String publicUuid) {
        // Redis Stream에 저장할 데이터
        Map<String, String> data = Map.of(
            "roomNo", roomNo.toString(),
            "publicUuid", publicUuid
        );

        // chat_read_stream에 메시지 추가
        redisTemplate.opsForStream().add("chat_read_stream", data);
    }
}
  • Redis Stream은 String 객체만 저장 가능
  • 따라서 이벤트 데이터를 Map<String, String> 형태로 변환 후 전달
  • Stream은 큐처럼 순서대로 메시지를 저장
    • 메시지가 XADD 명령으로 Stream에 추가되면, 내부적으로 항목(Item) 단위로 순서대로 저장
    • 항목(Item) 단위로 ID + Field-Value 쌍으로 저장
    • XREAD 또는 XREADGROUP으로 순서대로 읽을 수 있음
  • 일반 큐와 달리 소비 후에도 기록이 남아 재처리 가능
  • Consumer Group을 사용하면 여러 소비자가 동시에 읽으면서 안정적 재처리 가능
  • 따라서 이벤트 기반 처리, 재시도, 지연 처리 등에 적합

Redis로 1초마다 pull하는 메서드 구현

@Service
@RequiredArgsConstructor
@Slf4j
public class ChatReadEventConsumer {
    private final RedisTemplate<String, String> redisTemplate;
    private final MessageReadStatusRepository messageReadStatusRepository;

    private static final String STREAM_KEY = "chat_read_stream";
    private static final String GROUP = "group1";
    private static final String CONSUMER = "consumer1";
    private static final int BATCH_SIZE = 10;

    @Scheduled(fixedDelay = 1000)
    public void consume() {
        // Pending 메시지 먼저 처리 (재처리)
        PendingMessagesSummary pending = redisTemplate.opsForStream().pending(STREAM_KEY, GROUP);

        // Pending 메시지별 처리
        pending.getPendingMessagesPerConsumer().forEach((consumer, msgId) -> {
            // msgId는 Long 하나
            List<MapRecord<String, Object, Object>> records =
                    redisTemplate.opsForStream()
                            .range(STREAM_KEY, Range.closed(msgId.toString(), msgId.toString()));

                // records 처리
                processMessages(records);
        });
        // 새로 들어온 메시지 읽기
        List<MapRecord<String, Object, Object>> newMessages =
                redisTemplate.opsForStream().read(
                        Consumer.from(GROUP, CONSUMER),
                        StreamReadOptions.empty().count(BATCH_SIZE),
                        StreamOffset.fromStart(STREAM_KEY));

        processMessages(newMessages);
    }

    private void processMessages(List<MapRecord<String, Object, Object>> messages) {
        for (MapRecord<String, Object, Object> msg : messages) {
            try {
                Long roomNo = Long.valueOf((String) msg.getValue().get("roomNo"));
                UUID publicUuid = UUID.fromString((String) msg.getValue().get("publicUuid"));

                // 읽음 처리
                messageReadStatusRepository.markMessageRead(roomNo, publicUuid);

                // ACK: 처리 완료 표시
                redisTemplate.opsForStream().acknowledge(STREAM_KEY, GROUP, msg.getId());

            } catch (Exception e) {
                // 실패 로그 기록, 다음 스케줄에서 Pending으로 재처리
                log.error("읽음 처리 실패, 재시도 예정: msgId={}, error={}", msg.getId(), e.getMessage());
            }
        }
    }
}

🔹 설명

  • Consumer Group 사용

    • group1 : 그룹 이름
    • consumer1 : 해당 그룹 내 소비자 이름
    • 여러 소비자가 동시에 Stream을 읽고 처리 가능
  • MapRecord<String, Object, Object>

타입 의미
String Stream의 key, 즉 Stream 이름(chat_read_stream)
Object (첫 번째) 항목의 Field 이름 (roomNo, publicUuid 등)
Object (두 번째) 항목의 Value ("123" 같은 String 값, Redis Stream은 String만 저장 가능)
  • 메시지 읽기
    • StreamReadOptions.empty().count(10) : 최대 10개씩 읽기
    • StreamOffset.fromStart("chat_read_stream") : Stream 처음부터 읽기
  • Pending 메시지 처리
    • 처리 실패 시 ACK를 하지 않은 메시지는 Pending 상태로 남음
    • 다음 스케줄에서 Pending 메시지를 먼저 읽어 재처리 가능
    • 이를 통해 안정적인 재처리와 중복 처리 방지 가능
  • 데이터 변환
    • Redis Stream은 String 객체만 저장 → Map에서 읽어 Long, UUID로 변환
  • 읽음 처리
    • messageReadStatusRepository.markMessageRead() 호출
  • ACK
    • redisTemplate.opsForStream().acknowledge()로 처리 완료 표시 → 재시도 방지
  • 재시도
    • 예외 발생 시 로그 기록, 다음 스케줄에서 다시 처리 가능

⚙️ 특징

  • 주기적으로 Stream을 읽는 Pull 방식 소비자
  • Pending 메시지 먼저 처리하여 안정적인 재처리 보장
  • Redis Stream + Consumer Group으로 안정적 재처리 가능
  • 메시지 처리 실패 시 다시 읽어서 재처리 가능, 중복 처리 방지 위해 ACK 필요

📡 이벤트 스트리밍

⚙️ 개발 환경 구축

Websocket

Kafka

Redis

Debezium

🧩 기능 구현

이벤트리스너(Spring)

Pub/Sub 기반 메세지 처리(Redis)

이벤트 스트리밍 처리(Kafka)

CDC(Debezium)

🧠 개념

🎯 설계 패턴

🏗️ 아키텍처

📚 기술 스택

Clone this wiki locally