SSE 실시간 알림 (8/8)

이전 편: [스프링] 7. 알림 서비스와 프론트엔드 통합

이 문서가 시리즈의 마지막 편입니다.

지금까지 만든 알림 시스템에는 두 가지 약점이 있다. 연결이 살아있는지 확인할 방법이 없고, 연결이 끊긴 사이에 발생한 알림은 유실된다. 이번 편에서는 Heartbeat로 연결 상태를 감시하고, 이벤트 저장소를 만들어 유실된 이벤트를 복구하는 기능을 구현한다. 자세한 원리는 개념 4편에서 다루고, 여기서는 구현에 집중한다.

@EnableScheduling 활성화

Heartbeat와 이벤트 정리는 주기적으로 실행되어야 한다. 스프링의 @Scheduled를 쓰려면 먼저 스케줄링을 활성화해야 한다.

@SpringBootApplication
@EnableScheduling
public class SseNotificationApplication {
    public static void main(String[] args) {
        SpringApplication.run(SseNotificationApplication.class, args);
    }
}

@EnableScheduling을 메인 클래스에 추가한다. 이게 없으면 @Scheduled 어노테이션을 붙여도 아무 일도 일어나지 않는다.

Heartbeat 구현

왜 Heartbeat가 필요한가

SSE 연결은 HTTP 기반이다. 클라이언트와 서버 사이에 프록시, 로드밸런서, 방화벽이 있을 수 있다. 이들은 일정 시간 데이터가 흐르지 않으면 유휴 연결로 판단하고 끊어버린다. 서버 입장에서는 연결이 끊긴 줄 모르고 emitter를 계속 들고 있게 된다.

Heartbeat는 주기적으로 가벼운 데이터를 보내서 두 가지 문제를 해결한다:

  • 중간 장비가 연결을 끊지 않도록 트래픽을 유지한다.
  • 전송이 실패하면 죽은 연결을 감지하고 정리한다.
sequenceDiagram participant C as Client (Browser) participant P as Proxy/LB participant S as Server (Spring) Note over C,S: SSE 연결 수립 상태 loop 30초마다 (fixedRate) S->>S: cleanUp() 실행 S->>C: SSE Event: name="ping", data="ping" P-->>P: 유휴 타임아웃 초기화 C->>C: Console: "Ping 수신" 로그 출력 end Note right of S: 전송 실패 시 (IOException)
해당 Emitter 제거 및 정리

단계 1 lastActivityTime 추가

SseEmitterService에 각 사용자의 마지막 활동 시간을 기록하는 Map을 추가한다.

private final Map<String, LocalDateTime> lastActivityTime = new ConcurrentHashMap<>();

emitter와 마찬가지로 ConcurrentHashMap을 쓴다. @Scheduled 메서드와 일반 요청 처리 스레드가 동시에 접근할 수 있기 때문이다.

단계 2 ping 메서드

모든 연결에 가벼운 ping 이벤트를 보내는 메서드다.

public void ping() {
    emitters.forEach((userId, emitter) -> {
        try {
            emitter.send(SseEmitter.event().name("ping").data("ping"));
            lastActivityTime.put(userId, LocalDateTime.now());
        } catch (IOException e) {
            emitters.remove(userId);
            lastActivityTime.remove(userId);
            emitter.completeWithError(e);
        }
    });
}

ping을 보내는 진짜 목적은 실패를 유도하는 것이다. 이미 끊긴 연결에 send()를 시도하면 IOException이 발생한다. 이때 catch 블록에서 해당 emitter를 제거한다. 살아있는 연결이면 lastActivityTime을 갱신해서 "이 연결은 아직 살아있다"고 기록한다.

단계 3 cleanUp 스케줄러

30초마다 실행되는 정리 작업이다.

@Scheduled(fixedRate = 30000)
public void cleanUp() {
    LocalDateTime now = LocalDateTime.now();
    emitters.keySet().forEach(userId -> {
        LocalDateTime lastActivity = lastActivityTime.get(userId);
        if (lastActivity != null && lastActivity.plusMinutes(5).isBefore(now)) {
            SseEmitter emitter = emitters.get(userId);
            if (emitter != null) {
                emitter.complete();
                emitters.remove(userId);
                lastActivityTime.remove(userId);
            }
        }
    });
    ping();
}

