이번 프로젝트에서는 사용자가 생성한 투자 전략과 백테스트 결과를 LLM에 전달하고, 위험 등급과 위험 요인을 분석하는 기능을 구현했다.
AI 분석은 일반적인 DB 조회와 달리 LLM API 호출 시간이 필요하다. 따라서 사용자가 분석을 요청할 때 LLM 응답이 끝날 때까지 HTTP 요청을 유지하는 대신, 분석 요청과 실제 AI 처리를 분리하는 비동기 구조를 적용했다.
1. 처음 설계한 AI 분석 흐름
AI 분석 요청이 들어오면 먼저 PostgreSQL에 분석 데이터를 PENDING 상태로 저장한다.
POST /api/v1/ai/analyses
│
▼
전략 / 백테스트 검증
│
▼
AiRiskAnalysis 저장
status = PENDING
│
▼
이벤트 발행
│
▼
@Async
│
▼
LLM API 호출
│
▼
응답 파싱 및 검증
│
┌────┴────┐
▼ ▼
COMPLETED FAILED
이렇게 구성하면 HTTP 요청에서 LLM 처리가 끝날 때까지 기다릴 필요가 없다.
요청 트랜잭션에서는 분석 요청만 저장하고 실제 AI 분석은 별도의 비동기 작업으로 처리한다.
2. 왜 AFTER_COMMIT 이벤트를 사용했는가?
분석 데이터가 DB에 정상적으로 저장된 이후에만 AI 분석을 시작하고 싶었다.
그래서 분석 생성 트랜잭션에서 이벤트를 발행하고:
AiRiskAnalysis savedAnalysis =
aiRiskAnalysisCommandRepository.save(analysis);
eventPublisher.publishEvent(
new AiRiskAnalysisRequestedEvent(
savedAnalysis.getId(),
strategy,
backtest
)
);
리스너에서는:
@TransactionalEventListener(
phase = TransactionPhase.AFTER_COMMIT
)
을 사용했다.
중요한 점은 이벤트를 발행한 순간 바로 AI 분석을 시작하는 것이 아니라, 해당 트랜잭션이 실제 COMMIT된 이후 처리한다는 것이다.
이를 통해 DB 저장이 롤백되었는데 AI 분석만 실행되는 상황을 방지할 수 있다.
3. @Async를 붙였다고 끝이 아니었다
처음에는 이벤트 처리 메서드 자체에 @Async를 적용하는 형태를 생각했다.
문제는 Thread Pool이 처리할 수 있는 양을 넘어서는 순간 발생한다.
예를 들어 Executor가:
corePoolSize = 2
maxPoolSize = 5
queueCapacity = 50
이라면 작업이 몰렸을 때 실행 중인 작업과 Queue가 모두 차면서 새로운 비동기 작업이 거부될 수 있다.
이때 발생할 수 있는 예외가 TaskRejectedException이다.
문제는 단순히 요청 하나가 실패하는 것이 아니었다.
DB에 PENDING 저장
↓
COMMIT
↓
이벤트 발생
↓
@Async 작업 제출
↓
Thread Pool 포화
↓
TaskRejectedException
↓
AI 분석 시작조차 못함
이미 DB Transaction은 COMMIT되었기 때문에 분석 데이터는 존재한다.
그런데 비동기 작업은 실행되지 않았다.
결국:
status = PENDING
인 데이터가 계속 남을 수 있다.
4. 이벤트 수신과 비동기 실행을 분리했다
이를 해결하기 위해 이벤트를 받는 책임과 비동기 작업을 실행하는 책임을 분리했다.
이벤트 리스너
@TransactionalEventListener(
phase = TransactionPhase.AFTER_COMMIT
)
public void handle(AiRiskAnalysisRequestedEvent event) {
try {
asyncProcessor.process(event);
} catch (TaskRejectedException e) {
resultService.fail(
event.analysisId(),
AiAnalysisFailureType.INTERNAL_ERROR,
"AI 분석 작업 실행이 거부되었습니다."
);
}
}
리스너는 동기적으로 이벤트를 받은 다음 AsyncProcessor에게 작업 제출을 요청한다.
비동기 Processor
@Async("aiAnalysisExecutor")
public void process(AiRiskAnalysisRequestedEvent event) {
processor.process(event);
}
이렇게 분리하면 비동기 작업을 Executor에 제출하는 시점의 실패를 이벤트 리스너가 감지할 수 있다.
따라서:
Thread Pool 정상
↓
비동기 AI 분석 수행
Thread Pool 포화
↓
TaskRejectedException
↓
FAILED / INTERNAL_ERROR
로 처리할 수 있고, 실행되지 못한 분석이 계속 PENDING으로 남는 문제를 줄일 수 있다.
5. 그런데 서버가 여러 대라면?
여기까지 구현한 뒤 한 가지 문제가 더 보였다.
현재 사용하는 ApplicationEventPublisher는 Kafka 같은 메시지 브로커가 아니다.
Spring Application Event는 해당 애플리케이션 인스턴스 내부에서 동작한다.
예를 들어 Trading Service를 세 대로 확장했다고 가정하면:
Load Balancer
/ | \
/ | \
Trading-1 Trading-2 Trading-3
│
요청 수신
│
PENDING 저장
│
ApplicationEvent
│
▼
Trading-1 내부
Trading-1에서 발생한 Spring Event를 Trading-2, Trading-3가 받아 처리하는 것은 아니다.
다만 이것 자체는 현재 기능에서는 반드시 문제가 되는 것은 아니다.
요청을 받은 Trading-1이 자기 인스턴스에서 AI 분석까지 정상적으로 수행하면 되기 때문이다.
진짜 문제는 장애가 발생했을 때의 신뢰성이다.
6. COMMIT 이후 서버가 죽는다면?
다음 상황을 생각해볼 수 있다.
Trading-1
│
├─ PENDING 저장
│
├─ COMMIT 성공
│
├─ 이벤트 처리
│
💥 서버 장애
PostgreSQL에는 이미:
status = PENDING
데이터가 저장되어 있다.
하지만 Spring Application Event는 애플리케이션 내부 이벤트이므로 서버가 죽었다고 해서 다른 인스턴스가 해당 이벤트를 이어받는 구조가 아니다.
서버를 재시작한다고 기존 이벤트가 복구되는 것도 아니다.
결국 DB에는 분석 요청이 존재하지만 이를 처리할 작업은 사라지는 상황이 가능하다.
Thread Pool 포화 문제는 TaskRejectedException을 잡아서 해결할 수 있었지만, 프로세스 자체가 종료되는 문제까지 현재 구조로 해결할 수는 없다.
7. 다중 인스턴스 환경에서는 Kafka를 고려할 수 있다
이 문제를 개선한다면 Spring 내부 이벤트 대신 메시지 브로커를 이용하는 방법을 생각할 수 있다.
예를 들어:
POST AI 분석 요청
↓
PENDING 저장
↓
Kafka
ai.analysis.requested
↓
Consumer Group
↓
┌─────────┬─────────┬─────────┐
▼ ▼ ▼
Trading-1 Trading-2 Trading-3
여러 Trading Service 인스턴스를 같은 Consumer Group으로 묶으면 메시지를 Consumer Group 내 한 Consumer가 처리하도록 구성할 수 있다.
한 인스턴스에 종속된 Spring Event와 달리 서비스 인스턴스 간 작업 분배와 장애 복구를 고려할 수 있는 구조가 된다.
다만 여기서 또 하나의 문제가 생긴다.
8. DB 저장은 성공했는데 Kafka 발행 전에 죽으면?
단순히 Spring Event를 Kafka로 교체한다고 모든 문제가 해결되는 것은 아니다.
1. PENDING DB 저장 성공
2. COMMIT
3. 💥 서버 장애
4. Kafka publish 실행 못함
이 상황에서도 DB에는 PENDING이 있지만 Kafka에는 메시지가 없다.
즉,
DB Transaction
과
Kafka Message Publish
는 서로 다른 작업이기 때문에 둘 사이의 원자성을 어떻게 보장할 것인지라는 새로운 문제가 생긴다.
여기까지 신뢰성을 높이려면 Transactional Outbox Pattern 같은 구조도 고려할 수 있다.
Business Table
p_ai_risk_analyses
+
Outbox Table
AI_ANALYSIS_REQUESTED
│
▼
같은 DB Transaction으로 저장
│
▼
Outbox Publisher
│
▼
Kafka
│
▼
Consumer
분석 데이터와 발행할 이벤트 정보를 같은 DB Transaction으로 저장하고, 이후 Outbox 데이터를 Kafka로 발행하는 방식이다.
여기에 Consumer의 중복 처리까지 고려한다면 멱등성 설계도 필요해진다.
9. 그래서 MVP에서는 어디까지 구현했는가?
이번 프로젝트의 MVP에서는 다음 구조를 선택했다.
Spring Application Event
+
@TransactionalEventListener(AFTER_COMMIT)
+
@Async
+
TaskRejectedException 처리
목적은 LLM 호출로 인해 HTTP 요청이 오래 대기하지 않도록 요청 처리와 AI 분석을 분리하는 것이다.
현재 구조에서도 Thread Pool 포화로 작업 제출 자체가 실패하는 경우에는 분석 상태를 FAILED로 변경해 무한정 PENDING으로 남는 것을 방지했다.
반면 다음과 같은 한계는 남아 있다.
다중 인스턴스 간 이벤트 공유 X
서버 장애 시 인메모리 이벤트 유실 가능
PENDING 작업 복구 기능 X
따라서 서비스 규모가 커지고 AI 분석 요청에 대한 신뢰성 보장이 중요해진다면 다음 단계로 발전시킬 수 있다.
현재
Spring Event + @Async
↓
다중 인스턴스
Kafka + Consumer Group
↓
이벤트 발행 신뢰성
Transactional Outbox
↓
소비 신뢰성
멱등 처리 + Retry + DLT
이번 구현에서 중요한 점은 단순히 @Async를 붙이는 데서 끝나지 않고, 비동기 작업이 실행되기 전 실패할 수 있는 구간과 서버 인스턴스 자체가 사라질 수 있는 구간을 구분하게 되었다는 점이다.
비동기 처리를 도입하면 요청 응답 시간을 줄일 수 있지만, 동시에 작업이 정말 실행되었는지, 실패했다면 상태를 어떻게 복구할 것인지, 서버가 여러 대가 되었을 때 누가 작업을 처리할 것인지까지 함께 설계해야 한다.
'Project > 주식 자동매매 프로젝트' 카테고리의 다른 글
| [동시성] 동시 요청에서 ACTIVE 데이터는 어떻게 하나만 보장할까? (0) | 2026.09.04 |
|---|