Post

비동기 이벤트 작업 동적 스케쥴러 + 배치처리로 개선하기

🎬 Intro

실시간 비동기 이벤트 처리 작업을 동적 스케쥴러 + 배치로 전환하는 과정을 정리해봅니다.

결론

결론부터 말씀드리면 기존 비동기 이벤트로 클러스터링 되던 로직을 동적 스케쥴러 + 배치처리를 통해 성능을 개선하였습니다.

구체적으로 POST 요청 1번당 외부 API 호출 1회 되던 부분을 동적 주기 배치로 최대 100건씩 묶어 처리하도록 변경했습니다.

결과적으로 외부 API 호출 횟수 -> 최대 1/100로 감소하였고, 약 7초이상의 처리 지연0.06초로 개선하였습니다.

현재 상황

  • POST 요청 1번 당(피드백 작성) 비동기 이벤트 형식으로 벡터 계산을 위한 외부 API 호출

현재 상황이 왜 문제가 되는가?

AI 클러스터링하는 기능은 비슷한 피드백을 군집으로 묶어주는 역할을 합니다. 따라서 통계의 의미에 가까운 상황이므로 해당 기능은 실시간성이 필요 없다. 가 없었습니다. 실시간성이 필요 없다면 굳이 비동기 이벤트로 로직을 처리할 이유가 없습니다.

img.png img.png

또한 클러스터링 작업의 경우 외부 API 호출이 필요하므로 갑작스럽게 피드백이 많이 작성된다면, 위와 같이 처리시간이 점점 느려지는 병목이 생기게 됩니다.

개선 방향

실시간 비동기 이벤트로 처리하던 클러스터링 로직을 동적 스케줄러 기반 배치 처리로 전환하였습니다. 배치 단위와 실행 주기는 외부 API 응답 시간에 따라 동적으로 조정하였습니다.

✅ 스케쥴러 주기 + 배치 단위 설정

스케쥴러와 배치 처리 로직을 구현하기에 앞서 스케쥴러 주기와 배치 단위를 먼저 결정해야합니다.

배치 단위 = 100개

img.png voyage ai 공식 문서에 따르면 저희 서비스가 사용하고 있는 voyage-3.5 모델의 경우 limit이 TPM = 8,000,000 / RPM = 2,000 입니다. TPM은 분당 요청할 수 있는 토큰 수, RPM은 분당 요청할 수 있는 요청 수를 의미합니다.

  • 한글 1글자 당 3개의 토큰 사용
  • 매 강의시 100글자 이내의 피드백이 50개 생성 됨

위 상황을 고려했을때 50개의 피드백을 요청한다면, 3 * 100 * 50 = 15,000 토큰이 사용 됩니다. 만약 강의가 10개라고 가정해도 TPM은 여유로운 상황입니다. 따라서 서비 메모리와 RPM을 고려해서 배치 단위는 100개로 설정하였습니다.

스케쥴러 주기

배치 단위 100개일 경우 평균 응답 시간이 5초 입니다. 이를 고려하여 기본 스케쥴러 주기를 1분 으로 하고 다음과 같이 동적으로 계산하여 설정하였습니다.

1
2
응답 시간 ≤ 5,000ms  : 폴링 주기 = 60,000ms (고정)
응답 시간 > 5,000ms  : 폴링 주기 = 60,000ms × (응답 시간 / 5,000ms)

✅ 스케쥴러 + 배치 로직 구현

스케쥴러

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@SchedulerLock(name = "feedbackClusteringBatch", lockAtMostFor = "PT30M")
public void clusterUnclusteredFeedbacks() {
    long startTime = System.currentTimeMillis();
    log.info("[클러스터링 배치] 시작");
    List<FeedbackClusteringQueue> queues = queueRepository.findPendingQueues(BATCH_SIZE);
    
    /*
  
      클러스터링 작업
      
     */
  
    long responseTime = System.currentTimeMillis() - startTime; 
    intervalCalculator.updateExecutionTime(responseTime);
}
  • BATCH_SIZE = 100으로 하여 100개씩 클러스터링 작업을 묶어서 진행합니다
  • 클러스터링 배치 작업 이후 응답 시간을 계산하여 스케쥴링 주기를 설정합니다

스케쥴러 등록

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@Configuration
@RequiredArgsConstructor
public class FeedbackClusteringSchedulingConfig implements SchedulingConfigurer {

    private final FeedbackClusteringBatchScheduler feedbackClusteringBatchScheduler;
    private final FeedbackClusteringAdaptiveTrigger feedbackClusteringAdaptiveTrigger;

