forked from MinePing/Final
-
Notifications
You must be signed in to change notification settings - Fork 0
debezium cdc
qbsb147 edited this page Mar 21, 2026
·
2 revisions
메시지가 읽혔음을 감지하면 해당 채팅방 구독자들에게 실시간으로 상태를 전파합니다.
즉, 읽음 처리 이벤트 → Kafka → Consumer → WebSocket 전파까지 연결하는 흐름입니다.
즉, 읽음 처리 이벤트 → Kafka → Consumer → WebSocket 전파까지 연결하는 흐름입니다.
Debezium MySQL Source Connector에서 설정한 값을 기준으로 토픽을 구독합니다.
-
connect-mysql-source.properties예시
database.server.name=example_server
database.include.list=example_db
table.include.list=example_table- Consumer 측
@KafkaListener적용
@KafkaListener(
topics = "example_server.example_db.example_table",
groupId = "chat_service",
containerFactory = "kafkaStringListenerFactory"
)
public void consumeChatReadDebeziumEvent(String payload) throws IOException {
chatService.chatReadDebezium(payload);
}💡 설명:
Debezium이 MySQL 테이블 변화를 감지하면, 토픽에 메시지를 발행하고 해당 메서드가 이를 소비합니다.
public void chatReadDebezium(String payload) throws IOException {
JsonNode root = objectMapper.readTree(payload);
String op = root.get("op").asText(); // u=update, c=create, d=delete
if (!"u".equals(op)) return; // update 이벤트만 처리
JsonNode after = root.path("after");
if (after.isMissingNode() || after.isNull()) {
log.warn("after is null. payload={}", payload);
return;
}
// 📌 주요 데이터 추출
Long roomNo = after.path("room_no").asLong();
Long messageNo = after.path("message_no").asLong();
String uuidStr = after.path("public_uuid").asText();
if (uuidStr == null) {
log.warn("public_uuid null. payload={}", payload);
return;
}
// 🧩 UUID 변환
byte[] bytes = Base64.getDecoder().decode(uuidStr);
ByteBuffer bb = ByteBuffer.wrap(bytes);
UUID publicUuid = new UUID(bb.getLong(), bb.getLong());
boolean isRead = after.path("is_read").asBoolean(false);
JsonNode batchNode = after.get("batch_in");
if (batchNode == null || batchNode.isNull()) return;
Long batchIn = batchNode.asLong();
if(!isRead) return;
// 🔑 idempotent 처리 (중복 처리 방지)
String idempotentKey = String.valueOf(batchIn);
Boolean firstProcess = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "", Duration.ofHours(1));
if(Boolean.TRUE.equals(firstProcess)){
// 📡 WebSocket 전송용 DTO 생성
MessageReadStatusDto.Response messageReadStatusDto = MessageReadStatusDto.Response.builder()
.room_no(roomNo)
.message_no(messageNo)
.public_uuid(publicUuid)
.type(SocketEnums.type.READ)
.build();
// 👥 채팅방 구독자 조회
Set<WebSocketSession> sessions = chatRoomWebSocketHandler.getRoomSesstions(roomNo);
for(WebSocketSession session : sessions) {
if (session != null && session.isOpen()) {
String json = objectMapper.writeValueAsString(messageReadStatusDto);
session.sendMessage(new TextMessage(json)); // 🔔 읽음 상태 전파
}
}
}
}public void chatReadDebezium(String payload) throws IOException {
JsonNode root = objectMapper.readTree(payload);
String op = root.get("op").asText(); // u=update, c=create, d=delete
if (!"u".equals(op)) return; // update 이벤트만 처리- Debezium이 보내는 이벤트는 c(create), u(update), d(delete)가 있음
- 메시지 읽음 상태는 업데이트 이벤트이므로 u만 처리
JsonNode after = root.path("after");
if (after.isMissingNode() || after.isNull()) {
log.warn("after is null. payload={}", payload);
return;
}- Debezium 이벤트는 before와 after 데이터 구조를 가짐
- after가 null이면 의미 있는 업데이트가 없으므로 종료
Long roomNo = after.path("room_no").asLong();
Long messageNo = after.path("message_no").asLong();
String uuidStr = after.path("public_uuid").asText();
if (uuidStr == null) {
log.warn("public_uuid null. payload={}", payload);
return;
}- 채팅방 번호(roomNo), 메시지 번호(messageNo), 메시지 식별 UUID(publicUuid)를 가져옴
- UUID가 없으면 읽음 전파 불가 → 종료
byte[] bytes = Base64.getDecoder().decode(uuidStr);
ByteBuffer bb = ByteBuffer.wrap(bytes);
UUID publicUuid = new UUID(bb.getLong(), bb.getLong());- MySQL에 저장된 UUID는 Base64 인코딩되어 있음
- 이를 UUID 객체로 변환하여 DTO에 넣을 수 있도록 준비
boolean isRead = after.path("is_read").asBoolean(false);
JsonNode batchNode = after.get("batch_in");
if (batchNode == null || batchNode.isNull()) return;
Long batchIn = batchNode.asLong();
if(!isRead) return;- is_read가 false이면 전파할 필요 없음 → 종료
- batch_in은 중복 전파 방지를 위한 고유 키
String idempotentKey = String.valueOf(batchIn);
Boolean firstProcess = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "", Duration.ofHours(1));- Redis setIfAbsent를 사용해 이미 처리된 이벤트는 무시
- TTL 1시간 설정 → 1시간 내 동일 이벤트 중복 전파 방지
if(Boolean.TRUE.equals(firstProcess)){
MessageReadStatusDto.Response messageReadStatusDto = MessageReadStatusDto.Response.builder()
.room_no(roomNo)
.message_no(messageNo)
.public_uuid(publicUuid)
.type(SocketEnums.type.READ)
.build();- 채팅방 번호, 메시지 번호, 메시지 UUID, 타입(READ)을 DTO로 묶음
Set<WebSocketSession> sessions = chatRoomWebSocketHandler.getRoomSesstions(roomNo);
for(WebSocketSession session : sessions) {
if (session != null && session.isOpen()) {
String json = objectMapper.writeValueAsString(messageReadStatusDto);
session.sendMessage(new TextMessage(json)); // 🔔 읽음 상태 전파
}
}
}
}- 채팅방에 연결된 모든 WebSocket 세션 조회
- 연결이 살아있으면 DTO를 JSON으로 변환 후 메시지 전송
- 결과적으로 읽음 상태가 실시간으로 반영됨
💡 정리
- Debezium → Kafka → Consumer 흐름으로 이벤트 수신
- update 이벤트만 필터링
- is_read 체크 + null 체크
- Redis를 이용한 idempotency 처리
- WebSocket으로 채팅방 구독자에게 읽음 상태 전송