Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions .claude/launch.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
{
"version": "0.0.1",
"configurations": [
{
"name": "backend",
"runtimeExecutable": "./backend/gradlew",
"runtimeArgs": ["-p", "backend", "bootRun", "--args=--spring.profiles.active=local"],
"port": 8080
}
]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
package com.opensource.docgrid.domain.rag.config;

import java.util.concurrent.Semaphore;
import java.util.concurrent.SynchronousQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.CustomizableThreadFactory;

/**
* RAG 답변 생성을 동시에 몇 건까지 처리할지 제한하는 실행 자원을 구성한다 (#340).
*
* <p>{@code domain/worker/config/WorkerExecutionConfig}(이 코드베이스에서 유일했던 커스텀
* 스레드풀 선례)를 그대로 본떴다 — {@code core=max}인 {@link ThreadPoolExecutor} +
* {@link SynchronousQueue}(큐잉 없음) + {@link ThreadPoolExecutor.AbortPolicy}(꽉 차면 즉시
* 거부, 조용히 쌓아두지 않음)로 "설정된 동시성만 즉시 실행"을 보장한다.
*
* <p>{@code embedding_jobs}의 {@code WorkerExecutionSlotPool}(별도 클래스, 종료 플래그 +
* introspection 메서드 포함)까지는 필요 없다 — RAG는 별도 워커 등록/우아한 종료 조율이나
* 대시보드 노출 요구가 없어서, 같은 안전 성질(로컬 슬롯을 먼저 확보한 뒤에만 DB claim을
* 시도해 "claim은 됐는데 실행할 스레드가 없는" 상태를 만들지 않는 것)을 순수
* {@link Semaphore}만으로 재현한다.
*
* <p>{@code indexing.worker.enabled}로 켜고 끌 수 있는 인덱싱 워커와 달리, RAG Worker는
* {@code RagSchedulingConfig}와 동일하게 조건 없이 항상 켜져 있어야 하는 검색 API 핵심
* 경로라 {@code @ConditionalOnProperty}를 붙이지 않는다.
*/
@Configuration
public class RagExecutionConfig {

public static final String RAG_WORKER_JOB_EXECUTOR = "ragWorkerJobExecutor";
public static final String RAG_WORKER_SLOTS = "ragWorkerSlots";

/**
* 동시에 최대 {@code rag.worker.max-concurrency}건까지만 즉시 실행하는 무대기 Executor를 만든다.
*/
@Bean(name = RAG_WORKER_JOB_EXECUTOR, destroyMethod = "shutdownNow")
public ThreadPoolExecutor ragWorkerJobExecutor(
@Value("${rag.worker.max-concurrency:2}") int maxConcurrency
) {
validateMaxConcurrency(maxConcurrency);
return new ThreadPoolExecutor(
maxConcurrency,
maxConcurrency,
0L,
TimeUnit.MILLISECONDS,
new SynchronousQueue<>(),
new CustomizableThreadFactory("rag-worker-job-"),
new ThreadPoolExecutor.AbortPolicy()
);
}

/**
* DB claim을 시도하기 전에 먼저 확보해야 하는 로컬 실행 슬롯. 공정 모드(fair)로 만들어
* 폴링 주기가 겹칠 때 대기가 한쪽으로 몰리지 않게 한다 — {@code WorkerExecutionSlotPool}의
* 선택과 동일하다.
*/
@Bean(name = RAG_WORKER_SLOTS)
public Semaphore ragWorkerSlots(@Value("${rag.worker.max-concurrency:2}") int maxConcurrency) {
validateMaxConcurrency(maxConcurrency);
return new Semaphore(maxConcurrency, true);
}

/**
* {@code max-concurrency=0}(또는 음수)은 예외 없이 Worker를 영구 대기 상태로 만들 수 있어
* ({@link Semaphore}는 permit 0으로도 생성 자체는 허용) 기동 시점에 바로 실패시킨다.
*/
private void validateMaxConcurrency(int maxConcurrency) {
if (maxConcurrency < 1) {
throw new IllegalArgumentException(
"rag.worker.max-concurrency는 1 이상이어야 합니다: " + maxConcurrency
);
}
}
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
package com.opensource.docgrid.domain.rag.entity;

import java.time.LocalDateTime;

import com.opensource.docgrid.domain.search.entity.SearchQuery;
import com.opensource.docgrid.domain.search.enums.ResultStatus;
import com.opensource.docgrid.global.common.entity.BaseEntity;
Expand Down Expand Up @@ -88,6 +90,12 @@ public class RagResponse extends BaseEntity {
@Column(name = "error_message", columnDefinition = "TEXT")
private String errorMessage;

// 병렬 Worker가 이 job을 이미 집었는지 표시한다(#340). status만으로는 "대기 중"과 "누가 이미
// 처리 중"을 구분할 수 없어서(둘 다 PROCESSING) 별도로 둔다. RagResponseClaimService의 짧은
// claim 트랜잭션 안에서만 채워지며, 그 밖의 완료 확정 경로(조건부 UPDATE)는 이 컬럼을 건드리지 않는다.
@Column(name = "claimed_at")
private LocalDateTime claimedAt;

@Builder
public RagResponse(SearchQuery query, String answerText, String llmProvider, String llmModelName,
String promptText, Integer inputTokenCount, Integer outputTokenCount, Integer latencyMs,
Expand All @@ -103,4 +111,13 @@ public RagResponse(SearchQuery query, String answerText, String llmProvider, Str
this.status = status;
this.errorMessage = errorMessage;
}

/**
* 이 job을 지금 이 Worker가 처리하기 시작했다는 표시를 남긴다. {@code RagResponseClaimService}의
* 짧은 claim 트랜잭션 안에서만 호출되어야 한다 — dirty checking으로 반영되므로, 이 엔티티가
* detached된 뒤(다른 트랜잭션/스레드로 넘어간 뒤)에 호출하면 반영되지 않는다.
*/
public void markClaimed(LocalDateTime claimedAt) {
this.claimedAt = claimedAt;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,17 +22,49 @@
public interface RagResponseRepository extends JpaRepository<RagResponse, Long> {

/**
* 주어진 상태(보통 PROCESSING)인 것들 중 가장 오래 기다린 것 하나를 반환한다 — RagJobWorker가
* 1초마다 폴링하며 이 메서드로 FIFO 큐를 구현한다. Worker가 1개뿐이라 별도 락/claim 없이도
* 안전하다.
*
* <p>{@code query}/{@code query.user}를 {@link EntityGraph}로 미리 fetch한다 — Worker가
* 이 메서드로 job을 꺼낸 트랜잭션이 끝난 뒤(WebSocket push 시점)에
* {@code job.getQuery().getUser().getEmail()}에 접근해도 두 연관관계 모두 LAZY라서 자칫
* {@code LazyInitializationException}이 날 수 있는데, 미리 로딩해두면 그 문제가 없다.
* PROCESSING 중 아직 아무 Worker도 집지 않은(claim 안 된) 것 하나를 골라 행 잠금을 건다(#340).
* {@code claimed_at IS NULL} 조건이 "대기 중"과 "이미 처리 중"을 구분하는 유일한 신호다 —
* status만으로는 둘 다 PROCESSING이라 구분이 안 된다. {@code FOR UPDATE SKIP LOCKED}로
* 여러 Worker가 동시에 이 쿼리를 날려도 이미 잠긴 행은 건너뛰고 그다음 미잠금 행을 잡아온다.
* {@code created_at}이 같은 밀리초를 공유할 수 있는 동시 접수 상황을 대비해 {@code id}를
* 2차 정렬 기준으로 둔다. {@link RagResponseClaimService}가 이 메서드로 잠근 행을 같은 짧은
* 트랜잭션 안에서 즉시 {@link RagResponse#markClaimed}로 확정하고 커밋해, 락을 오래 들고
* 있지 않는다({@code embedding_jobs}의 claim 패턴과 동일).
*/
@Query(value = """
SELECT * FROM rag_responses
WHERE status = 'PROCESSING' AND claimed_at IS NULL
ORDER BY created_at ASC, id ASC
LIMIT 1
FOR UPDATE SKIP LOCKED
""", nativeQuery = true)
Optional<RagResponse> findNextUnclaimedProcessingForUpdate();

/**
* claim된 job을 실제로 처리하는 Worker 스레드가 쓰는 조회. {@code query}/{@code query.user}를
* {@link EntityGraph}로 미리 fetch해, 처리가 끝난 뒤(WebSocket push 시점) 트랜잭션 밖에서
* {@code job.getQuery().getUser().getEmail()}에 접근해도 {@code LazyInitializationException}이
* 나지 않게 한다. 기존 {@code findById}는 그대로 두고 이름을 다르게 둔 이유는, {@code RagFacade}가
* 이미 쓰고 있는 평범한 {@code findById(jobId)} 호출의 동작을 이번 변경으로 건드리지 않기 위함이다.
*/
@EntityGraph(attributePaths = {"query", "query.user"})
Optional<RagResponse> findFirstByStatusOrderByCreatedAtAsc(ResultStatus status);
Optional<RagResponse> findWithQueryAndUserById(Long id);

/**
* 앱 재시작 복구 전용(#340 CodeRabbit 리뷰 반영). 이전 프로세스가 claim한 채 완료하지 못하고
* 죽은 job은 {@code claimed_at}이 채워진 상태로 DB에 남는다 — 이 상태로는
* {@link #findNextUnclaimedProcessingForUpdate}가 절대 다시 집어주지 않아, 스위퍼의
* {@code stale-threshold} 강제종료(fallback)만 기다리게 된다. 재시작 직후 한 번,
* PROCESSING인데 claim만 남아있는 행의 claim을 전부 풀어 새 Worker가 다시 시도할 수 있게
* 한다. "인스턴스는 항상 1개"라는 이 프로젝트의 전제 위에서만 안전하다 — 이 메서드가
* 실행되는 시점엔 다른 프로세스가 진짜로 처리 중일 수 없으므로, claim이 남아있는 행은
* 전부 죽은 이전 프로세스의 흔적이다.
*/
@Modifying(clearAutomatically = true)
@Query("UPDATE RagResponse r SET r.claimedAt = NULL "
+ "WHERE r.status = com.opensource.docgrid.domain.search.enums.ResultStatus.PROCESSING "
+ "AND r.claimedAt IS NOT NULL")
int releaseAllClaimsOnStartup();

/** 특정 검색 요청(queryId)에 대한 RAG 답변을 찾는다. GET /search/{queryId} 재조회에 쓰인다. */
Optional<RagResponse> findByQuery_Id(Long queryId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,7 @@ public RagEnqueueOutcome enqueue(
* extractive fallback을 채운 채 FAILED로 확정한다.
*
* <p>{@code job} 객체가 아니라 {@code jobId}만 받아 이 메서드 자신의 트랜잭션 안에서 다시
* 조회하는 이유: RagJobWorker가 {@code findFirstByStatusOrderByCreatedAtAsc()}로 꺼낸
* 조회하는 이유: RagJobWorker가 claim 단계({@code RagResponseClaimService}, #340)에서 꺼낸
* job은 그 조회 시점에 트랜잭션이 끝나 detached 상태다. 원래(#218) 이 detached 인스턴스를
* 그대로 받아 필드만 바꾸면 dirty checking이 감지 못해 DB에 반영되지 않는 버그가 있었는데,
* 지금은 완료 처리 자체가 dirty checking에 의존하지 않는다({@link
Expand Down Expand Up @@ -229,7 +229,11 @@ public boolean processJob(Long jobId) {
}
// 답변과 citation 저장이 모두 끝난 동일 Transaction의 커밋 이후 성공 Counter를 기록한다.
applicationEventPublisher.publishEvent(new RagJobCompletionMetricEvent(Outcome.SUCCESS));
log.info("[RAG] done queryId={} responseId={} latencyMs={}", queryId, job.getId(), result.latencyMs());
// promptTokens/answerTokens을 함께 남겨, 느린 job이 프롬프트를 읽느라(prefill) 오래 걸린 건지
// 답변을 쓰느라(decode) 오래 걸린 건지 로그만으로 구분할 수 있게 한다 — 병렬화(#340) 이후
// 요청당 작업량을 어느 쪽부터 줄여야 할지 판단하는 근거 자료.
log.info("[RAG] done queryId={} responseId={} latencyMs={} promptTokens={} answerTokens={}",
queryId, job.getId(), result.latencyMs(), result.inputTokenCount(), result.outputTokenCount());
return true;
}

Expand Down
Loading