    @Override
    public void configureTasks(final ScheduledTaskRegistrar taskRegistrar) {
        taskRegistrar.addTriggerTask(
                feedbackClusteringBatchScheduler::clusterUnclusteredFeedbacks,
                feedbackClusteringAdaptiveTrigger
        );
    }
}
  • 클러스터링 스케쥴러를 Spring TaskScheduler에 등록하는 설정 클래스 입니다

스케쥴러 주기 재설정

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
@Component
@RequiredArgsConstructor
public class FeedbackClusteringAdaptiveTrigger implements Trigger {

  private final FeedbackClusteringSchedulerTimeCalculator feedbackClusteringSchedulerTimeCalculator;

  @Override
  public Instant nextExecution(final TriggerContext triggerContext) {
    Instant lastCompletion = triggerContext.lastCompletion();
    Instant now = Instant.now();
    Instant baseTime = (lastCompletion != null) ? lastCompletion : now;

    Duration interval = feedbackClusteringSchedulerTimeCalculator.calculateNextInterval();
    return baseTime.plus(interval);
  }
}
  • Spring TaskScheduler에 등록된 클러스터링의 스케쥴러 주기를 재설정하는 클래스 입니다.
  • 다음 스케쥴러 실행 시간은 현재 시간 + 동적으로 계산된 스케줄링 주기 입니다.

벌크 연산

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
@Slf4j
@Component
@RequiredArgsConstructor
public class EmbeddingClusterBulkInserter {

    private final EmbeddingClusterRepository embeddingClusterRepository;
    private final FeedbackEmbeddingClusterRepository feedbackEmbeddingClusterRepository;

    @Transactional(propagation = Propagation.REQUIRES_NEW)
    public void bulkInsertEmbeddingClusters(List<EmbeddingCluster> newEmbeddingClusters) {
        log.info("EmbeddingClusters Bulk Insert Start");
        embeddingClusterRepository.bulkInsert(newEmbeddingClusters);
    }

    @Transactional(propagation = Propagation.REQUIRES_NEW)
    public void bulkInsertFeedbackEmbeddingClusters(List<FeedbackEmbeddingCluster> feedbackEmbeddingClusters) {
        log.info("FeedbackEmbeddingClusters Bulk Insert Start");
        feedbackEmbeddingClusterRepository.bulkInsert(feedbackEmbeddingClusters);
    }
}
  • 배치 작업이 완료된 클러스터링을 저장하기 위해 JDBC 벌크 연산을 사용하였습니다.

img.png

마지막으로 로직에 사용된 컴포넌트를 보자면 총 8개 였고, 각 역할은 다음과 같습니다.

  • FeedbackClusteringBatchScheduler
    • 2시간마다 실행되는 스케줄러, PENDING 상태인 피드백을 100개씩 조회하고 조직별로 그룹핑하여 배치 처리를 조율
  • FeedbackClusteringBatchService
    • 100개 피드백의 임베딩 추출 -> 클러스터링 -> 상태 업데이트 및 재시도 카운팅 처리
  • FeedbackClusteringService
    • EmbeddingClustering 라벨 생성
  • EmbeddingClusteringBulkInserter
    • FeedbackEmbeddingCluster Bulk Insert
    • EmbeddingCluster Bulk Insert
  • EmbeddingClusterRepository
    • EmbeddingCluster 도메인 JPA 레파지토리
  • EmbeddingClusterCustomRepository
    • EmbeddingCluster 도메인 JDBC 레파지토리 (벌크 연산)
  • FeedbackEmbeddingClusterRepository
    • FeedbackEmbeddingClusterRepository 도메인 JPA 레파지토리
  • FeedbackEmbeddingClusterCustomRepository
    • FeedbackEmbeddingClusterCustomRepository 도메인 JDBC 레파지토리 (벌크 연산)

성과

✅ 반복적인 외부 API 호출 시간 제거 -> 1/100회

Before

img.png

After

img.png

✅ 약 7초 이상 처리 지연 -> 0.06초로 개선

Before

img.png

After

img.png

마무리 하며

이번 리팩토링은 성능 개선에 집중하였습니다. 이를 위해 스케줄러 + 배치 처리, JDBC 벌크 연산을 활용해보았고, 그 과정에서 shedLock 개념도 학습해볼 수 있었습니다.

완벽한 설계는 아니지만, 팀 내에서 고민이었던 병목 현상을 해결할 수 있었어서 뿌듯했습니다.

This post is licensed under CC BY 4.0 by the author.