두 가지 일을 한다:

  • 비활성 연결 정리 : 5분 이상 활동이 없는 연결을 강제로 닫는다. 클라이언트가 disconnect를 호출하지 않고 브라우저를 닫았을 때를 대비하는 안전장치다.
  • ping 전송 : 살아있는 모든 연결에 ping을 보낸다. 죽은 연결은 이 과정에서 IOException으로 걸러진다.

fixedRate = 30000이므로 30초마다 실행된다. 이 주기는 중간 장비의 유휴 타임아웃보다 짧아야 의미가 있다. 대부분의 프록시가 60초~120초 타임아웃을 갖고 있으므로 30초가 적당하다.

stateDiagram-v2 state "주기적 관리 작업" as Scheduled { [*] --> cleanUp : 30초 간격 [*] --> cleanupOldEvents : 10분 간격 } state cleanUp { [*] --> CheckInactivity : 비활성 체크 CheckInactivity --> RemoveEmitter : 5분 이상 무활동 RemoveEmitter --> SendPing : Emitter 제거 CheckInactivity --> SendPing : 활동 중 SendPing --> DetectDeath : ping 전송 DetectDeath --> RemoveEmitter : IOException 발생 } state cleanupOldEvents { [*] --> FilterOldEvents : 생성 시간 조회 FilterOldEvents --> DeleteFromStore : 30분 초과 이벤트 삭제 DeleteFromStore --> DeleteIndex : Index 데이터도 함께 제거 }

SseEvent 모델 작성

이벤트 유실을 복구하려면 보낸 이벤트를 어딘가에 저장해야 한다. 저장할 이벤트 객체를 먼저 정의한다.

package com.codeit.notification.model;

@Getter
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class SseEvent {
    private String eventId;
    private String userId;
    private String eventName;
    private Object data;
    private LocalDateTime createdAt;

    public static SseEvent create(String eventId, String userId, String eventName, Object data) {
        return SseEvent.builder()
                .eventId(eventId).userId(userId).eventName(eventName).data(data)
                .createdAt(LocalDateTime.now()).build();
    }
}

SseEvent는 DB 엔티티가 아니다. 메모리에 임시로 저장하는 객체다. 필드 구성은 다음과 같다:

  • eventId : 이벤트의 고유 식별자. 클라이언트가 "어디까지 받았는지"를 추적하는 데 쓴다.
  • userId : 수신 대상 사용자
  • eventName : notification, announcement 같은 이벤트 타입
  • data : 실제 전송 데이터
  • createdAt : 생성 시각. 오래된 이벤트 정리에 사용한다.

SseEventStore 구현

단계 1 클래스 뼈대와 설정값

@Service
@Slf4j
public class SseEventStore {
    private final AtomicLong eventIdGenerator = new AtomicLong(0);
    private final Map<String, List<SseEvent>> eventStore = new ConcurrentHashMap<>();
    private final Map<String, SseEvent> eventIndex = new ConcurrentHashMap<>();
    private static final int MAX_EVENTS_PER_USER = 100;
    private static final int EVENT_RETENTION_MINUTES = 30;
}

세 가지 필드가 있다:

  • eventIdGenerator : 이벤트 ID를 순차적으로 생성하는 카운터. AtomicLong이라서 멀티스레드 환경에서도 안전하다.
  • eventStore : 사용자별 이벤트 리스트. 유실 복구 시 특정 사용자의 이벤트만 빠르게 조회하기 위한 구조다.
  • eventIndex : eventId로 이벤트를 바로 찾을 수 있는 인덱스. 삭제 시 O(1) 접근에 쓴다.
classDiagram class SseEventStore { -Map eventStore (userId : List) -Map eventIndex (eventId : SseEvent) -AtomicLong idGenerator +saveEvent(userId, name, data) +getEventSince(userId, lastId) +cleanupOldEvents() } class SseEvent { +String eventId +String userId +String eventName +Object data +LocalDateTime createdAt } SseEventStore "1" --> "*" SseEvent : 사용자별 이벤트 목록 관리 note for SseEventStore "사용자당 최대 100개
30분 이상 경과 시 자동 삭제"

사용자당 최대 100개, 30분 이상 된 이벤트는 삭제한다. 메모리 기반 저장소이므로 무한정 쌓이면 안 된다.

단계 2 이벤트 저장

public String generateEventId() {
    return String.valueOf(eventIdGenerator.incrementAndGet());
}

