diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java index 3fa0d212..88b9f95e 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/config/RagExecutionConfig.java @@ -13,20 +13,18 @@ /** * RAG 답변 생성을 동시에 몇 건까지 처리할지 제한하는 실행 자원을 구성한다 (#340). * - *
{@code domain/worker/config/WorkerExecutionConfig}(이 코드베이스에서 유일했던 커스텀 - * 스레드풀 선례)를 그대로 본떴다 — {@code core=max}인 {@link ThreadPoolExecutor} + - * {@link SynchronousQueue}(큐잉 없음) + {@link ThreadPoolExecutor.AbortPolicy}(꽉 차면 즉시 - * 거부, 조용히 쌓아두지 않음)로 "설정된 동시성만 즉시 실행"을 보장한다. + *
Ollama 호출은 한 건에 수십 초가 걸리고 그동안 호출한 스레드는 응답을 기다리며 멈춘다. + * 동시에 N건을 처리하려면 그렇게 기다릴 스레드가 N개 있어야 한다. 이 클래스는 그 스레드 N개 + * ({@link #ragWorkerJobExecutor})와, 그중 몇 개가 놀고 있는지 세는 카운터 + * ({@link #ragWorkerSlots})를 만든다. 둘 다 {@code rag.worker.max-concurrency} 하나에서 + * 크기를 받는다. 실제로 job을 집고 처리하는 흐름은 {@code RagJobWorker}에 있다. * - *
{@code embedding_jobs}의 {@code WorkerExecutionSlotPool}(별도 클래스, 종료 플래그 + - * introspection 메서드 포함)까지는 필요 없다 — RAG는 별도 워커 등록/우아한 종료 조율이나 - * 대시보드 노출 요구가 없어서, 같은 안전 성질(로컬 슬롯을 먼저 확보한 뒤에만 DB claim을 - * 시도해 "claim은 됐는데 실행할 스레드가 없는" 상태를 만들지 않는 것)을 순수 - * {@link Semaphore}만으로 재현한다. + *
{@code domain/worker/config/WorkerExecutionConfig}(이 코드베이스에서 유일했던 커스텀 + * 스레드풀 선례)를 본떴다. 다만 {@code WorkerExecutionSlotPool}(종료 플래그 + introspection 포함) + * 까지는 필요 없어 순수 {@link Semaphore}로 줄였다. * - *
{@code indexing.worker.enabled}로 켜고 끌 수 있는 인덱싱 워커와 달리, RAG Worker는 - * {@code RagSchedulingConfig}와 동일하게 조건 없이 항상 켜져 있어야 하는 검색 API 핵심 - * 경로라 {@code @ConditionalOnProperty}를 붙이지 않는다. + *
{@code indexing.worker.enabled}로 켜고 끌 수 있는 인덱싱 워커와 달리, RAG Worker는 검색 + * API의 핵심 경로라 {@code @ConditionalOnProperty}를 붙이지 않는다. */ @Configuration public class RagExecutionConfig { @@ -35,7 +33,19 @@ public class RagExecutionConfig { public static final String RAG_WORKER_SLOTS = "ragWorkerSlots"; /** - * 동시에 최대 {@code rag.worker.max-concurrency}건까지만 즉시 실행하는 무대기 Executor를 만든다. + * 실제로 Ollama를 호출하고 기다리는 RAG 전용 스레드 N개. + * + *
{@code RagJobWorker.processNext()}는 job을 여기에 {@code execute()}로 던지고 즉시 돌아온다. + * 그래서 폴링하는 스케줄러 스레드는 Ollama 호출에 막히지 않는다. + * + *
생성자 인자: core=max=N이라 스레드는 정확히 N개로 고정된다. {@link SynchronousQueue}는 + * 대기줄이 없다는 뜻 — 스레드가 다 바쁠 때 job을 던지면 풀 안에 쌓아두지 않고 즉시 거부한다. + * 이미 DB에 "처리 중"이라고 적힌 job이 풀 안에서 몰래 대기하면 스위퍼가 그걸 모르고 강제 + * 종료하므로, 대기는 DB 테이블에서만 한다. {@link ThreadPoolExecutor.AbortPolicy}는 그 거부를 + * {@code RejectedExecutionException}으로 알린다 — 조용히 버리거나(DiscardPolicy) 폴러가 직접 + * 실행하는(CallerRunsPolicy) 것보다 예외를 받아 슬롯을 돌려주는 편이 맞다. + * {@link CustomizableThreadFactory}는 로그에서 구분되게 스레드 이름을 {@code rag-worker-job-N} + * 으로 붙인다. */ @Bean(name = RAG_WORKER_JOB_EXECUTOR, destroyMethod = "shutdownNow") public ThreadPoolExecutor ragWorkerJobExecutor( @@ -54,9 +64,21 @@ public ThreadPoolExecutor ragWorkerJobExecutor( } /** - * DB claim을 시도하기 전에 먼저 확보해야 하는 로컬 실행 슬롯. 공정 모드(fair)로 만들어 - * 폴링 주기가 겹칠 때 대기가 한쪽으로 몰리지 않게 한다 — {@code WorkerExecutionSlotPool}의 - * 선택과 동일하다. + * {@link #ragWorkerJobExecutor}의 스레드 N개 중 몇 개가 놀고 있는지 세는 카운터. DB에서 job을 + * 꺼내기(claim) 전에 이 값을 보고, 0이면 DB를 건드리지 않는다. + * + *
왜 필요한가: claim은 DB에 {@code claimed_at}을 적는 행위라, 꺼낸 뒤 Executor가 거부하면 + * "DB엔 처리 중인데 실제론 아무도 안 하는" 유령 job이 남는다 — 다음 폴링은 + * {@code claimed_at IS NULL}만 찾으니 다시 집지도 않는다. 그런데 {@link ThreadPoolExecutor}는 + * "지금 노는 스레드 있어?"를 믿을 만하게 물어볼 방법이 없고 {@code execute()}를 던져 봐야 안다. + * 그래서 노는 스레드 수를 직접 세는 카운터를 옆에 두고, "슬롯 확보 → DB claim → execute" + * 순서를 강제한다. permit 수와 스레드 수가 같은 설정값에서 나오므로 슬롯이 남았으면 노는 + * 스레드도 반드시 있다. + * + *
{@link Semaphore}는 OS의 그 세마포어다. {@code tryAcquire()}가 P(0이면 블로킹 대신 즉시 + * false), {@code release()}가 V. P는 폴러 스레드가, V는 수십 초 뒤 워커 스레드가 + * {@code finally}에서 부른다 — 잡은 쪽과 푸는 쪽이 달라도 되는 게 락이 아니라 세마포어인 + * 이유다. 공정 모드({@code true})는 {@code tryAcquire()}에는 효과가 없고 선례를 따라 둔 것이다. */ @Bean(name = RAG_WORKER_SLOTS) public Semaphore ragWorkerSlots(@Value("${rag.worker.max-concurrency:2}") int maxConcurrency) { diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/controller/RagWebSocketController.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/controller/RagWebSocketController.java index 4423f1cf..1b96678b 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/rag/controller/RagWebSocketController.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/controller/RagWebSocketController.java @@ -31,8 +31,10 @@ public void notifyAnswerReady(String userEmail, Long queryId) { messagingTemplate.convertAndSendToUser(userEmail, RAG_ANSWER_QUEUE, new RagAnswerReadyEvent(queryId)); } - // 완료 알림의 최소 트리거 페이로드 — 답변 본문은 담지 않는다. 프론트가 이 이벤트를 받으면 - // 항상 GET /search/{queryId}로 다시 조회해야 하며, 이 record 자체를 최종 상태로 신뢰하면 안 된다. + /** + * 완료 알림의 최소 트리거 페이로드 — 답변 본문은 담지 않는다. 프론트가 이 이벤트를 받으면 + * 항상 GET /search/{queryId}로 다시 조회해야 하며, 이 record 자체를 최종 상태로 신뢰하면 안 된다. + */ private record RagAnswerReadyEvent(Long queryId) { } } diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/dto/RagAnswer.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/dto/RagAnswer.java index 62486fb5..2783acdd 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/rag/dto/RagAnswer.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/dto/RagAnswer.java @@ -10,7 +10,8 @@ *
비동기 Job 큐 전환(#218) 이전에는 {@code RagFacade.generate()}가 성공 경로에서 후보 목록을 * citation으로 변환해 이 타입을 만들었다. #218 이후 {@code RagFacade.processJob()}은 성공/실패 * 결과를 {@link com.opensource.docgrid.domain.rag.entity.RagResponse}에 직접 영속화하고 - * {@code void}를 반환하도록 바뀌어, 그 변환 로직은 더 이상 필요 없어져 제거됐다. 지금 유일하게 + * "실제로 확정이 일어났는지"만 {@code boolean}으로 반환하도록 바뀌어(#288), 그 변환 로직은 + * 더 이상 필요 없어져 제거됐다. 지금 유일하게 * 쓰이는 경로는 {@link #noContext} — {@code RagFacade.enqueue()}가 검색 후보 0건(NO_CONTEXT)일 때 * LLM 호출 없이 즉시 만드는 응답이다. */ diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java index 3c6498c1..383551b8 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/entity/RagResponse.java @@ -57,9 +57,11 @@ public class RagResponse extends BaseEntity { @JoinColumn(name = "query_id", nullable = false) private SearchQuery query; - // PROCESSING 상태로 처음 저장될 때는 아직 값이 없다 — Worker의 정상 완료(completeSuccess/ - // completeFailed) 또는 RagJobTimeoutSweeper의 강제 종료(forceFailIfProcessing) 중 먼저 - // 확정되는 쪽이 채운다. + /** + * PROCESSING 상태로 처음 저장될 때는 아직 값이 없다 — Worker의 정상 완료(completeSuccess/ + * completeFailed) 또는 RagJobTimeoutSweeper의 강제 종료(forceFailIfProcessing) 중 먼저 + * 확정되는 쪽이 채운다. + */ @Column(name = "answer_text", columnDefinition = "TEXT") private String answerText; @@ -90,9 +92,11 @@ public class RagResponse extends BaseEntity { @Column(name = "error_message", columnDefinition = "TEXT") private String errorMessage; - // 병렬 Worker가 이 job을 이미 집었는지 표시한다(#340). status만으로는 "대기 중"과 "누가 이미 - // 처리 중"을 구분할 수 없어서(둘 다 PROCESSING) 별도로 둔다. RagResponseClaimService의 짧은 - // claim 트랜잭션 안에서만 채워지며, 그 밖의 완료 확정 경로(조건부 UPDATE)는 이 컬럼을 건드리지 않는다. + /** + * 병렬 Worker가 이 job을 이미 집었는지 표시한다(#340). status만으로는 "대기 중"과 "누가 이미 + * 처리 중"을 구분할 수 없어서(둘 다 PROCESSING) 별도로 둔다. RagResponseClaimService의 짧은 + * claim 트랜잭션 안에서만 채워지며, 그 밖의 완료 확정 경로(조건부 UPDATE)는 이 컬럼을 건드리지 않는다. + */ @Column(name = "claimed_at") private LocalDateTime claimedAt; diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java index db4defe6..c07c1496 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/repository/RagResponseRepository.java @@ -18,18 +18,52 @@ /** * RagResponse 엔티티에 대한 JPA Repository. + * + *
{@code rag_responses} 테이블은 답변 저장소이면서 동시에 RagJobWorker의 job 큐다 — 별도 큐 + * 테이블 없이 {@code status=PROCESSING} 행이 곧 대기/처리 중인 job이다. 그래서 이 Repository에는 + * 평범한 조회 외에 큐 소비를 위한 쿼리가 섞여 있다. + * + *
동시성 방침: 잠금은 job을 집는 순간({@link #findNextUnclaimedProcessingForUpdate})에만 짧게 + * 걸고, 완료 확정은 잠금 없이 {@code WHERE status = PROCESSING} 조건부 UPDATE + * ({@link #completeSuccessIfProcessing}/{@link #forceFailIfProcessing})로 한다. 이 프로젝트엔 + * {@code @Version}이 없어 엔티티를 불러와 save()하면 나중 쓰기가 무조건 이기는데, SQL의 WHERE절이 + * "먼저 끝난 결과를 덮어쓰지 마라"를 대신한다. + * + *
+ * 메서드 호출자 시점 + * findNextUnclaimedProcessingForUpdate Claim 서비스 폴링마다 (job 집기) + * findWithQueryAndUserById Worker claim 직후 (알림용 email) + * releaseAllClaimsOnStartup Worker 기동 1회 (죽은 claim 복구) + * findByStatusAndCreatedAtBefore Sweeper 15초마다 (90초 넘은 job 탐색) + * forceFailIfProcessing Sweeper+Worker FAILED 확정 + * completeSuccessIfProcessing Worker SUCCESS 확정 + **/ public interface RagResponseRepository extends JpaRepository
이 쿼리 혼자서는 소유권이 안 생긴다 — 잠금은 트랜잭션이 끝나면 풀리므로, + * {@link RagResponseClaimService}가 같은 짧은 트랜잭션 안에서 {@link RagResponse#markClaimed}로 + * {@code claimed_at}을 채우고 커밋해야 그 뒤로도 다른 Worker가 못 집는다. + * + *
왜 잠금 없이 {@code claimed_at}만으로는 부족한가: 두 Worker가 정말 같은 순간에
+ * "1번 job의 claimed_at이 NULL이네"를 각자 확인하면, 그 확인 결과를 써서 UPDATE하려는
+ * 찰나에 둘 다 같은(NULL) 값을 본 상태라 둘 다 자기가 집은 줄 안다 — "확인하고 쓰는" 그
+ * 틈에 경합이 생긴다. {@code FOR UPDATE}는 그 틈에 한 트랜잭션만 행을 보게 만들어 이
+ * 경합을 원천 차단하고, {@code claimed_at}은 잠금이 풀린 뒤의 장기 소유권 표시를 맡는다.
*/
@Query(value = """
SELECT * FROM rag_responses
@@ -104,7 +138,12 @@ List {@code RagJobTimeoutSweeper}뿐 아니라 {@code RagResponseCommandService.completeFailed()}
- * (RagJobWorker가 Ollama 호출 실패를 처리하는 정상 경로)도 이 메서드를 그대로 재사용한다 —
- * 둘 다 "PROCESSING인 job을 FAILED + 문구로 확정한다"는 동일한 SQL이 필요하고, 반대로
- * RagJobTimeoutSweeper가 먼저 이 job을 확정해버렸다면 RagJobWorker 쪽 시도도 똑같이
- * 무시돼야 하기 때문이다(#288).
+ * 호출자는 둘이다 — RagJobTimeoutSweeper의 타임아웃 강제 종료와, Worker의 Ollama 호출 실패
+ * 처리({@code RagResponseCommandService.completeFailed}). SQL 모양이 같아 공유하며, 어느 쪽이
+ * 먼저 끝냈든 나중 쪽은 똑같이 무시돼야 하기 때문이기도 하다(#288).
*
- * {@code clearAutomatically}: 벌크 UPDATE는 영속성 컨텍스트를 거치지 않고 DB에 직접
- * 실행되므로, 같은 트랜잭션에서 이 job 엔티티를 이미 로딩해둔 상태라면 그 캐시된 인스턴스가
- * 여전히 갱신 전 값을 들고 있다 — 이후 같은 트랜잭션에서 다시 조회해도 DB가 아니라 그 캐시를
- * 돌려줘 최신 상태를 못 본다. {@code clearAutomatically = true}로 UPDATE 직후 영속성
- * 컨텍스트를 비워 이 문제를 막는다.
+ * {@code clearAutomatically = true}: 벌크 UPDATE는 영속성 컨텍스트(1차 캐시)를 거치지 않고
+ * DB로 바로 간다. 같은 트랜잭션에서 이 엔티티를 이미 읽어 뒀다면 캐시엔 옛 값이 남고, 그 뒤
+ * {@code findById}는 DB가 아니라 캐시를 돌려줘 갱신 전 값을 본다. UPDATE 직후 캐시를 비워 이를
+ * 막는다(#286 테스트에서 실제로 걸렸던 문제).
*/
@Modifying(clearAutomatically = true)
@Query("UPDATE RagResponse r SET r.status = com.opensource.docgrid.domain.search.enums.ResultStatus.FAILED, "
@@ -139,11 +173,12 @@ int forceFailIfProcessing(@Param("id") Long id, @Param("answerText") String answ
@Param("errorMessage") String errorMessage);
/**
- * PROCESSING 상태인 job을 SUCCESS + 생성 결과로 확정한다. {@link #forceFailIfProcessing}과
- * 대칭되는 목적이다 — RagJobTimeoutSweeper가 이 job을 먼저 FAILED로 강제 종료했다면,
- * RagJobWorker의 뒤늦은 정상 완료 시도가 그 결과를 조건 없이 덮어써버리는 경합(#288)을
- * 막는다. {@code WHERE ... AND status = PROCESSING} 조건 덕분에, 스위퍼가 먼저 확정해
- * 이 UPDATE 시점에 status가 이미 FAILED라면 영향받은 행이 0건이 된다.
+ * 아직 PROCESSING일 때만 SUCCESS + 생성 결과로 확정한다. {@link #forceFailIfProcessing}의 반대
+ * 방향 — 스위퍼가 먼저 FAILED로 끝낸 job을 Worker의 뒤늦은 정상 완료가 덮어쓰는 경합(#288)을
+ * 막는다. 0건이면 호출자(RagFacade.processJob)는 citation 저장도 건너뛴다.
+ *
+ * SET 절의 5개 필드가 완료 시 채우는 전부다 — 옛 {@code markSuccess()}가 하던 일을 SQL이
+ * 대신한다. {@code claimed_at}은 건드리지 않아 완료 뒤에도 "언제 집혔는지"가 남는다.
*/
@Modifying(clearAutomatically = true)
@Query("UPDATE RagResponse r SET r.status = com.opensource.docgrid.domain.search.enums.ResultStatus.SUCCESS, "
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/OllamaClient.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/OllamaClient.java
index 72f5cac0..13b1b784 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/OllamaClient.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/OllamaClient.java
@@ -37,8 +37,10 @@
@Service
public class OllamaClient {
- // eval_count(실제 생성된 토큰 수)가 num_predict에 도달했다는 건 모델이 할 말을 다 못 하고
- // 토큰 상한에 걸려 끊겼다는 확정적 신호다 — LLM이 스스로 이를 감지·보고하게 하는 것보다 신뢰할 수 있다.
+ /**
+ * eval_count(실제 생성된 토큰 수)가 num_predict에 도달했다는 건 모델이 할 말을 다 못 하고
+ * 토큰 상한에 걸려 끊겼다는 확정적 신호다 — LLM이 스스로 이를 감지·보고하게 하는 것보다 신뢰할 수 있다.
+ */
private static final String TRUNCATION_NOTICE =
"\n\n(※ 답변이 길어 일부 내용이 생략됐을 수 있습니다. 자세한 내용은 문서를 확인해주세요.)";
@@ -46,13 +48,17 @@ public class OllamaClient {
private static final ObjectMapper CHUNK_MAPPER = new ObjectMapper()
.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
- // 한국어 RAG 답변에 한자·히라가나·가타카나가 나올 일은 없다. qwen 계열의 code-switching으로
- // 섞여 나온 문자를 프롬프트 지시(모델이 무시할 수 있음)가 아닌 코드로 제거한다.
+ /**
+ * 한국어 RAG 답변에 한자·히라가나·가타카나가 나올 일은 없다. qwen 계열의 code-switching으로
+ * 섞여 나온 문자를 프롬프트 지시(모델이 무시할 수 있음)가 아닌 코드로 제거한다.
+ */
private static final Pattern FOREIGN_CJK_PATTERN =
Pattern.compile("[\\p{IsHan}\\p{IsHiragana}\\p{IsKatakana}]+");
- // 혼입이 이 글자 수를 넘으면 낱자 노이즈가 아니라 모델이 중국어로 넘어가 무너진 구간으로 판단하고,
- // 문자만 지워 구두점 뼈대를 남기는 대신 혼입 시작 지점에서 답변을 자른다.
+ /**
+ * 혼입이 이 글자 수를 넘으면 낱자 노이즈가 아니라 모델이 중국어로 넘어가 무너진 구간으로 판단하고,
+ * 문자만 지워 구두점 뼈대를 남기는 대신 혼입 시작 지점에서 답변을 자른다.
+ */
private static final int FOREIGN_CJK_CUT_THRESHOLD = 8;
private final String model;
@@ -134,9 +140,11 @@ public OllamaGenerateResult generate(String prompt) {
throw new DocGridException(ErrorCode.RAG_SERVICE_UNAVAILABLE);
}
- // done:true 없이 스트림이 끝나는 경우가 있다: 데드라인 조기 종료 외에도, Ollama의 PEG 파서가
- // 한글이 토큰 경계에서 바이트 단위로 쪼개진 출력을 파싱하지 못하고 생성을 취소하는 버그
- // (llama.cpp #24807)가 확인됐다. 발생 빈도를 추적할 수 있게 경고 로그를 남긴다.
+ /**
+ * done:true 없이 스트림이 끝나는 경우가 있다: 데드라인 조기 종료 외에도, Ollama의 PEG 파서가
+ * 한글이 토큰 경계에서 바이트 단위로 쪼개진 출력을 파싱하지 못하고 생성을 취소하는 버그
+ * (llama.cpp #24807)가 확인됐다. 발생 빈도를 추적할 수 있게 경고 로그를 남긴다.
+ */
boolean prematureEnd = !chunks.last().done();
if (prematureEnd && !chunks.deadlineExceeded()) {
log.warn("Ollama 스트림이 done 없이 조기 종료됨(서버 측 생성 취소 추정): 수신 텍스트 길이={}", chunks.answer().length());
@@ -191,8 +199,10 @@ private StreamChunks readStream(InputStream body, long deadline) throws IOExcept
}
}
} catch (IOException e) {
- // 스트림이 멈춰 read-timeout이 본문 연결을 끊는 경우 등. 이미 받은 부분 답변이 있으면
- // 버리지 않고 done 없는 조기 종료로 처리해 반환하고, 하나도 없을 때만 실패로 전파한다.
+ /**
+ * 스트림이 멈춰 read-timeout이 본문 연결을 끊는 경우 등. 이미 받은 부분 답변이 있으면
+ * 버리지 않고 done 없는 조기 종료로 처리해 반환하고, 하나도 없을 때만 실패로 전파한다.
+ */
if (answer.isEmpty()) {
throw e;
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagFacade.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagFacade.java
index d5e5f047..6a78dae6 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagFacade.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/RagFacade.java
@@ -36,10 +36,15 @@
* 이 클래스 자체엔 동시성 조율 코드가 없다 — "몇 건을 동시에 처리할지"는 RagJobWorker와
+ * {@code RagExecutionConfig}가 정한다(#340). 대신 RagJobWorker의 스레드 여러 개가 서로 다른
+ * jobId로 {@link #processJob}을 동시에 부를 수 있다는 전제로 짜여 있다 — 인스턴스 필드를
+ * 바꾸는 코드가 없고, 결과가 이미 다른 경로(스위퍼)에 뺏겼으면 boolean으로 조용히 물러난다.
+ *
* SearchFacade와 별도 트랜잭션으로 분리되어 있다(SearchController가 순차 호출) — 검색 DB 작업이
* enqueue()의 짧은 DB 작업과 하나의 커넥션을 오래 물고 있지 않도록 하기 위함이다.
*/
@@ -49,29 +54,28 @@
@Slf4j
public class RagFacade {
- /*
- * LLM_FALLBACK_PREFIX — Ollama 호출 실패 시 최상위 검색 후보 원문을 인용하며 붙이는
- * 안내 문구(buildExtractiveFallbackAnswer 참고).
- * UNEXPECTED_FAILURE_ANSWER_TEXT — processJob() 내부에서 예상 못한 예외(버그 등)로 실패했을
- * 때 쓰는 최소 안내 문구. extractive fallback과 달리 candidates를
- * 다시 불러오지 않는다 — 이미 한 번 예상 밖으로 실패한 상황에서
- * 추가 조회를 시도하다 또 실패할 위험을 만들지 않기 위함이다
- * (RagJobWorker 참고).
- * FALLBACK_EXCERPT_MAX_CODE_POINTS(300) — fallback 문구에 원문을 통째로 붙이면 답변이
- * 지나치게 길어져, 미리보기 수준으로만 잘라 보여준다.
- * MAX_PROMPT_CANDIDATES(3) — topK는 호출자가 1~20까지 자유롭게 요청할 수 있어
- * (SearchRequest), 후보 수를 그대로 프롬프트에 다 넣으면
- * prefill 시간이 예측 불가능해진다(#210). 검색 결과 수(topK)와
- * 별개로 프롬프트와 정상 답변 citation 모두 상위 후보를 이 값까지
- * 제한해, LLM에 제공하지 않은 후보가 출처로 추가되지 않게 한다.
- * NO_RELEVANT_DOC_PHRASE — PromptBuilder가 LLM에게 무관한 문서일 때 이 문구로만 답하도록
- * 지시한다(#65 INSTRUCTION 참고). 검색은 됐지만(candidates 존재)
- * LLM이 무관하다고 판단한 경우, 화면에 근거 문서를 같이 보여주면
- * 안내 문구와 모순돼 보인다.
- * TIMEOUT_ERROR_MESSAGE — RagJobTimeoutSweeper가 너무 오래 PROCESSING으로 남은 job을
- * 강제 종료할 때 error_message에 남기는 문구(#286). Ollama
- * 예외 메시지와 구분해, 나중에 로그/DB로 "진짜 실패"와 "큐
- * 적체로 인한 강제 종료"를 구분할 수 있게 한다.
+ /**
+ * 상수 7개 요약:
+ * 이 메서드 자체는 기다리는 코드가 없다 — claim 한 번은 수 ms 안에 끝나고, 실제 Ollama
+ * 호출(수십 초)은 {@link #executeClaimedJob}으로 넘겨 이 메서드는 곧바로 다음 슬롯을 보러
+ * 돌아간다. {@code tryAcquire()}를 {@code claimNext()}보다 먼저 부르는 순서도 이 때문에
+ * 중요하다 — 순서가 바뀌면 "DB엔 claim됐다고 적혔는데 넘길 스레드가 없는" 유령 job이
+ * 생긴다(클래스 Javadoc 참고).
+ *
* {@code while (tryAcquire())}만으로 반복 횟수 상한이 자동으로 걸린다 — Semaphore의 총
* permit 수가 이미 {@code max-concurrency}와 같아서, 별도 카운터 변수 없이도 이 루프가
* {@code max-concurrency}번보다 더 돌 수 없다.
@@ -114,8 +120,10 @@ public void processNext() {
try {
ragWorkerJobExecutor.execute(() -> executeClaimedJob(jobId));
} catch (RejectedExecutionException e) {
- // 슬롯을 먼저 확보했으므로 이론상 도달하지 않아야 하지만(Executor 정원 =
- // Semaphore 총 permit 수), 종료 절차 중 등 극단적 상황에 대비한 방어다.
+ /**
+ * 슬롯을 먼저 확보했으므로 이론상 도달하지 않아야 하지만(Executor 정원 =
+ * Semaphore 총 permit 수), 종료 절차 중 등 극단적 상황에 대비한 방어다.
+ */
ragWorkerSlots.release();
log.warn("[RAG-WORKER] 실행 제출이 거부됨 jobId={}", jobId);
return;
@@ -144,8 +152,10 @@ private void executeClaimedJob(Long jobId) {
try {
RagResponse job = ragResponseRepository.findWithQueryAndUserById(jobId).orElse(null);
if (job == null) {
- // claim 직후 이 job이 통째로 사라지는 건 극단적 상황(예: 테스트 데이터 정리)에서만
- // 가능하다 — processJob() 자신도 findById로 다시 조회하므로 여기서는 방어만 한다.
+ /**
+ * claim 직후 이 job이 통째로 사라지는 건 극단적 상황(예: 테스트 데이터 정리)에서만
+ * 가능하다 — processJob() 자신도 findById로 다시 조회하므로 여기서는 방어만 한다.
+ */
log.error("[RAG-WORKER] claim된 job을 찾을 수 없음 jobId={}", jobId);
return;
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/command/RagResponseClaimService.java b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/command/RagResponseClaimService.java
index c76c50ea..9bc6d988 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/rag/service/command/RagResponseClaimService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/rag/service/command/RagResponseClaimService.java
@@ -23,6 +23,16 @@
* RAG는 PENDING 같은 별도 대기 상태가 없어(생성 즉시 PROCESSING) {@code embedding_jobs}처럼
* status 전이로 claim을 표시할 수 없다 — 대신 {@code claimed_at} 컬럼을 그 신호로 쓴다.
*
+ * "job을 고르는 단계"와 "실제로 처리하는 단계"는 시간상 완전히 분리된다 — 이 클래스는
+ * 전자(수 ms)만 담당한다:
+ * {@code embedding_jobs}와 달리 별도의 분산 Worker 등록/생존 검증은 하지 않는다 — RAG
* Worker는 이 프로세스 안의 로컬 스레드일 뿐이라 그런 개념 자체가 없다.
*/
@@ -35,7 +45,9 @@ public class RagResponseClaimService {
private final Clock clock;
/**
- * 다음으로 처리할 PROCESSING job 하나를 claim한다.
+ * 다음으로 처리할 PROCESSING job 하나를 claim한다. 클래스 Javadoc의 "짧은 트랜잭션"이
+ * 정확히 이 메서드 호출 하나다 — 시작할 때 잠금이 걸리고, 리턴과 함께 커밋되며 claim이
+ * 확정되고 잠금이 풀린다.
*
* @return claim에 성공한 job의 id. 대기 중인 job이 없으면 빈 값.
*/
* 1. enqueue() — SearchController가 검색 직후 동기 호출. 프롬프트만 조립해 PROCESSING으로 저장하고
* 즉시 반환한다(LLM 호출 없음). candidates가 비어있으면(NO_CONTEXT) 여기서 바로 끝난다.
- * 2. processJob() — RagJobWorker가 PROCESSING row를 하나씩 꺼내 호출. 실제 OllamaClient 호출과
- * 결과 영속화(rag_responses, response_citations)를 담당한다.
+ * 2. processJob() — RagJobWorker가 PROCESSING row를 하나 claim할 때마다 호출. 실제 OllamaClient
+ * 호출과 결과 영속화(rag_responses, response_citations)를 담당한다.
*
*
+ *
+ *
*/
private static final String LLM_FALLBACK_PREFIX = "AI 답변 생성이 지연되고 있습니다. "
+ "가장 관련도 높은 문서에서 다음 내용을 찾았습니다:\n\n";
@@ -152,9 +156,9 @@ public RagEnqueueOutcome enqueue(
* 반환하고 Ollama를 아예 호출하지 않는다 — RagJobWorker가 이 job을 집어든 뒤, 여기서
* {@code findById}로 다시 읽기 전에 RagJobTimeoutSweeper가 먼저 강제 종료했을 수 있다.
* 이 조기 반환이 없으면 이미 끝난 job에도 Ollama 호출(수십 초)을 그대로 낭비하게 되는데,
- * Worker/GPU가 1개뿐이라 그 시간만큼 뒤에 대기 중인 다른 job까지 더 늦어진다 — 아래
- * completeSuccess/completeFailed의 조건부 UPDATE는 이 조기 체크 "이후"에 벌어지는 경합(더
- * 좁은 창)까지 막아주는 최종 방어선이다.
+ * Worker 슬롯이 {@code rag.worker.max-concurrency}개뿐이라(#340) 그 슬롯 하나가 낭비되는
+ * 동안 전체 동시 처리량이 그만큼 줄어든다 — 아래 completeSuccess/completeFailed의 조건부
+ * UPDATE는 이 조기 체크 "이후"에 벌어지는 경합(더 좁은 창)까지 막아주는 최종 방어선이다.
*/
public boolean processJob(Long jobId) {
RagResponse job = ragResponseRepository.findById(jobId)
@@ -168,9 +172,11 @@ public boolean processJob(Long jobId) {
try {
result = ollamaClient.generate(job.getPromptText());
} catch (DocGridException e) {
- // LLM 장애가 권한 검증을 통과한 벡터 검색 결과까지 숨기지 않도록, 최상위 후보 원문을
- // 그대로 인용해 최소한의 답을 제공한다(extractive fallback). 이 fallback은 비동기 전환
- // 이전과 달리 rag_responses에 그대로 영속화된다 — 나중에 GET/조회로 이 값을 그대로 돌려준다.
+ /**
+ * LLM 장애가 권한 검증을 통과한 벡터 검색 결과까지 숨기지 않도록, 최상위 후보 원문을
+ * 그대로 인용해 최소한의 답을 제공한다(extractive fallback). 이 fallback은 비동기 전환
+ * 이전과 달리 rag_responses에 그대로 영속화된다 — 나중에 GET/조회로 이 값을 그대로 돌려준다.
+ */
List
+ * claim 단계 (이 클래스, 수 ms) : 행 잠금 걸림 → claimed_at 기록 → 커밋과 동시에 잠금 풀림
+ * 처리 단계 (RagFacade, 수십 초) : 잠금 없이 Ollama 호출. claimed_at 값만으로 소유권 유지
+ *
+ * 잠금을 처리 단계까지 들고 있으면 그 job 하나 때문에 다른 claim 시도 전체가 수십 초씩 막힌다
+ * — 그래서 잠금은 claim 단계에서만 잠깐 쓰고, 풀린 뒤의 장기 소유권 표시는 값 하나
+ * ({@code claimed_at})로 넘긴다.
+ *
*