금학기 교환학교 리스트가 발표되면서 급증하는 서버 요청들을 모니터링 하던 도중,
특정 api 요청들에서 응답시간이 급격히 지연되는 문제를 발견하였다.


구글이 발표한 자료에 따르면 응답시간이 길어지면 길어질수록 사용자 이탈율이 증가함에 따라서 권장하는 max 응답시간은 2.5s이다.
이에 사용자가 불편함을 느낄 수 있다고 생각되는 api들의 응답시간을 개선하기위해, 캐싱을 도입하게 되었다.
데이터 캐싱을 위해서는 간단히 Map과 같은 local memory를 활용할 수도 있지만,
scale out 등 확장성 측면 및 도입 용이성 등을 고려하여 Redis 저장소를 사용하기로 하였다.
(Redis는 기존에 도입되어있었다.)
사실 스프링에서 제공하는 @Cacheable 어노테이션과 @CachePut @CacheEvict 등의 어노테이션을 활용하여 캐싱기능을 쉽게 적용할 수 있다.
하지만 기존 기능을 그대로 활용할 경우, 동작을 세밀하게 제어하기 어렵다.
우리 프로젝트에서 캐싱을 도입할 때, ttl 관리, prefix 기반으로 캐시아웃, Thundering Herd 문제해결 등이 필요한데 이를 해결할 수 없기에 캐시 매니저와 어노테이션을 커스텀하여 개발하기로 하였다.
캐시아웃 시점에 다수의 동시요청이 오게되면 과연 어떻게 동작할까?
기본 캐싱전략을 사용하게되면 모든 요청들이 DB에 접근하여 데이터를 가져오는 과정을 거치게 된다.
이에 DB 부하는 급증할 수 밖에 없다.
그렇다면 어떻게 동작하면 좋을까?
내가 생각한 방안은 특정 키로 lock을 사용하여 해당 lock을 잡은 요청만 DB에 접근하여 데이터를 가져오도록 하고,
나머지 요청들은 그동안 대기하다가 DB 접근한 요청이 데이터를 캐싱하고 대기중인 요청들에게 알리면
그때, 나머지 요청들은 캐시에서 데이터를 가져오는 것이였다.
이렇게 수행하게되면 수많은 동시요청이 하나의 DB 접근으로 해결될 수 있게되어 DB 부하를 충분히 줄여줄 수 있게된다.
위에서 언급한바와 같이 lock을 잡지 못한 나머지 모든 요청들은 대기상태로 들어가야한다.
하지만 어떤 방식으로 대기시키는게 좋을까?
우선은 lock 잡기에 실패한 요청의 경우 굳이 lock 잡기를 재시도할 필요가 없다.
단순히 캐싱되었을때 캐싱된 값을 받아오기만 하면 되기 때문이다.
따라서 재시도 로직이나 스핀락등은 제외한다.
Thread.sleep()을 사용하면 어떨까?
정말 단순히 접근하기 쉽지만 이역시 불필요한 오버헤드가 발생한다.
캐싱작업이 대기상태로 들어가고 0.0000000001s 이후에 완료될 수도 있는데 sleep을 진행하게되면 특정 시간동안 필수적으로 대기상태에 빠질 수 밖에 없게되어 불필요한 오버헤드가 발생한다.
Future를 사용하면 어떨까?
단순히 Future만 사용하게되면 특정 was(lock을 잡은 스레드가 존재하는)에서만 대기중인 요청들에게 알려 깨워줄 수 있게된다. 따라서 scale out을 고려할 때 아쉬운 해결방안이다.
따라서 우리는 scale out등 확장성, 오버헤드 등을 고려하여
Future를 사용하되, Redis의 pub/sub 기능을 활용하여 모든 was에서 대기중인 요청들에게 일괄적으로 알려줄 수 있고, 이벤트 발생 시 즉각 깨어나 불필요한 오버헤드를 줄일 수 있도록 하였다.
위의 과정들을 통해서 다수의 중복된 DB 조회 작업을 1개의 작업으로 해결할 수 있게 되었다.
하지만 캐시 도입 이후, 캐시 히트율에대한 고민을 빼놓을 수 없다.
대부분 warming이라는 작업을 통해서 자주 접근하는 데이터들을 캐시에 올려두어 hit율을 높이곤 하는데 이에 착안하여 hit율을 높여 보았다.
현재 데이터가 캐싱되면 추후 TTL이 끝나는 시점에 데이터가 캐시아웃되도록 구현되어있는데,
이때 TTL이 얼마남지 않은 시점 예를들어 설정 TTL이 10%이하로 남아있는 시점 등에 요청이 오게되면 locality를 고려하여 해당 키를 다시 갱신하도록 하여 캐시 hit율을 높이도록 구현하였다.
이 역시 다수의 요청이 동시에 발생할 수 있기에 lock을 사용하여 하나의 요청만 캐시 갱신 작업을 진행하도록 하고 나머지 요청들은 캐시에 저장되어있던 데이터를 바로 반환하도록 하여 불필요한 오버헤드를 줄였다.
커스텀 캐싱 어노테이션 선언
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface ThunderingHerdCaching {
String key();
String cacheManager();
long ttlSec();
}
커스텀 캐시아웃 어노테이션 선언
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface CustomCacheOut {
String[] key();
String cacheManager();
boolean prefix() default false;
}
CacheManager 인터페이스 선언
public interface CacheManager {
void put(String key, Object value, Long ttl);
Object get(String key);
void evict(String key);
void evictUsingPrefix(String key);
}
커스텀 캐시 매니저 구현부이며 추가적으로 prefix를 활용하여 다수의 key를 캐시아웃 시키는 기능을 구현한다.
@Component("customCacheManager")
public class CustomCacheManager implements CacheManager {
private final RedisTemplate<String, Object> redisTemplate;
@Autowired
public CustomCacheManager(RedisTemplate<String, Object> redisTemplate) {
redisTemplate.setKeySerializer(new StringRedisSerializer());
redisTemplate.setValueSerializer(new GenericJackson2JsonRedisSerializer());
this.redisTemplate = redisTemplate;
}
public void put(String key, Object object, Long ttl) {
redisTemplate.opsForValue().set(key, object, Duration.ofSeconds(ttl));
}
public Object get(String key) {
return redisTemplate.opsForValue().get(key);
}
public void evict(String key) {
redisTemplate.delete(key);
}
public void evictUsingPrefix(String key) {
Set<String> keys = redisTemplate.keys(key+"*");
if (keys != null && !keys.isEmpty()) {
redisTemplate.delete(keys);
}
}
}
대기 스레드를 관리하기 위한 CompletableFutureManager 정의
@Component
public class CompletableFutureManager {
private final Map<String, CompletableFuture<Void>> waitingRequests = new ConcurrentHashMap<>();
public CompletableFuture<Void> getOrCreateFuture(String key) {
return waitingRequests.computeIfAbsent(key, k -> new CompletableFuture<>());
}
public void completeFuture(String key) {
CompletableFuture<Void> future = waitingRequests.remove(key);
if (future != null) {
future.complete(null);
}
}
}
@CustomCacheOut 관련 Aspect 정의
@Aspect
@Component
@RequiredArgsConstructor
public class CachingAspect {
,,,
@Around("@annotation(defaultCacheOut)")
public Object cacheEvict(ProceedingJoinPoint joinPoint, DefaultCacheOut defaultCacheOut) throws Throwable {
CacheManager cacheManager = (CacheManager) applicationContext.getBean(defaultCacheOut.cacheManager());
for (String key : defaultCacheOut.key()) {
String cacheKey = redisUtils.generateCacheKey(key, joinPoint.getArgs());
boolean usingPrefix = defaultCacheOut.prefix();
if (usingPrefix) {
cacheManager.evictUsingPrefix(cacheKey);
}else{
cacheManager.evict(cacheKey);
}
}
return joinPoint.proceed();
}
}
@ThunderingHerdCaching 관련 Aspect 정의
@Aspect
@Component
@Slf4j
public class ThunderingHerdCachingAspect {
,,,
@Around("@annotation(thunderingHerdCaching)")
public Object cache(ProceedingJoinPoint joinPoint, ThunderingHerdCaching thunderingHerdCaching) {
CacheManager cacheManager = (CacheManager) applicationContext.getBean(thunderingHerdCaching.cacheManager());
String key = redisUtils.generateCacheKey(thunderingHerdCaching.key(), joinPoint.getArgs());
Long ttl = thunderingHerdCaching.ttlSec();
Object cachedValue = cacheManager.get(key);
if (cachedValue == null) {
log.info("Cache miss. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return createCache(joinPoint, cacheManager, ttl, key);
}
if (redisUtils.isCacheExpiringSoon(key, ttl, Double.valueOf(REFRESH_LIMIT_PERCENT.getValue()))) {
log.info("Cache hit, but TTL is expiring soon. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return refreshCache(cachedValue, ttl, key);
}
log.info("Cache hit. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return cachedValue;
}
private Object createCache(ProceedingJoinPoint joinPoint, CacheManager cacheManager, Long ttl, String key) {
return executeWithLock(
redisUtils.getCreateLockKey(key),
() -> {
log.info("생성락 흭득하였습니다. Key: {}, Thread: {}", key, Thread.currentThread().getName());
Object result = proceedJoinPoint(joinPoint);
cacheManager.put(key, result, ttl);
redisTemplate.convertAndSend(CREATE_CHANNEL.getValue(), key);
log.info("캐시 생성 후 채널에 pub 진행합니다. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return result;
},
() -> {
log.info("생성락 흭득에 실패하여 대기하러 갑니다. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return waitForCacheToUpdate(joinPoint, key);
}
);
}
private Object refreshCache(Object cachedValue, Long ttl, String key) {
return executeWithLock(
redisUtils.getRefreshLockKey(key),
() -> {
log.info("갱신락 흭득하였습니다. Key: {}, Thread: {}", key, Thread.currentThread().getName());
redisTemplate.opsForValue().getAndExpire(key, Duration.ofSeconds(ttl));
log.info("TTL 갱신을 마쳤습니다. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return cachedValue;
},
() -> {
log.info("갱신락 흭득에 실패하였습니다. 캐시의 값을 바로 반환합니다. Key: {}, Thread: {}", key, Thread.currentThread().getName());
return cachedValue;
}
);
}
private Object executeWithLock(String lockKey, Callable<Object> onLockAcquired, Callable<Object> onLockFailed) {
String lockValue = UUID.randomUUID().toString();
boolean lockAcquired = false;
try {
lockAcquired = tryAcquireLock(lockKey, lockValue);
if (lockAcquired) {
return onLockAcquired.call();
} else {
return onLockFailed.call();
}
} catch (Exception e) {
throw new RuntimeException("Error during executeWithLock", e);
} finally {
releaseLock(lockKey, lockValue, lockAcquired);
}
}
private boolean tryAcquireLock(String lockKey, String lockValue) {
return redisTemplate.opsForValue().setIfAbsent(lockKey, lockValue, Duration.ofMillis(Long.parseLong(LOCK_TIMEOUT_MS.getValue())));
}
private void releaseLock(String lockKey, String lockValue, boolean lockAcquired) {
if (lockAcquired && lockValue.equals(redisTemplate.opsForValue().get(lockKey))) {
redisTemplate.delete(lockKey);
log.info("락 반환합니다. Key: {}", lockKey);
}
}
private Object waitForCacheToUpdate(ProceedingJoinPoint joinPoint, String key) {
,,,
}
}
동작과정 확인 및 이해를 돕기 위한 테스트 로그
초기에 캐싱된 데이터가 없을 때, 다수의 요청이 오면 각각의 키마다 하나의 요청만 락을 흭득하여 DB 접근 후 데이터를 캐싱합니다.
데이터 캐싱 작업이 끝나면 publish를 진행하고 해당 이벤트를 수신한 스레드들이 대기에서 풀려나 캐시값을 가져와 반환합니다.

캐싱 이후 TTL이 많이 남았을 때 hit되는 상황입니다.

TTL이 얼마남지 않은 상황에 다수의 요청이 들어오면 키마다 하나의 요청만 락을 흭득하여 데이터를 갱신하고, 나머지 요청들은 바로 캐싱되어있던 값을 반환합니다.

50rps가 유지되는 상황 가정 (vuser 50 && 각각 초당 1회 요청)
before
api-server: p(95)=40s, p(99)=45s
local-server: p(95)=300ms, p(99)=400ms
after
api-server: p(95)=700ms, p(99)=950ms
local-server: p(95)=55ms, p(99)=70ms
항상 서비스 개선에 있어서 단순히 발생한 문제를 어떻게 더 효율적으로 해결해나갈지 고민해왔다.
하지만, 한번 더 개선하기위해, 나만의 전략을 세우기 위해서 크게 고민해온적은 없었던 것 같다.
이번 경험에서 얻은걸 토대로 앞으로는 문제 해결시에 '조금 더 개선할 순 없을까?'라는 질문을 스스로 해보는 습관을 가져보려한다.