Skip to content

Commit a1fc0bc

Browse files
authored
feat: 영상 분석 큐 기반 처리 및 장애 복구 메커니즘 도입 (#61)
* feat: VideoAnalysisQueuePort 큐 포트 및 인메모리 구현체 생성 (#60) * refactor: RateLimitBucketStore 메서드명 개선 및 waitAndConsume 추가 (#60) * feat: VideoAnalysisQueueConsumer 도입 및 EventListener를 enqueue 방식으로 전환 (#60) * feat: YouTube 자막 에러 분류 (429 Rate Limit / 자막 없음 / 알 수 없는 오류) (#60) * feat: VideoAnalysisTaskPersistencePort에 findPendingTasks 추가 (#60) * feat: 앱 시작 시 PENDING 건 큐 재적재 및 큐 소비자 시작 (#60) * feat: 방치된 PENDING 영상 분석 건 재시도 배치 스케줄러 추가 (#60) * test: 큐 소비자 테스트 및 EventListener 테스트 수정 (#60) * refactor: ktlintformat (#58) * fix: 누락 /api prefix 추가 (#58) * refactor: PlaceEnrichUseCase 도입으로 배치 모듈 의존성 역전 (#60) * fix: 영상 분석 큐 정확성 개선 (#60) - Gemini 토큰이 자막 추출 실패 시에도 소모되던 문제 수정 VideoAnalyzePort 를 extractTranscript / analyzeFromTranscript 로 분리하여 Gemini 호출 직전에만 토큰 차감 - 재시도 배치가 처리 중인 작업을 중복 재투입하던 문제 수정 PROCESSING 상태 도입으로 배치의 findPendingTasks 에서 자동 제외 * setting: MySQL 자격증명 GitHub Secrets 분리 및 CI/CD 통합 (#60) - docker-compose.mysql.yml 의 하드코딩된 비밀번호를 환경변수로 변경 - CI/CD 에서 mysql compose 를 함께 배포하되 이미 실행 중이면 스킵 - .env 생성 시 MYSQL_ROOT_PASSWORD 포함 - application-prod.yml 의 password 제거 (SPRING_DATASOURCE_PASSWORD 환경변수로 자동 주입) * fix: 영상 분석 가격 추정 정확도 개선 및 로그 간소화 (#60) * feat: YouTube 자막 추출 프록시 지원 추가 (#60) * fix: Java HttpClient 프록시 터널링 Basic 인증 허용 (#60) * fix: 키워드 수집 후 바로 분석하도록 변경 (#60) * setting: 캄보디아, 짐바브웨, 러시아 키워드 제거 (#60) * refactor: VideoQueue Caffeine 구현체 모듈로 이동 (#60) * fix: 영상 분석 재시도 정책 및 자막 오류 분류 개선 (#60)
1 parent 56426ea commit a1fc0bc

38 files changed

Lines changed: 1091 additions & 580 deletions

File tree

.github/workflows/cicd-release.yml

Lines changed: 24 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -64,24 +64,27 @@ jobs:
6464
host: ${{ secrets.SERVER_HOST }}
6565
username: ec2-user
6666
key: ${{ secrets.SERVER_KEY }}
67-
source: "docker/docker-compose.prod.yml,gcp-credentials.json"
67+
source: "docker/docker-compose.mysql.yml,docker/docker-compose.prod.yml,gcp-credentials.json"
6868
target: "~/linktrip/"
6969
strip_components: 0
7070
overwrite: true
7171

7272
- name: EC2 파일 정리 및 배포
7373
uses: appleboy/ssh-action@master
7474
env:
75+
MYSQL_ROOT_PASSWORD: ${{ secrets.MYSQL_ROOT_PASSWORD }}
7576
JWT_SECRET_KEY: ${{ secrets.JWT_SECRET_KEY }}
7677
DISCORD_WEBHOOK_ERROR_URL: ${{ secrets.DISCORD_WEBHOOK_ERROR_URL }}
7778
DISCORD_MENTION_USER_ID: ${{ secrets.DISCORD_MENTION_USER_ID }}
7879
GCP_PROJECT_ID: ${{ secrets.GCP_PROJECT_ID }}
7980
YOUTUBE_API_KEY: ${{ secrets.YOUTUBE_API_KEY }}
81+
YOUTUBE_PROXY_USERNAME: ${{ secrets.YOUTUBE_PROXY_USERNAME }}
82+
YOUTUBE_PROXY_PASSWORD: ${{ secrets.YOUTUBE_PROXY_PASSWORD }}
8083
with:
8184
host: ${{ secrets.SERVER_HOST }}
8285
username: ec2-user
8386
key: ${{ secrets.SERVER_KEY }}
84-
envs: JWT_SECRET_KEY,DISCORD_WEBHOOK_ERROR_URL,DISCORD_MENTION_USER_ID,GCP_PROJECT_ID,YOUTUBE_API_KEY
87+
envs: MYSQL_ROOT_PASSWORD,JWT_SECRET_KEY,DISCORD_WEBHOOK_ERROR_URL,DISCORD_MENTION_USER_ID,GCP_PROJECT_ID,YOUTUBE_API_KEY,YOUTUBE_PROXY_USERNAME,YOUTUBE_PROXY_PASSWORD
8588
script: |
8689
set -e
8790
@@ -99,33 +102,46 @@ jobs:
99102
100103
# SCP로 전송된 파일 정리 (strip_components=0이라 docker/ 폴더 안에 들어감)
101104
mkdir -p $DEPLOY_DIR
105+
mv -f $DEPLOY_DIR/docker/docker-compose.mysql.yml $DEPLOY_DIR/docker-compose.mysql.yml 2>/dev/null || true
102106
mv -f $DEPLOY_DIR/docker/docker-compose.prod.yml $DEPLOY_DIR/docker-compose.prod.yml 2>/dev/null || true
103107
rmdir $DEPLOY_DIR/docker 2>/dev/null || true
104108
chmod 600 $DEPLOY_DIR/gcp-credentials.json
105109
106110
# .env 파일 생성
107-
echo "JWT_SECRET_KEY=${JWT_SECRET_KEY}" > $DEPLOY_DIR/.env
111+
echo "MYSQL_ROOT_PASSWORD=${MYSQL_ROOT_PASSWORD}" > $DEPLOY_DIR/.env
112+
echo "JWT_SECRET_KEY=${JWT_SECRET_KEY}" >> $DEPLOY_DIR/.env
108113
echo "DISCORD_WEBHOOK_ERROR_URL=${DISCORD_WEBHOOK_ERROR_URL}" >> $DEPLOY_DIR/.env
109114
echo "DISCORD_MENTION_USER_ID=${DISCORD_MENTION_USER_ID}" >> $DEPLOY_DIR/.env
110115
echo "GCP_PROJECT_ID=${GCP_PROJECT_ID}" >> $DEPLOY_DIR/.env
111116
echo "YOUTUBE_API_KEY=${YOUTUBE_API_KEY}" >> $DEPLOY_DIR/.env
117+
echo "YOUTUBE_PROXY_USERNAME=${YOUTUBE_PROXY_USERNAME}" >> $DEPLOY_DIR/.env
118+
echo "YOUTUBE_PROXY_PASSWORD=${YOUTUBE_PROXY_PASSWORD}" >> $DEPLOY_DIR/.env
112119
chmod 600 $DEPLOY_DIR/.env
113120
114-
# 네트워크 생성 (없으면)
121+
cd $DEPLOY_DIR
122+
123+
# MySQL: 이미 실행 중이면 skip (데이터 볼륨 보존)
124+
if docker ps --format '{{.Names}}' | grep -q '^linktrip-mysql$'; then
125+
echo ">>> linktrip-mysql 이미 실행 중 → 스킵"
126+
else
127+
echo ">>> linktrip-mysql 신규 기동"
128+
docker compose -f docker-compose.mysql.yml up -d
129+
fi
130+
131+
# 네트워크 보장 (없으면 생성)
115132
docker network create linktrip-network 2>/dev/null || true
116133
117-
# 기존 컨테이너 정리 (docker run으로 만든 잔여 컨테이너)
134+
# 기존 app 컨테이너 정리 (docker run으로 만든 잔여 컨테이너)
118135
docker stop linktrip-app 2>/dev/null || true
119136
docker rm linktrip-app 2>/dev/null || true
120137
121-
# docker-compose 배포
122-
cd $DEPLOY_DIR
138+
# App: 매 배포마다 재기동
123139
docker compose -f docker-compose.prod.yml pull app
124140
docker compose -f docker-compose.prod.yml up -d --force-recreate
125141
docker image prune -f
126142
127143
echo ">>> 배포 완료"
128-
docker compose -f docker-compose.prod.yml ps
144+
docker ps
129145
130146
- name: 디스코드 배포 알림
131147
if: success()

docker/docker-compose.mysql.yml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ services:
66
environment:
77
- TZ=Asia/Seoul
88
- MYSQL_DATABASE=linktrip
9-
- MYSQL_ROOT_PASSWORD=12345678
9+
- MYSQL_ROOT_PASSWORD=${MYSQL_ROOT_PASSWORD}
1010
command: >
1111
--character-set-server=utf8mb4
1212
--collation-server=utf8mb4_0900_ai_ci

docker/docker-compose.prod.yml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ services:
99
environment:
1010
- SPRING_PROFILES_ACTIVE=prod
1111
- TZ=Asia/Seoul
12+
- JAVA_OPTS=-Djdk.http.auth.tunneling.disabledSchemes=
13+
- SPRING_DATASOURCE_PASSWORD=${MYSQL_ROOT_PASSWORD}
1214
- JWT_SECRET_KEY=${JWT_SECRET_KEY}
1315
- JWT_EXPIRATION_MS=${JWT_EXPIRATION_MS:-86400000}
1416
- DISCORD_WEBHOOK_ERROR_URL=${DISCORD_WEBHOOK_ERROR_URL}
@@ -17,6 +19,8 @@ services:
1719
- GCP_CREDENTIALS_PATH=/app/config/gcp-credentials.json
1820
- GCP_VERTEX_AI_LOCATION=us-central1
1921
- YOUTUBE_API_KEY=${YOUTUBE_API_KEY}
22+
- YOUTUBE_PROXY_USERNAME=${YOUTUBE_PROXY_USERNAME}
23+
- YOUTUBE_PROXY_PASSWORD=${YOUTUBE_PROXY_PASSWORD}
2024
volumes:
2125
- ./gcp-credentials.json:/app/config/gcp-credentials.json:ro
2226
- /etc/localtime:/etc/localtime:ro

linktrip-application/src/main/kotlin/com/linktrip/application/domain/video/PlaceEnrichService.kt

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
package com.linktrip.application.domain.video
22

3+
import com.linktrip.application.port.input.PlaceEnrichUseCase
34
import com.linktrip.application.port.output.external.GooglePlacesPort
45
import com.linktrip.application.port.output.persistence.PlaceEnrichPersistencePort
56
import com.linktrip.application.port.output.persistence.TravelItineraryItemPersistencePort
@@ -20,10 +21,10 @@ class PlaceEnrichService(
2021
private val travelItineraryItemPersistencePort: TravelItineraryItemPersistencePort,
2122
private val videoAnalysisTaskPersistencePort: VideoAnalysisTaskPersistencePort,
2223
private val placeEnrichDispatcher: CoroutineDispatcher,
23-
) {
24-
fun enrichPlaces(
24+
) : PlaceEnrichUseCase {
25+
override fun enrichPlaces(
2526
videoAnalysisTaskId: String,
26-
destination: String? = null,
27+
destination: String?,
2728
) {
2829
val items = travelItineraryItemPersistencePort.findRetryableItems(videoAnalysisTaskId)
2930

@@ -60,7 +61,7 @@ class PlaceEnrichService(
6061
PlaceEnrichResult(itemId = item.id, place = null, success = false)
6162
}
6263

63-
fun retryAll() {
64+
override fun retryAll() {
6465
val videoAnalysisTaskIds = travelItineraryItemPersistencePort.findVideoAnalysisTaskIdsWithRetryableItems()
6566

6667
if (videoAnalysisTaskIds.isEmpty()) {
Lines changed: 213 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,213 @@
1+
package com.linktrip.application.domain.video
2+
3+
import com.linktrip.application.domain.trip.TripPlanService
4+
import com.linktrip.application.port.output.external.VideoAnalysisNotificationPort
5+
import com.linktrip.application.port.output.external.VideoAnalyzePort
6+
import com.linktrip.application.port.output.persistence.TripPlanRequestPersistencePort
7+
import com.linktrip.application.port.output.persistence.VideoAnalysisTaskPersistencePort
8+
import com.linktrip.application.port.output.queue.VideoAnalysisQueuePort
9+
import com.linktrip.application.port.output.ratelimit.RateLimitBucketStore
10+
import com.linktrip.application.port.output.ratelimit.RateLimitPolicy
11+
import com.linktrip.common.exception.ExceptionCode
12+
import com.linktrip.common.exception.LinktripException
13+
import mu.KotlinLogging
14+
import org.springframework.stereotype.Component
15+
16+
private val logger = KotlinLogging.logger {}
17+
18+
@Component
19+
class VideoAnalysisQueueConsumer(
20+
private val videoAnalysisQueuePort: VideoAnalysisQueuePort,
21+
private val videoAnalyzePort: VideoAnalyzePort,
22+
private val videoAnalysisResultSaver: VideoAnalysisResultSaver,
23+
private val placeEnrichService: PlaceEnrichService,
24+
private val videoAnalysisTaskPersistencePort: VideoAnalysisTaskPersistencePort,
25+
private val videoAnalysisNotificationPort: VideoAnalysisNotificationPort,
26+
private val tripPlanRequestPort: TripPlanRequestPersistencePort,
27+
private val tripPlanService: TripPlanService,
28+
private val rateLimitBucketStore: RateLimitBucketStore,
29+
) {
30+
fun startConsuming() {
31+
Thread({ consumeLoop() }, "VideoAnalysisQueueConsumer").apply {
32+
isDaemon = true
33+
start()
34+
}
35+
logger.info { "영상 분석 큐 소비자 시작" }
36+
}
37+
38+
private fun consumeLoop() {
39+
while (!Thread.currentThread().isInterrupted) {
40+
try {
41+
val event = videoAnalysisQueuePort.dequeue() ?: continue
42+
processAnalysis(event)
43+
} catch (_: InterruptedException) {
44+
Thread.currentThread().interrupt()
45+
logger.info { "영상 분석 큐 소비자 중단" }
46+
break
47+
} catch (e: Exception) {
48+
logger.error(e) { "영상 분석 큐 소비자 오류" }
49+
}
50+
}
51+
}
52+
53+
private fun processAnalysis(event: VideoAnalyzeEvent) {
54+
val startTime = System.currentTimeMillis()
55+
logger.info { "영상 분석 시작: id=${event.videoAnalysisTaskId}, url=${event.youtubeUrl}" }
56+
57+
// 처리 중 상태로 전환하여 재시도 배치가 중복 재투입하지 못하게 한다.
58+
videoAnalysisTaskPersistencePort.updateStatus(
59+
event.videoAnalysisTaskId,
60+
VideoAnalysisTaskStatus.PROCESSING,
61+
)
62+
63+
val destination: String?
64+
val title: String?
65+
66+
try {
67+
// 1단계: 자막 추출 (YouTube). 이 단계 실패는 Gemini 토큰을 소모하지 않는다.
68+
val transcript = videoAnalyzePort.extractTranscript(event.youtubeUrl)
69+
70+
// 2단계: Gemini 분석. 실제 Gemini 호출 직전에만 토큰을 차감한다.
71+
rateLimitBucketStore.waitAndConsume(RATE_LIMIT_KEY, RateLimitPolicy.GEMINI_API)
72+
val result = videoAnalyzePort.analyzeFromTranscript(transcript, event.youtubeUrl)
73+
74+
if (!result.valid) {
75+
logger.warn { "유효하지 않은 영상: id=${event.videoAnalysisTaskId}" }
76+
videoAnalysisTaskPersistencePort.updateValidAndStatus(
77+
event.videoAnalysisTaskId,
78+
valid = false,
79+
VideoAnalysisTaskStatus.INVALID,
80+
)
81+
return
82+
}
83+
84+
destination = result.destination
85+
title = result.title
86+
val itineraryItems = toItineraryItems(event.videoAnalysisTaskId, result)
87+
val timelines = toTimelines(event.videoAnalysisTaskId, result)
88+
videoAnalysisResultSaver.save(
89+
event.videoAnalysisTaskId,
90+
itineraryItems,
91+
summary = result.summary,
92+
estimatedMinCost = result.estimatedMinCost,
93+
estimatedMaxCost = result.estimatedMaxCost,
94+
costBasis = result.costBasis,
95+
hashtags = result.hashtags,
96+
timelines = timelines,
97+
destination = result.destination,
98+
)
99+
100+
val elapsed = System.currentTimeMillis() - startTime
101+
logger.info {
102+
"영상 분석 완료: id=${event.videoAnalysisTaskId}, destination=$destination, " +
103+
"${elapsed}ms, items=${itineraryItems.size}"
104+
}
105+
} catch (e: LinktripException) {
106+
when (e.exceptionCode) {
107+
ExceptionCode.BAD_GATEWAY_YOUTUBE -> {
108+
logger.warn { "YouTube Rate Limit으로 재시도 대기: id=${event.videoAnalysisTaskId}" }
109+
// 즉시 재큐하지 않고, 재시도 배치가 다음 주기에 다시 적재하도록 한다.
110+
videoAnalysisTaskPersistencePort.updateStatus(
111+
event.videoAnalysisTaskId,
112+
VideoAnalysisTaskStatus.PENDING,
113+
)
114+
}
115+
ExceptionCode.BAD_REQUEST_VIDEO -> {
116+
logger.warn { "자막 없는 영상 INVALID 처리: id=${event.videoAnalysisTaskId}" }
117+
videoAnalysisTaskPersistencePort.updateValidAndStatus(
118+
event.videoAnalysisTaskId,
119+
valid = false,
120+
VideoAnalysisTaskStatus.INVALID,
121+
)
122+
}
123+
else -> {
124+
logger.error(e) { "영상 분석 실패: id=${event.videoAnalysisTaskId}" }
125+
videoAnalysisTaskPersistencePort.updateStatus(
126+
event.videoAnalysisTaskId,
127+
VideoAnalysisTaskStatus.FAILED,
128+
)
129+
}
130+
}
131+
return
132+
} catch (e: Exception) {
133+
logger.error(e) { "영상 분석 실패: id=${event.videoAnalysisTaskId}" }
134+
videoAnalysisTaskPersistencePort.updateStatus(
135+
event.videoAnalysisTaskId,
136+
VideoAnalysisTaskStatus.FAILED,
137+
)
138+
return
139+
}
140+
141+
processPendingRequests(event.videoAnalysisTaskId, title ?: destination)
142+
enrichPlaces(event.videoAnalysisTaskId, destination)
143+
144+
val memberIds = tripPlanRequestPort.findMemberIdsByVideoAnalysisTaskId(event.videoAnalysisTaskId)
145+
videoAnalysisNotificationPort.notifyAnalysisComplete(event.videoAnalysisTaskId, memberIds)
146+
}
147+
148+
private fun processPendingRequests(
149+
videoAnalysisTaskId: String,
150+
destination: String?,
151+
) {
152+
val requests = tripPlanRequestPort.findUnprocessedByVideoAnalysisTaskId(videoAnalysisTaskId)
153+
if (requests.isEmpty()) return
154+
155+
val planTitle = destination ?: "여행 계획"
156+
requests.forEach { request ->
157+
runCatching {
158+
tripPlanService.createFromAnalysisIfAbsent(request.memberId, videoAnalysisTaskId, planTitle)
159+
request.markProcessed()
160+
}.onFailure { e ->
161+
logger.error(e) {
162+
"여행 계획 자동 생성 실패: taskId=$videoAnalysisTaskId, requestId=${request.id}"
163+
}
164+
}
165+
}
166+
tripPlanRequestPort.saveAll(requests)
167+
168+
logger.info {
169+
"여행 계획 자동 생성: taskId=$videoAnalysisTaskId, " +
170+
"성공=${requests.count { it.processed }}/${requests.size}"
171+
}
172+
}
173+
174+
private fun enrichPlaces(
175+
videoAnalysisTaskId: String,
176+
destination: String?,
177+
) {
178+
try {
179+
val startTime = System.currentTimeMillis()
180+
placeEnrichService.enrichPlaces(videoAnalysisTaskId, destination)
181+
val elapsed = System.currentTimeMillis() - startTime
182+
logger.info { "장소 보강 소요시간: id=$videoAnalysisTaskId, ${elapsed}ms" }
183+
} catch (e: Exception) {
184+
logger.warn(e) { "장소 보강 실패, 분석 결과는 유지: id=$videoAnalysisTaskId" }
185+
}
186+
}
187+
188+
private fun toItineraryItems(
189+
videoAnalysisTaskId: String,
190+
result: VideoAnalysisResult,
191+
): List<TravelItineraryItem> =
192+
result.days.flatMap { daySchedule ->
193+
daySchedule.items.map { item ->
194+
TravelItineraryItem.from(videoAnalysisTaskId, daySchedule, item)
195+
}
196+
}
197+
198+
private fun toTimelines(
199+
videoAnalysisTaskId: String,
200+
result: VideoAnalysisResult,
201+
): List<VideoTimeline> =
202+
result.timeline.map { timelineItem ->
203+
VideoTimeline.create(
204+
videoAnalysisTaskId = videoAnalysisTaskId,
205+
timestampSeconds = timelineItem.timestampSeconds,
206+
description = timelineItem.description,
207+
)
208+
}
209+
210+
companion object {
211+
private const val RATE_LIMIT_KEY = "gemini-api"
212+
}
213+
}

linktrip-application/src/main/kotlin/com/linktrip/application/domain/video/VideoAnalysisTaskStatus.kt

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package com.linktrip.application.domain.video
22

33
enum class VideoAnalysisTaskStatus {
44
PENDING,
5+
PROCESSING,
56
COMPLETED,
67
INVALID,
78
FAILED,

0 commit comments

Comments
 (0)