forked from MinePing/Final
-
Notifications
You must be signed in to change notification settings - Fork 0
kafka eventStreaming
qbsb147 edited this page Mar 21, 2026
·
4 revisions
웹소켓(WebSocket)으로 들어오는 요청을 비동기로 처리하여 블로킹을 최소화하고, 동시에 많은 트래픽을 안정적으로 처리할 수 있도록 설계합니다.
Kafka를 통해 Producer → Topic → Consumer 구조로 이벤트를 전달하고, Consumer에서 서비스 로직을 처리하는 방식입니다.
sequenceDiagram
participant WS as WebSocket Client
participant P as Kafka Producer
participant K as Kafka Topic
participant C as Kafka Consumer
participant S as ChatService
WS->>P: 접속/메세지/읽음 이벤트 전달
P->>K: 토픽별 메시지 전송
K->>C: 메시지 소비
C->>S: 서비스 로직 실행 (메세지 전송, 상태 업데이트)
//클라이언트가 WebSocket 연결을 시도하면 호출되는 메서드
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
// 접속 시 token으로부터 객체를 추출하고 이를 전달
chatKafkaProducer.sendUserStatus(userStateEvent);
}
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
// 클라이언트로부터 넘어온 메세지를 Kafka로 전달
chatKafkaProducer.sendChatMessage(messageDto);
}
//WebSocket 클라이언트가 연결을 끊었을 때 자동으로 호출
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
//연결 종료 시 정보를 전달
chatKafkaProducer.sendUserStatus(userStateEvent);
}- WebSocket 이벤트를 Kafka Topic에 전달
- Topic 종류:
- chat_message → 채팅 메시지
- user_status → 유저 상태 변화
- chat_read → 메시지 읽음 이벤트
@Component
@RequiredArgsConstructor
public class ChatKafkaProducerImpl implements ChatKafkaProducer {
private final KafkaTemplate<String, Object> kafkaTemplate;
public void sendChatMessage(MessageDto.Request messageDto){
kafkaTemplate.send("chat_message", messageDto);
}
public void sendUserStatus(ChatEvent.UserStateEvent userStateEvent) {
kafkaTemplate.send("user_status", userStateEvent);
}
public void sendChatReadEvent(MessageReadStatusDto.Request chatReadEvent) {
kafkaTemplate.send("chat_read", chatReadEvent);
}
}✅ spring.kafka.producer.properties.enable.idempotence: true
- 멱등성(Idempotence) 활성화
- Consumer에서 예외 발생 시 재처리 가능 → Exactly-Once 보장
- Topic별 메시지를 수신하고 Service 로직 실행
- 각 KafkaListener는 독립적으로 이벤트 처리
@Slf4j
@Component
@RequiredArgsConstructor
public class ChatKafkaConsumerImpl implements ChatKafkaConsumer{
private final ChatService chatService;
@KafkaListener(topics = "chat_message", groupId = "chat_service")
public void consumeChatMessage(MessageDto.Request messageDto) throws IOException {
chatService.consumeMessage(messageDto);
}
@KafkaListener(topics = "user_status", groupId = "chat_service")
public void consumeUserStatus(ChatEvent.UserStateEvent userStateEvent) throws IOException {
chatService.sessionStateChange(userStateEvent);
}
@KafkaListener(topics = "chat_read", groupId = "chat_service")
public void consumeChatReadEvent(MessageReadStatusDto.Request chatReadEvent) {
try {
chatService.readMessage(chatReadEvent);
} catch (Exception e){
throw e;
}
}
@KafkaListener(topics = "mineping_server.mineping.message_read_status",
groupId = "chat_service",
containerFactory = "kafkaStringListenerFactory")
public void consumeChatReadDebeziumEvent(String payload) throws IOException {
chatService.chatReadDebezium(payload);
}
}Consumer에서 메시지를 수신하면, 구독하고 있는 클라이언트(WebSocket Session)에게 전송합니다.
@Override
private void send(MessageDto.Request chatMessageDto) throws IOException {
for (String subscriber : subscribers) {
try {
UUID subscriberUuid = UUID.fromString(subscriber);
WebSocketSession subscriberSession = onlineWebSocketHandler.getSession(subscriberUuid);
}
}
}
public void handle(MessageDto.Request chatMessageDto) throws IOException {
//redis에서 나를 구독하는 사람들을 조회
//웹소켓에서 session들을 가져와 해당 메세지를 전송
send(chatMessageDto);
}
@Override
public void sessionStateChange(ChatEvent.UserStateEvent userStateEvent) throws IOException {
//redis에서 나를 구독하는 사람들을 조회
//웹소켓에서 session들을 가져와 해당 메세지를 전송
userStateHandler.broadCast(publicUuid);
}- 블로킹 최소화: WebSocket 이벤트를 바로 처리하지 않고 Kafka를 통해 비동기 전달
- Exactly-Once 처리: Producer에서 멱등성을 활성화하여 Consumer 재처리 가능
- 구독자 기반 전송: Redis를 통해 구독자 목록 조회 후 WebSocket으로 메시지 전송
- 확장성: Kafka Topic 단위로 확장 가능, 이벤트 타입 추가 용이