public String saveEvent(String userId, String eventName, Object data) {
    String eventId = generateEventId();
    SseEvent event = SseEvent.create(eventId, userId, eventName, data);

    eventStore.computeIfAbsent(userId, k -> Collections.synchronizedList(new ArrayList<>()))
              .add(event);
    eventIndex.put(eventId, event);

    List<SseEvent> userEvents = eventStore.get(userId);
    if (userEvents.size() > MAX_EVENTS_PER_USER) {
        SseEvent removedEvent = userEvents.remove(0);
        eventIndex.remove(removedEvent.getEventId());
    }
    return eventId;
}

computeIfAbsent로 해당 사용자의 리스트가 없으면 새로 만든다. Collections.synchronizedList로 감싸는 이유는 여러 스레드에서 동시에 이벤트를 추가할 수 있기 때문이다.

사용자당 100개를 초과하면 가장 오래된 이벤트를 제거한다. FIFO 방식이다. 제거할 때 eventIndex에서도 함께 삭제해서 두 자료구조의 동기화를 유지한다.

단계 3 이벤트 조회

클라이언트가 재연결할 때 "마지막으로 받은 이벤트 ID" 이후의 이벤트를 조회하는 메서드다.

public List<SseEvent> getEventSince(String userId, String lastEventId) {
    List<SseEvent> userEvents = eventStore.get(userId);
    if (userEvents == null || userEvents.isEmpty()) return Collections.emptyList();
    if (lastEventId == null || lastEventId.isEmpty()) return new ArrayList<>(userEvents);

    return userEvents.stream()
            .filter(event -> Long.parseLong(event.getEventId()) > Long.parseLong(lastEventId))
            .collect(Collectors.toList());
}

eventId가 순차적인 숫자이므로 크기 비교가 가능하다. lastEventId보다 큰 ID를 가진 이벤트만 필터링하면 그게 클라이언트가 놓친 이벤트다.

lastEventId가 null이면 저장된 이벤트 전체를 반환한다. 최초 연결 시에는 놓친 이벤트가 없으므로 이 경우가 호출되지는 않지만, 방어적으로 처리한다.

단계 4 오래된 이벤트 정리

@Scheduled(fixedRate = 600000)
public void cleanupOldEvents() {
    LocalDateTime cutoffTime = LocalDateTime.now().minusMinutes(EVENT_RETENTION_MINUTES);
    for (Map.Entry<String, List<SseEvent>> entry : eventStore.entrySet()) {
        List<SseEvent> userEvents = entry.getValue();
        List<SseEvent> eventsToRemove = userEvents.stream()
                .filter(event -> event.getCreatedAt().isBefore(cutoffTime))
                .collect(Collectors.toList());
        for (SseEvent event : eventsToRemove) {
            userEvents.remove(event);
            eventIndex.remove(event.getEventId());
        }
        if (userEvents.isEmpty()) eventStore.remove(entry.getKey());
    }
}

10분마다(fixedRate = 600000) 30분 이상 된 이벤트를 정리한다. 메모리 누수를 방지하는 장치다.

사용자의 이벤트가 모두 삭제되면 해당 사용자의 키도 eventStore에서 제거한다. 빈 리스트를 들고 있을 이유가 없다.

SseEmitterService 수정

기존 SseEmitterService에 EventStore를 통합한다. 변경 포인트는 세 곳이다.

sequenceDiagram participant C as Client (Browser) participant S as SseEmitterService participant ES as SseEventStore C->>S: /connect?userId=user1&lastEventId=10 (재연결 요청) S->>ES: getEventSince("user1", "10") ES-->>S: List [Event 11, Event 12] S->>C: SSE "connect" 메시지 전송 Note over S,C: 놓친 이벤트 전송 시작 S->>C: SSE id="11", name="notification", data="..." S->>C: SSE id="12", name="notification", data="..." C->>C: LocalStorage: lastEventId = 12 갱신

단계 1 의존성 추가

@Service
@Slf4j
@RequiredArgsConstructor
public class SseEmitterService {
    @Value("${notification.timeout}")
    private Long timeout;

    private final SseEventStore sseEventStore;
    private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
    private final Map<String, LocalDateTime> lastActivityTime = new ConcurrentHashMap<>();
}

SseEventStore를 주입받고, lastActivityTime Map을 추가한다.

단계 2 createEmitter에 lastEventId 파라미터 추가

