1편에서 알림 도메인을 만들었다. 이번 편에서는 SSE의 핵심인 SseEmitter를 관리하는 서비스를 구현한다. 사용자별 연결을 저장하고, 특정 사용자에게 이벤트를 보내고, 전체에게 브로드캐스트하는 기능까지 만든다.
왜 Emitter 관리 서비스가 필요한가
개념 1편에서 봤듯이, SseEmitter는 컨트롤러에서 만들어서 리턴하면 끝이 아니다. "user1에게 알림을 보내야 하는데, user1의 emitter가 어디 있지?" 하는 문제가 생긴다.
컨트롤러의 connect() 메서드에서 emitter를 만들고 리턴하면, 그 메서드 실행이 끝난 뒤에는 emitter에 대한 참조가 사라진다. 나중에 알림을 보내려고 해도 해당 사용자의 emitter를 찾을 방법이 없다.
그래서 emitter를 별도 저장소에 보관하고, 사용자 ID로 조회할 수 있게 만드는 서비스가 필요하다. 이 서비스가 연결의 생성, 조회, 전송, 종료를 모두 담당한다.
클래스 구조와 저장소 선언
com.codeit.notification.service 패키지에 SseEmitterService를 만든다. 먼저 클래스의 뼈대부터 잡자.
package com.codeit.notification.service;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Service
@Slf4j
public class SseEmitterService {
@Value("${notification.timeout}")
private Long timeout;
private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
}
두 가지 핵심 필드가 있다.
@Value("${notification.timeout}")는 1편에서 application.yml에 넣어둔 notification.timeout 값을 주입받는다. emitter 생성 시 타임아웃으로 사용할 값이다. 코드에 숫자를 하드코딩하지 않고 설정 파일에서 가져오니, 운영 환경에서 타임아웃을 바꿀 때 코드를 수정할 필요가 없다.
ConcurrentHashMap이 왜 일반 HashMap이 아닌지가 중요하다. 웹 서버는 여러 요청을 동시에 처리한다. user1이 연결하는 요청과 user2가 연결하는 요청이 동시에 들어올 수 있다. 일반 HashMap은 동시 접근에 안전하지 않다. 두 스레드가 동시에 put을 하면 데이터가 꼬일 수 있다. ConcurrentHashMap은 내부적으로 동시성을 보장해서 이런 문제가 없다. 자세한 내용은 개념 2편에서 다룬다.
createEmitter 구현
사용자가 SSE 연결을 요청하면 호출되는 메서드다. 단계별로 나눠서 살펴보자.
단계 1 기존 연결 처리
public SseEmitter createEmitter(String userId) {
SseEmitter emitter = new SseEmitter(timeout);
if (emitters.containsKey(userId)) {
SseEmitter oldEmitter = emitters.get(userId);
oldEmitter.complete();
log.info("기존 SSE 연결 종료 - 사용자: {}", userId);
}
emitters.put(userId, emitter);
새 emitter를 만들기 전에, 같은 사용자의 기존 연결이 있는지 먼저 확인한다. 있으면 complete()로 깔끔하게 종료한 뒤 새 emitter로 교체한다.
왜 이렇게 하느냐면, 사용자가 브라우저 탭을 닫고 다시 열거나 새로고침하면 새 연결 요청이 들어온다. 이때 기존 emitter가 아직 Map에 남아있으면 좀비 연결이 된다. 실제로는 클라이언트 쪽이 끊겼는데 서버는 모르고 있는 상태다. 이 좀비를 정리하지 않으면 메시지 전송 시 IOException이 발생한다.
단계 2 생명주기 콜백 등록
emitter.onCompletion(() -> {
emitters.remove(userId);
log.info("SSE 연결 완료 - 사용자: {}, 남은 연결 수: {}", userId, emitters.size());
});
emitter.onTimeout(() -> {
emitters.remove(userId);
log.info("SSE 연결 타임아웃 - 사용자: {}", userId);
});
emitter.onError((e) -> {
emitters.remove(userId);
log.info("SSE 연결 에러 - 사용자: {}, 에러: {}", userId, e.getMessage());
});
세 가지 콜백 모두 같은 일을 한다. Map에서 해당 emitter를 제거하는 것이다.
onCompletion:complete()호출이나 정상 종료 시 실행된다.onTimeout: 설정한 타임아웃 시간이 지나면 실행된다.onError: 전송 중 에러가 발생하면 실행된다. 클라이언트가 연결을 끊었을 때 주로 발생한다.
이 콜백들이 없으면 어떻게 될까? 연결이 끊긴 emitter가 Map에 계속 남아있게 된다. 나중에 이 emitter에 메시지를 보내면 IOException이 터진다. 콜백으로 즉시 정리해야 Map이 항상 활성 연결만 보유하게 된다.
단계 3 초기 메시지 전송
try {
emitter.send(SseEmitter.event()
.name("connect")
.data("SSE 연결이 수립되었습니다."));
} catch (IOException e) {
log.error("초기 메시지 전송 실패 - 사용자: {}", userId, e);
emitters.remove(userId);
emitter.completeWithError(e);
}
return emitter;
}
연결이 수립되자마자 connect라는 이름의 이벤트를 보낸다. 이 초기 메시지가 왜 필요할까?
SSE 연결은 HTTP 응답 스트림을 열어두는 방식인데, 아무 데이터도 보내지 않으면 일부 브라우저나 프록시가 연결이 실패한 것으로 판단할 수 있다. 특히 Nginx 같은 리버스 프록시는 일정 시간 내에 응답이 없으면 타임아웃으로 처리한다. 초기 메시지를 보내면 "이 연결은 살아있다"는 신호가 된다.
전송에 실패하면 completeWithError()로 emitter를 에러 상태로 종료하고, Map에서도 제거한다.
sendToUser 구현
특정 사용자에게 이벤트를 보내는 메서드다.
public void sendToUser(String userId, String eventName, Object data) {
SseEmitter emitter = emitters.get(userId);
if (emitter == null) {
log.warn("SSE 연결을 찾을 수 없음 - 사용자: {}", userId);
return;
}
try {
emitter.send(SseEmitter.event()
.name(eventName)
.data(data));
log.info("알림 전송 성공 - 사용자: {}, 이벤트: {}", userId, eventName);
} catch (IOException e) {
log.error("알림 전송 실패 - 사용자: {}", userId, e);
emitters.remove(userId);
emitter.completeWithError(e);
}
}
흐름은 단순하다.
- Map에서 userId로 emitter를 꺼낸다.
- 없으면 로그만 남기고 리턴한다. 사용자가 SSE에 연결하지 않은 상태일 수 있다. 이건 에러가 아니다.
- 있으면 이벤트를 전송한다.
- 전송 실패 시 emitter를 정리한다.
data 파라미터가 Object 타입인 점에 주목하자. SseEmitter.event().data()는 객체를 받으면 Jackson이 자동으로 JSON 직렬화한다. DTO를 넘기면 JSON 문자열로 변환되어 클라이언트에 전달된다. 직접 ObjectMapper를 쓸 필요가 없다.
IOException이 발생하는 대표적인 경우는 클라이언트가 이미 연결을 끊은 상태에서 메시지를 보내려 할 때다. 브라우저 탭을 닫거나, 네트워크가 끊기거나, eventSource.close()를 호출한 경우다. 이때 Map에서 제거하지 않으면 같은 에러가 반복된다.
broadcast 구현
모든 연결된 사용자에게 같은 메시지를 보내는 메서드다. 시스템 공지나 서비스 점검 안내 같은 전체 알림에 사용한다.
public void broadcast(String eventName, Object data) {
log.info("브로드캐스트 시작 - 이벤트: {}, 대상: {}명", eventName, emitters.size());
emitters.forEach((userId, emitter) -> {
try {
emitter.send(SseEmitter.event()
.name(eventName)
.data(data));
} catch (IOException e) {
log.error("알림 전송 실패 - 사용자: {}", userId, e);
emitters.remove(userId);
emitter.completeWithError(e);
}
});
}
forEach 안에서 각 emitter마다 개별적으로 try-catch를 감싸는 게 핵심이다. 만약 try-catch를 바깥에 두면, user1에게 전송이 실패했을 때 user2, user3에게는 아예 보내지도 못하고 메서드가 끝난다. 한 명의 전송 실패가 나머지 전체를 막아서는 안 된다.
ConcurrentHashMap의 forEach는 순회 중에 remove를 호출해도 ConcurrentModificationException이 발생하지 않는다. 일반 HashMap이었다면 순회 중 삭제가 불가능하다. 이것도 ConcurrentHashMap을 쓰는 이유 중 하나다.
유틸리티 메서드
연결을 종료하거나, 현재 연결 상태를 확인하는 보조 메서드들이다.
public void closeEmitter(String userId) {
SseEmitter emitter = emitters.get(userId);
if (emitter != null) {
emitter.complete();
emitters.remove(userId);
}
}
public void closeAllEmitters() {
emitters.forEach((userId, emitter) -> emitter.complete());
emitters.clear();
}
public int getConnectedUserCount() {
return emitters.size();
}
public boolean isConnected(String userId) {
return emitters.containsKey(userId);
}
closeEmitter: 특정 사용자의 연결을 서버 쪽에서 종료한다. 사용자가 로그아웃했을 때 호출할 수 있다.closeAllEmitters: 전체 연결을 종료한다. 서버를 내리기 전에 호출하면 클라이언트가 깔끔하게 재연결을 시도할 수 있다.getConnectedUserCount: 현재 접속자 수를 리턴한다. 모니터링이나 관리자 페이지에서 활용한다.isConnected: 특정 사용자가 SSE에 연결 중인지 확인한다. 알림을 보내기 전에 연결 여부를 미리 체크할 때 쓴다.
closeEmitter에서 complete() 호출 후 remove도 하는데, 사실 complete()가 호출되면 onCompletion 콜백에서도 remove가 실행된다. 그래도 명시적으로 remove를 호출하는 이유는 콜백 실행 타이밍이 즉시가 아닐 수 있기 때문이다. 방어적으로 양쪽 모두에서 제거하는 게 안전하다.
이번 편 최종 전체 코드
SseEmitterService.java
package com.codeit.notification.service;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Service
@Slf4j
public class SseEmitterService {
@Value("${notification.timeout}")
private Long timeout;
private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
public SseEmitter createEmitter(String userId) {
SseEmitter emitter = new SseEmitter(timeout);
if (emitters.containsKey(userId)) {
SseEmitter oldEmitter = emitters.get(userId);
oldEmitter.complete();
log.info("기존 SSE 연결 종료 - 사용자: {}", userId);
}
emitters.put(userId, emitter);
emitter.onCompletion(() -> {
emitters.remove(userId);
log.info("SSE 연결 완료 - 사용자: {}, 남은 연결 수: {}", userId, emitters.size());
});
emitter.onTimeout(() -> {
emitters.remove(userId);
log.info("SSE 연결 타임아웃 - 사용자: {}", userId);
});
emitter.onError((e) -> {
emitters.remove(userId);
log.info("SSE 연결 에러 - 사용자: {}, 에러: {}", userId, e.getMessage());
});
try {
emitter.send(SseEmitter.event()
.name("connect")
.data("SSE 연결이 수립되었습니다."));
} catch (IOException e) {
log.error("초기 메시지 전송 실패 - 사용자: {}", userId, e);
emitters.remove(userId);
emitter.completeWithError(e);
}
return emitter;
}
public void sendToUser(String userId, String eventName, Object data) {
SseEmitter emitter = emitters.get(userId);
if (emitter == null) {
log.warn("SSE 연결을 찾을 수 없음 - 사용자: {}", userId);
return;
}
try {
emitter.send(SseEmitter.event()
.name(eventName)
.data(data));
log.info("알림 전송 성공 - 사용자: {}, 이벤트: {}", userId, eventName);
} catch (IOException e) {
log.error("알림 전송 실패 - 사용자: {}", userId, e);
emitters.remove(userId);
emitter.completeWithError(e);
}
}
public void broadcast(String eventName, Object data) {
log.info("브로드캐스트 시작 - 이벤트: {}, 대상: {}명", eventName, emitters.size());
emitters.forEach((userId, emitter) -> {
try {
emitter.send(SseEmitter.event()
.name(eventName)
.data(data));
} catch (IOException e) {
log.error("알림 전송 실패 - 사용자: {}", userId, e);
emitters.remove(userId);
emitter.completeWithError(e);
}
});
}
public void closeEmitter(String userId) {
SseEmitter emitter = emitters.get(userId);
if (emitter != null) {
emitter.complete();
emitters.remove(userId);
}
}
public void closeAllEmitters() {
emitters.forEach((userId, emitter) -> emitter.complete());
emitters.clear();
}
public int getConnectedUserCount() {
return emitters.size();
}
public boolean isConnected(String userId) {
return emitters.containsKey(userId);
}
}
자주 하는 실수
HashMap은 동시 접근에 안전하지 않다. 웹 서버는 여러 요청을 동시에 처리하기 때문에, 두 사용자가 동시에 연결하면 데이터가 꼬일 수 있다. 심하면 무한 루프에 빠지기도 한다. 반드시 ConcurrentHashMap을 사용한다. 동시성 이슈에 대한 자세한 설명은 개념 2편을 참고한다.
[!DANGER] 초기 메시지를 보내지 않음
createEmitter에서 emitter를 만들고 바로 리턴하면, 일부 브라우저나 리버스 프록시가 응답이 없다고 판단해서 연결을 끊을 수 있다. 연결 직후 더미 이벤트라도 하나 보내야 "이 연결은 살아있다"는 신호가 된다. connect 이벤트를 보내는 건 관례이기도 하다.
[!DANGER] broadcast에서 try-catch를 forEach 바깥에 둠
forEach 전체를 하나의 try-catch로 감싸면, 첫 번째 전송 실패에서 예외가 던져지고 나머지 사용자에게는 전송 자체가 시도되지 않는다. 각 emitter마다 개별 try-catch로 감싸야 한 명의 실패가 전체에 영향을 주지 않는다.