Skip to content

Commit b7a8585

Browse files
authored
fix: 대기열 입장 순서 보존 및 진단 로그 추가 (#229)
1 parent e00377d commit b7a8585

3 files changed

Lines changed: 43 additions & 12 deletions

File tree

src/main/java/com/threestar/trainus/domain/lesson/issue/LegacyLessonAdmissionScheduler.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ public void admitUsers() {
4646

4747
private void processAdmissionForLesson(Long lessonId) {
4848
// 해당 레슨 대기열에서 인원 추출
49-
Set<String> requestIds = waitingRoomService.dequeue(lessonId, ADMIT_BATCH_SIZE);
49+
List<String> requestIds = waitingRoomService.dequeue(lessonId, ADMIT_BATCH_SIZE);
5050

5151
if (requestIds.isEmpty()) {
5252
return;
@@ -69,7 +69,7 @@ private void processAdmissionForLesson(Long lessonId) {
6969
if (parts.length < 3)
7070
continue;
7171

72-
String requestId = requestIds.toArray(new String[0])[i];
72+
String requestId = requestIds.get(i);
7373
Long userId = Long.parseLong(parts[2]);
7474
Long originalTimestamp = parts.length >= 4 ? Long.parseLong(parts[3]) : System.currentTimeMillis();
7575

src/main/java/com/threestar/trainus/domain/lesson/issue/LessonAdmissionScheduler.java

Lines changed: 36 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
import java.util.List;
55
import java.util.Map;
66
import java.util.Set;
7+
import java.util.concurrent.atomic.AtomicInteger;
78

89
import org.springframework.beans.factory.annotation.Qualifier;
910
import org.springframework.context.annotation.Profile;
@@ -56,7 +57,7 @@ private void processAdmissionForLesson(Long lessonId) {
5657
}
5758

5859
// 해당 레슨 대기열에서 인원 추출
59-
Set<String> requestIds = waitingRoomService.dequeue(lessonId, ADMIT_BATCH_SIZE);
60+
List<String> requestIds = waitingRoomService.dequeue(lessonId, ADMIT_BATCH_SIZE);
6061

6162
if (requestIds.isEmpty()) {
6263
return;
@@ -68,22 +69,44 @@ private void processAdmissionForLesson(Long lessonId) {
6869
.toList();
6970
List<String> statusInfos = mqRedisTemplate.opsForValue().multiGet(statusKeys);
7071

71-
if (statusInfos == null) return;
72+
if (statusInfos == null) {
73+
log.warn(
74+
"Admission diagnostics. lessonId={}, dequeuedCount={}, statusKeyCount={}, statusInfos=null. Admission skipped after dequeue.",
75+
lessonId, requestIds.size(), statusKeys.size());
76+
return;
77+
}
78+
79+
AtomicInteger statusNullCount = new AtomicInteger();
80+
AtomicInteger invalidStatusCount = new AtomicInteger();
81+
AtomicInteger xaddAttemptCount = new AtomicInteger();
7282

7383
// Pipelining 방식으로 요청 상태 변경 / Stream ADD 로직 일괄 처리
7484
mqRedisTemplate.executePipelined(new SessionCallback<Object>() {
7585
@Override
7686
public Object execute(RedisOperations operations) {
7787
for (int i = 0; i < statusKeys.size(); i++) {
7888
String info = statusInfos.get(i);
79-
if (info == null) continue;
89+
if (info == null) {
90+
statusNullCount.incrementAndGet();
91+
continue;
92+
}
8093

8194
String[] parts = info.split(":");
82-
if (parts.length < 3) continue;
83-
84-
String requestId = requestIds.toArray(new String[0])[i];
85-
Long userId = Long.parseLong(parts[2]);
86-
Long originalTimestamp = parts.length >= 4 ? Long.parseLong(parts[3]) : System.currentTimeMillis();
95+
if (parts.length < 3) {
96+
invalidStatusCount.incrementAndGet();
97+
continue;
98+
}
99+
100+
String requestId = requestIds.get(i);
101+
Long userId;
102+
Long originalTimestamp;
103+
try {
104+
userId = Long.parseLong(parts[2]);
105+
originalTimestamp = parts.length >= 4 ? Long.parseLong(parts[3]) : System.currentTimeMillis();
106+
} catch (NumberFormatException e) {
107+
invalidStatusCount.incrementAndGet();
108+
continue;
109+
}
87110

88111
// 상태 변경 (SET)
89112
String statusKey = statusKeys.get(i);
@@ -97,11 +120,16 @@ public Object execute(RedisOperations operations) {
97120
content.put("requestId", requestId);
98121
content.put("timestamp", String.valueOf(originalTimestamp));
99122
operations.opsForStream().add(LessonApplyStreamConstant.STREAM_KEY, content);
123+
xaddAttemptCount.incrementAndGet();
100124
}
101125
return null;
102126
}
103127
});
104128

129+
log.info(
130+
"Admission diagnostics. lessonId={}, dequeuedCount={}, statusKeyCount={}, statusNullCount={}, invalidStatusCount={}, xaddAttemptCount={}",
131+
lessonId, requestIds.size(), statusKeys.size(), statusNullCount.get(), invalidStatusCount.get(),
132+
xaddAttemptCount.get());
105133
log.debug("Admitted {} users to MQ via pipeline for lesson: {}", requestIds.size(), lessonId);
106134
}
107135
}

src/main/java/com/threestar/trainus/domain/lesson/issue/LessonWaitingRoomService.java

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
package com.threestar.trainus.domain.lesson.issue;
22

33
import java.time.Duration;
4+
import java.util.Comparator;
5+
import java.util.List;
46
import java.util.Optional;
57

68
import org.springframework.beans.factory.annotation.Qualifier;
@@ -51,12 +53,13 @@ public Optional<Long> getRank(Long lessonId, String requestId) {
5153
}
5254

5355
// 특정 레슨 대기열에서 가장 오래된 N명을 꺼내기
54-
public java.util.Set<String> dequeue(Long lessonId, long count) {
56+
public List<String> dequeue(Long lessonId, long count) {
5557
String waitingRoomKey = String.format(LessonApplyStreamConstant.WAITING_ROOM_KEY, lessonId);
5658
return coreRedisTemplate.opsForZSet()
5759
.popMin(waitingRoomKey, count)
5860
.stream()
61+
.sorted(Comparator.comparing(tuple -> tuple.getScore() == null ? Double.MAX_VALUE : tuple.getScore()))
5962
.map(tuple -> tuple.getValue())
60-
.collect(java.util.stream.Collectors.toSet());
63+
.toList();
6164
}
6265
}

0 commit comments

Comments
 (0)