public SseEmitter createEmitter(String userId, String lastEventId) {
    SseEmitter emitter = new SseEmitter(timeout);
    if (emitters.containsKey(userId)) {
        emitters.get(userId).complete();
    }
    emitters.put(userId, emitter);

    emitter.onCompletion(() -> { emitters.remove(userId); });
    emitter.onTimeout(() -> { emitters.remove(userId); });
    emitter.onError((e) -> { emitters.remove(userId); });

    try {
        emitter.send(SseEmitter.event().name("connect").data("SSE 연결이 수립되었습니다."));
        restoreMissedEvents(userId, lastEventId, emitter);
    } catch (IOException e) {
        emitters.remove(userId);
        emitter.completeWithError(e);
    }
    return emitter;
}

기존 버전과의 차이는 lastEventId 파라미터가 추가된 것이다. 연결 수립 직후 restoreMissedEvents()를 호출해서 놓친 이벤트를 즉시 재전송한다.

단계 3 누락 이벤트 복원 메서드

private void restoreMissedEvents(String userId, String lastEventId, SseEmitter emitter) {
    if (lastEventId == null || lastEventId.isEmpty()) return;
    try {
        List<SseEvent> missedEvents = sseEventStore.getEventSince(userId, lastEventId);
        for (SseEvent event : missedEvents) {
            emitter.send(SseEmitter.event()
                    .id(event.getEventId())
                    .name(event.getEventName())
                    .data(event.getData()));
        }
    } catch (IOException e) {
        log.error("이벤트 복원 실패!", e);
    }
}

lastEventId가 없으면 최초 연결이므로 복원할 게 없다. 값이 있으면 EventStore에서 그 이후의 이벤트를 가져와서 하나씩 다시 보낸다. 이때 각 이벤트에 원래의 eventId를 붙여서 보내야 클라이언트가 중복 수신을 감지할 수 있다.

단계 4 sendToUser에 이벤트 ID 부여

public void sendToUser(String userId, String eventName, Object data) {
    String eventId = sseEventStore.saveEvent(userId, eventName, data);
    SseEmitter emitter = emitters.get(userId);
    if (emitter == null) { return; }
    try {
        emitter.send(SseEmitter.event().id(eventId).name(eventName).data(data));
    } catch (IOException e) {
        emitters.remove(userId);
        emitter.completeWithError(e);
    }
}

기존 버전에서 두 가지가 바뀌었다:

  • 전송 전에 sseEventStore.saveEvent()로 이벤트를 저장한다. 전송이 실패해도 저장소에는 남아있다.
  • SseEmitter.event().id(eventId)로 이벤트에 고유 ID를 부여한다. 클라이언트의 EventSource가 이 ID를 lastEventId로 자동 추적한다.

emitter가 null이면 해당 사용자가 연결되어 있지 않다는 뜻이다. 이때는 전송을 건너뛰지만, 이벤트는 이미 저장되어 있으므로 나중에 재연결하면 복구할 수 있다.

SseController 수정

@GetMapping(value = "/connect", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter connect(@RequestParam String userId,
                          @RequestParam(required = false) String lastEventId) {
    log.info("SSE Connecting to user {}", userId);
    return sseEmitterService.createEmitter(userId, lastEventId);
}

lastEventId@RequestParam(required = false)로 받는다. 최초 연결 시에는 이 값이 없고, 재연결 시에만 클라이언트가 전달한다.

프론트엔드 수정

단계 1 lastEventId 관리

let lastEventId = null;

function connect() {
    const userId = document.getElementById('userIdInput').value.trim();
    if (!userId) { alert('사용자 ID를 입력해주세요.'); return; }
    currentUserId = userId;
    if (eventSource) eventSource.close();

    let connectUrl = `/api/sse/connect?userId=${userId}`;
    lastEventId = localStorage.getItem(`lastEventId_${currentUserId}`);
    if (lastEventId) {
        connectUrl += `&lastEventId=${lastEventId}`;
    }
    eventSource = new EventSource(connectUrl);
}

localStorage에 사용자별로 마지막 이벤트 ID를 저장한다. 브라우저를 닫았다 다시 열어도 이 값이 남아있으므로, 재접속 시 놓친 이벤트를 복구할 수 있다.

localStorage인가? sessionStorage는 탭을 닫으면 사라진다. SSE 연결이 끊기는 가장 흔한 상황이 바로 브라우저/탭을 닫는 것이므로, 세션보다 오래 유지되는 localStorage가 적합하다.

단계 2 이벤트 수신 시 ID 갱신

eventSource.addEventListener('notification', (event) => {
    if (event.lastEventId) {
        lastEventId = event.lastEventId;
        localStorage.setItem(`lastEventId_${currentUserId}`, lastEventId);
    }
    const notification = JSON.parse(event.data);
    addNotificationToList(notification, true);
    updateUnreadCount();
    showNotificationToast(notification);
});

event.lastEventId는 서버가 SseEmitter.event().id(eventId)로 설정한 값이다. EventSource가 자동으로 이 값을 추적하지만, 브라우저를 닫으면 사라지므로 localStorage에도 별도로 저장한다.

단계 3 ping 리스너

eventSource.addEventListener('ping', (event) => {
    console.log('Ping 수신', new Date().toLocaleTimeString());
});

ping 이벤트는 처리할 로직이 없다. 로그만 찍어서 Heartbeat가 잘 동작하는지 확인하는 용도다. 프로덕션에서는 이 로그도 제거해도 무방하다.

동작 확인

Heartbeat 확인

  1. 서버를 실행하고 브라우저에서 user1로 SSE 연결한다.
  2. 30초를 기다린다.
  3. 브라우저 Console에 Ping 수신 로그가 찍히는지 확인한다.
  4. 서버 로그에서도 30초마다 cleanUp이 실행되는 것을 확인할 수 있다.

이벤트 유실 복구 확인

이 확인이 핵심이다. 단계별로 따라해 보자.

flowchart TD A[알림 발생] --> B[SseEventStore.saveEvent] B --> C[eventId 생성 및
Memory Store 저장] C --> D{사용자 연결 중?} D -- Yes --> E[SseEmitter.send
(Event ID 포함)] D -- No --> F[전송 대기] E --> G[클라이언트 수신] G --> H[LocalStorage에
lastEventId 업데이트] F --> I[재연결 요청 시
lastEventId 전달] I --> J[EventStore에서
이후 이벤트 복원] J --> E
  1. user1로 SSE 연결한다.
  2. 테스트 알림을 2~3개 전송한다. 알림이 정상 수신되는지 확인한다.
  3. 브라우저 개발자 도구 Console에서 localStorage를 확인한다. lastEventId_user1 값이 저장되어 있어야 한다.
  4. SSE 연결을 끊는다 (disconnect 버튼 또는 브라우저 탭 닫기).
  5. 서버 측에서 user1에게 알림을 2개 더 보낸다. 연결이 끊어진 상태이므로 실시간 수신은 안 된다.
  6. 다시 user1로 연결한다.
  7. 연결 직후 4단계에서 끊기기 전까지 받지 못한 알림 2개가 자동으로 도착하는지 확인한다.

5단계에서 보낸 알림이 SseEventStore에 저장되어 있다가, 6단계에서 재연결 시 lastEventId 이후의 이벤트로 복원된 것이다.

5단계에서 서버 측 알림 전송하는 방법

연결이 끊긴 사용자에게 알림을 보내려면 다른 브라우저 탭에서 REST API를 직접 호출하거나, Postman/curl로 POST /api/notifications를 요청한다. 사용자가 오프라인이어도 NotificationService.createAndSendNotification()은 DB 저장까지는 성공하고, SSE 전송만 건너뛴다.

이번 편 최종 전체 코드

SseEvent.java

package com.codeit.notification.model;

import lombok.*;
import java.time.LocalDateTime;

@Getter
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class SseEvent {
    private String eventId;
    private String userId;
    private String eventName;
    private Object data;
    private LocalDateTime createdAt;

    public static SseEvent create(String eventId, String userId, String eventName, Object data) {
        return SseEvent.builder()
                .eventId(eventId).userId(userId).eventName(eventName).data(data)
                .createdAt(LocalDateTime.now()).build();
    }
}

SseEventStore.java

package com.codeit.notification.service;

import com.codeit.notification.model.SseEvent;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;

@Service
@Slf4j
public class SseEventStore {
    private final AtomicLong eventIdGenerator = new AtomicLong(0);
    private final Map<String, List<SseEvent>> eventStore = new ConcurrentHashMap<>();
    private final Map<String, SseEvent> eventIndex = new ConcurrentHashMap<>();
    private static final int MAX_EVENTS_PER_USER = 100;
    private static final int EVENT_RETENTION_MINUTES = 30;

    public String generateEventId() {
        return String.valueOf(eventIdGenerator.incrementAndGet());
    }

    public String saveEvent(String userId, String eventName, Object data) {
        String eventId = generateEventId();
        SseEvent event = SseEvent.create(eventId, userId, eventName, data);
        eventStore.computeIfAbsent(userId, k -> Collections.synchronizedList(new ArrayList<>())).add(event);
        eventIndex.put(eventId, event);

        List<SseEvent> userEvents = eventStore.get(userId);
        if (userEvents.size() > MAX_EVENTS_PER_USER) {
            SseEvent removedEvent = userEvents.remove(0);
            eventIndex.remove(removedEvent.getEventId());
        }
        return eventId;
    }

    public List<SseEvent> getEventSince(String userId, String lastEventId) {
        List<SseEvent> userEvents = eventStore.get(userId);
        if (userEvents == null || userEvents.isEmpty()) return Collections.emptyList();
        if (lastEventId == null || lastEventId.isEmpty()) return new ArrayList<>(userEvents);
        return userEvents.stream()
                .filter(event -> Long.parseLong(event.getEventId()) > Long.parseLong(lastEventId))
                .collect(Collectors.toList());
    }

    @Scheduled(fixedRate = 600000)
    public void cleanupOldEvents() {
        LocalDateTime cutoffTime = LocalDateTime.now().minusMinutes(EVENT_RETENTION_MINUTES);
        for (Map.Entry<String, List<SseEvent>> entry : eventStore.entrySet()) {
            List<SseEvent> userEvents = entry.getValue();
            List<SseEvent> eventsToRemove = userEvents.stream()
                    .filter(event -> event.getCreatedAt().isBefore(cutoffTime))
                    .collect(Collectors.toList());
            for (SseEvent event : eventsToRemove) {
                userEvents.remove(event);
                eventIndex.remove(event.getEventId());
            }
            if (userEvents.isEmpty()) eventStore.remove(entry.getKey());
        }
    }
}

SseEmitterService.java

package com.codeit.notification.service;

import com.codeit.notification.model.SseEvent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

@Service
@Slf4j
@RequiredArgsConstructor
public class SseEmitterService {
    @Value("${notification.timeout}")
    private Long timeout;

    private final SseEventStore sseEventStore;
    private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
    private final Map<String, LocalDateTime> lastActivityTime = new ConcurrentHashMap<>();

    public SseEmitter createEmitter(String userId, String lastEventId) {
        SseEmitter emitter = new SseEmitter(timeout);
        if (emitters.containsKey(userId)) {
            emitters.get(userId).complete();
        }
        emitters.put(userId, emitter);

        emitter.onCompletion(() -> { emitters.remove(userId); });
        emitter.onTimeout(() -> { emitters.remove(userId); });
        emitter.onError((e) -> { emitters.remove(userId); });

        try {
            emitter.send(SseEmitter.event().name("connect").data("SSE 연결이 수립되었습니다."));
            restoreMissedEvents(userId, lastEventId, emitter);
        } catch (IOException e) {
            emitters.remove(userId);
            emitter.completeWithError(e);
        }
        return emitter;
    }

    private void restoreMissedEvents(String userId, String lastEventId, SseEmitter emitter) {
        if (lastEventId == null || lastEventId.isEmpty()) return;
        try {
            List<SseEvent> missedEvents = sseEventStore.getEventSince(userId, lastEventId);
            for (SseEvent event : missedEvents) {
                emitter.send(SseEmitter.event()
                        .id(event.getEventId())
                        .name(event.getEventName())
                        .data(event.getData()));
            }
        } catch (IOException e) {
            log.error("이벤트 복원 실패!", e);
        }
    }

    public void sendToUser(String userId, String eventName, Object data) {
        String eventId = sseEventStore.saveEvent(userId, eventName, data);
        SseEmitter emitter = emitters.get(userId);
        if (emitter == null) { return; }
        try {
            emitter.send(SseEmitter.event().id(eventId).name(eventName).data(data));
        } catch (IOException e) {
            emitters.remove(userId);
            emitter.completeWithError(e);
        }
    }

    public void broadcast(String eventName, Object data) {
        emitters.forEach((userId, emitter) -> {
            try {
                emitter.send(SseEmitter.event().name(eventName).data(data));
            } catch (IOException e) {
                emitters.remove(userId);
                emitter.completeWithError(e);
            }
        });
    }

    @Scheduled(fixedRate = 30000)
    public void cleanUp() {
        LocalDateTime now = LocalDateTime.now();
        emitters.keySet().forEach(userId -> {
            LocalDateTime lastActivity = lastActivityTime.get(userId);
            if (lastActivity != null && lastActivity.plusMinutes(5).isBefore(now)) {
                SseEmitter emitter = emitters.get(userId);
                if (emitter != null) {
                    emitter.complete();
                    emitters.remove(userId);
                    lastActivityTime.remove(userId);
                }
            }
        });
        ping();
    }

    public void ping() {
        emitters.forEach((userId, emitter) -> {
            try {
                emitter.send(SseEmitter.event().name("ping").data("ping"));
                lastActivityTime.put(userId, LocalDateTime.now());
            } catch (IOException e) {
                emitters.remove(userId);
                lastActivityTime.remove(userId);
                emitter.completeWithError(e);
            }
        });
    }

    public void closeEmitter(String userId) {
        SseEmitter emitter = emitters.get(userId);
        if (emitter != null) { emitter.complete(); emitters.remove(userId); }
    }

    public int getConnectedUserCount() { return emitters.size(); }

    public boolean isConnected(String userId) { return emitters.containsKey(userId); }
}

SseController.java 변경 부분

@GetMapping(value = "/connect", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter connect(@RequestParam String userId,
                          @RequestParam(required = false) String lastEventId) {
    log.info("SSE Connecting to user {}", userId);
    return sseEmitterService.createEmitter(userId, lastEventId);
}

프론트엔드 변경 부분

let lastEventId = null;

function connect() {
    const userId = document.getElementById('userIdInput').value.trim();
    if (!userId) { alert('사용자 ID를 입력해주세요.'); return; }
    currentUserId = userId;
    if (eventSource) eventSource.close();

    let connectUrl = `/api/sse/connect?userId=${userId}`;
    lastEventId = localStorage.getItem(`lastEventId_${currentUserId}`);
    if (lastEventId) {
        connectUrl += `&lastEventId=${lastEventId}`;
    }
    eventSource = new EventSource(connectUrl);

    eventSource.addEventListener('connect', (event) => {
        console.log('SSE 연결 성공:', event.data);
        updateConnectionStatus(true);
    });

    eventSource.addEventListener('notification', (event) => {
        if (event.lastEventId) {
            lastEventId = event.lastEventId;
            localStorage.setItem(`lastEventId_${currentUserId}`, lastEventId);
        }
        const notification = JSON.parse(event.data);
        addNotificationToList(notification, true);
        updateUnreadCount();
        showNotificationToast(notification);
    });

    eventSource.addEventListener('announcement', (event) => {
        const announcement = JSON.parse(event.data);
        addNotificationToList(announcement, true);
        showNotificationToast(announcement);
    });

    eventSource.addEventListener('ping', (event) => {
        console.log('Ping 수신', new Date().toLocaleTimeString());
    });

    eventSource.onerror = (error) => {
        console.error('SSE 에러:', error);
        updateConnectionStatus(false);
    };
}

자주 하는 실수

@EnableScheduling을 빼먹음

@Scheduled를 붙여도 @EnableScheduling이 없으면 스케줄러가 동작하지 않는다. Heartbeat도 안 보내지고 이벤트 정리도 안 된다. 애플리케이션이 에러 없이 정상 실행되기 때문에 문제를 늦게 발견하는 경우가 많다.

[!DANGER] lastEventId를 sessionStorage에 저장

sessionStorage는 탭을 닫으면 사라진다. SSE 재연결이 필요한 가장 흔한 상황이 브라우저를 닫았다 다시 여는 것이므로, 반드시 localStorage를 써야 한다. sessionStorage를 쓰면 탭을 닫는 순간 lastEventId가 사라져서 이벤트 복구가 불가능해진다.

[!DANGER] 이벤트 저장 없이 ID만 부여

SseEmitter.event().id(eventId)로 ID를 붙여 보내기만 하고, 실제로 이벤트를 저장하지 않는 실수다. ID가 있어도 저장소에 이벤트가 없으면 재연결 시 복원할 데이터가 없다. sendToUser에서 반드시 sseEventStore.saveEvent()를 먼저 호출해야 한다.