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
1 change: 1 addition & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ dependencies {
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testImplementation 'org.springframework.security:spring-security-test'
testImplementation 'org.postgresql:postgresql'
testImplementation 'org.awaitility:awaitility'
testCompileOnly 'org.projectlombok:lombok'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
testAnnotationProcessor 'org.projectlombok:lombok'
Expand Down
554 changes: 554 additions & 0 deletions docs/design/kangcheolung-#151-ragops-dashboard-realtime-update.md

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
package com.opensource.docgrid.domain.dashboard.config;

import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableScheduling;

/**
* Dashboard debounce push 스케줄러 전용 스케줄링 활성화.
*
* <p>{@code WorkerSchedulingConfig}의 {@code @EnableScheduling}은
* {@code indexing.worker.enabled=true}일 때만 켜지는 조건부 설정이고 기본값은 {@code false}다.
* 그 설정에 얹혀가면 Worker 기능이 꺼진 기본 상태에서 대시보드 push 스케줄러도 같이 동작하지
* 않게 된다. 대시보드 push는 자동 Worker On/Off와 무관하게(관리자 수동 재처리만으로도) 항상
* 동작해야 하므로 조건 없이 별도로 스케줄링을 켠다. {@code @EnableScheduling}을 여러 설정
* 클래스에서 선언해도 스프링이 안전하게 처리한다.
*/
@Configuration
@EnableScheduling
public class DashboardSchedulingConfig {
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package com.opensource.docgrid.domain.dashboard.event;

import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;

import com.opensource.docgrid.domain.dashboard.controller.DashboardWebSocketController;
import com.opensource.docgrid.domain.dashboard.service.query.DashboardQueryService;

import lombok.RequiredArgsConstructor;

/**
* 짧은 주기로 {@link DashboardUpdateFlag}를 확인해서, 그 사이 상태 전이가 있었을 때만 대시보드
* 집계를 다시 계산해 push하는 debounce 스케줄러.
*
* <p>이벤트가 몇 번 들어왔든 한 주기(기본 300ms)당 최대 1번만 {@code getSummary()}(집계 쿼리 약
* 9개)와 push를 실행한다. Burst 상황(전체 재처리, 시스템 장애로 다건 실패)에서 이벤트 개수만큼
* DB를 두드리는 걸 막는 게 이 클래스의 유일한 목적이다.
*
* <p>이 클래스는 새 로직을 거의 안 만들고, 이미 있는 세 부품을 "언제 조합해서 실행할지"만
* 정한다: "바뀌었는지 확인"은 {@link DashboardUpdateFlag}, "최신 집계 계산"은
* {@code DashboardQueryService}(대시보드 집계 조회), "WebSocket 전송"은
* {@code DashboardWebSocketController}(WebSocket 전송)가 이미 만들어 둔 것을 그대로 가져다 쓴다.
*/
@Component
@RequiredArgsConstructor
public class DashboardPushScheduler {

private final DashboardUpdateFlag dashboardUpdateFlag;
private final DashboardQueryService dashboardQueryService;
private final DashboardWebSocketController dashboardWebSocketController;

/**
* {@code fixedDelayString}이라 "이전 실행이 끝난 시점부터" 설정된 간격 뒤에 다음 실행이
* 잡힌다({@code fixedRate}처럼 "시작 시점 기준 고정 주기"가 아니다) — 이 메서드가 드물게
* 오래 걸려도 다음 실행과 겹치거나 밀리지 않는 안전한 쪽을 선택했다. 주기 값은
* {@code application.yml}의 {@code dashboard.push.debounce-interval-ms}에서 읽는다.
*/
@Scheduled(fixedDelayString = "${dashboard.push.debounce-interval-ms}")
public void pushIfDirty() {
// 1. 그 사이 상태 전이가 있었는지 확인하면서 동시에 플래그를 내린다(원자적).
if (dashboardUpdateFlag.consumeIfDirty()) {
// 2. 있었을 때만 최신 집계를 처음부터 다시 계산해서 push한다. 없었으면 이 블록 자체가
// 실행되지 않으므로 DB 조회도 push도 전혀 일어나지 않는다.
try {
dashboardWebSocketController.sendDashboardUpdate(dashboardQueryService.getSummary());
} catch (RuntimeException exception) {
// 집계나 push가 실패하면 이미 소비해버린 dirty 신호를 되살린다 — 안 그러면 다음
// 상태 전이가 안 들어오는 한 대시보드가 갱신 신호를 영영 잃어버린 채로 남는다.
dashboardUpdateFlag.markDirty();
throw exception;
}
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
package com.opensource.docgrid.domain.dashboard.event;

import java.util.concurrent.atomic.AtomicBoolean;

import org.springframework.stereotype.Component;

/**
* "대시보드 집계가 최신이 아니다"라는 사실만 기억하는 스레드 안전한 플래그.
*
* <p>여러 Worker Thread가 동시에 {@link #markDirty()}를 호출해도 안전하며, DB 조회나 WebSocket
* push 같은 무거운 작업은 전혀 하지 않는다. 실제 집계 재계산과 push는 이 플래그를 주기적으로
* 확인하는 스케줄러({@code DashboardPushScheduler})가 전담한다. 짧은 시간에 여러 상태 전이가
* 몰려도(burst) 플래그는 계속 true로만 유지되므로, 스케줄러 입장에서는 몇 번 세워졌는지와 무관하게
* "그 사이 뭔가 바뀌었다"는 사실 하나만 확인하면 된다 — 이 방식으로 이벤트 개수만큼 집계 쿼리가
* 늘어나는 걸 막는다.
*
* <p>실제 흐름 예시:
* <pre>
* [시작] dirty = false
*
* Worker가 job A를 claim → 리스너가 markDirty() 호출 → dirty = true
* ... 0~300ms 사이 ...
* 스케줄러가 300ms마다 도는 시점 → consumeIfDirty() 호출
* → dirty가 true였으니 true 반환, 동시에 dirty = false로 리셋
* → 스케줄러: "바뀐 게 있었네" → getSummary() 계산 → WebSocket push
*
* ... 다음 300ms, 아무도 markDirty()를 안 부른 경우 ...
* 스케줄러가 다시 consumeIfDirty() 호출 → dirty가 false였으니 false 반환
* → 스케줄러: "바뀐 거 없네" → 아무것도 안 하고 다음 주기까지 대기
* </pre>
*/
@Component
public class DashboardUpdateFlag {

// 앱이 막 시작된 시점엔 아직 어떤 Job도 상태가 안 바뀌었으므로 "최신 상태"인 false로 시작한다.
private final AtomicBoolean dirty = new AtomicBoolean(false);

/**
* "화면이 최신이 아니다"는 표시등을 켠다. {@code EmbeddingJobStatusChangedEventListener}가
* Job 상태가 바뀔 때마다 호출한다. 이미 켜져 있는 상태에서 또 호출해도 결과는 그대로 켜진
* 상태 하나뿐이라, burst로 여러 번 호출돼도 스케줄러 입장에서는 구분이 안 된다(의도된 동작).
*/
public void markDirty() {
dirty.set(true);
}

/**
* 플래그가 서 있으면 원자적으로 내리면서 {@code true}를 반환한다.
* {@code DashboardPushScheduler}가 debounce 주기마다 호출해서 "그 사이 뭔가 바뀌었는지"를
* 확인하고, 바뀌었으면 그때 가서 집계를 다시 계산해 push한다.
*
* <p>"확인"과 "초기화"를 한 번의 원자 연산으로 묶어야, 스케줄러가 확인하는 순간과 다음 이벤트가
* 플래그를 세우는 순간이 겹쳐도 변경 신호를 놓치지 않는다.
*/
public boolean consumeIfDirty() {
return dirty.compareAndSet(true, false);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package com.opensource.docgrid.domain.dashboard.event;

import org.springframework.stereotype.Component;
import org.springframework.transaction.event.TransactionPhase;
import org.springframework.transaction.event.TransactionalEventListener;

import com.opensource.docgrid.domain.embedding.event.EmbeddingJobStatusChangedEvent;

import lombok.RequiredArgsConstructor;

/**
* Embedding Job 상태 전이 이벤트를 받아 대시보드 갱신 플래그만 세우는 구독자.
*
* <p>이 클래스 자체는 대시보드를 갱신하지 않는다 — "커밋이 확실히 된 경우에만 갱신이 필요하다는
* 신호를 남기는" 역할까지만 한다. 실제 집계 재계산과 WebSocket push는 이 플래그를 나중에 확인하는
* {@code DashboardPushScheduler}가 한다.
*
* <p>{@code phase = AFTER_COMMIT}이라 이벤트를 발행한 Transaction이 실제로 커밋된 뒤에만 호출된다
* — 롤백되면 이 메서드 자체가 실행되지 않으므로 확정되지 않은 상태 변화로 대시보드가 갱신되는 일이
* 없다. 여기서 바로 집계를 계산하거나 push하지 않고 {@link DashboardUpdateFlag#markDirty()}만
* 호출하는 이유는 {@code DashboardUpdateFlag}의 클래스 설명 참고 — burst 상황에서 이벤트 개수만큼
* DB 조회가 늘어나는 걸 막기 위함이다.
*/
@Component
@RequiredArgsConstructor
public class EmbeddingJobStatusChangedEventListener {

private final DashboardUpdateFlag dashboardUpdateFlag;

@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void onEmbeddingJobStatusChanged(EmbeddingJobStatusChangedEvent event) {
dashboardUpdateFlag.markDirty();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
package com.opensource.docgrid.domain.embedding.event;

/**
* Embedding Job의 상태가 바뀌었음을 알리는 마커 이벤트.
*
* <p>대시보드가 최신 집계를 다시 계산해야 한다는 신호일 뿐이라 어떤 상태에서 어떤 상태로
* 바뀌었는지는 담지 않는다. 구독 측은 이벤트 발생 시점에 항상 전체 집계를 새로 계산하므로
* 세부 상태를 실어봐야 쓰이지 않는다.
*/
public record EmbeddingJobStatusChangedEvent(Long jobId) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import java.util.Objects;
import java.util.Set;

import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.util.StringUtils;
Expand All @@ -25,6 +26,7 @@
import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus;
import com.opensource.docgrid.domain.embedding.enums.EmbeddingStatus;
import com.opensource.docgrid.domain.embedding.event.EmbeddingJobStatusChangedEvent;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository;
import com.opensource.docgrid.domain.worker.entity.EmbeddingJobAttempt;
Expand Down Expand Up @@ -72,6 +74,7 @@ public class DocumentIndexingCompletionService {
private final IndexingEventRepository indexingEventRepository;
private final EmbeddingJobOwnershipValidator ownershipValidator;
private final Clock clock;
private final ApplicationEventPublisher applicationEventPublisher;

/**
* 현재 Claim 실행이 생성한 전체 Embedding Set을 문서의 검색 가능 상태로 확정한다.
Expand Down Expand Up @@ -134,6 +137,10 @@ public DocumentIndexingCompletionResponse complete(
durationMs
);

// 6. 대시보드가 최신 집계를 다시 계산하도록 상태 전이를 알린다. AFTER_COMMIT 구독자만
// 반응하므로 이 Transaction이 실제로 커밋된 뒤에만 push로 이어진다.
applicationEventPublisher.publishEvent(new EmbeddingJobStatusChangedEvent(embeddingJob.getId()));

log.info(
"문서 인덱싱 완료: jobId={}, attemptId={}, documentId={}, versionId={}, "
+ "chunkCount={}, embeddingCount={}, durationMs={}",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import java.util.Optional;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

Expand All @@ -18,6 +19,7 @@
import com.opensource.docgrid.domain.embedding.dto.response.DocumentIndexingFailureResponse;
import com.opensource.docgrid.domain.embedding.entity.EmbeddingJob;
import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus;
import com.opensource.docgrid.domain.embedding.event.EmbeddingJobStatusChangedEvent;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository;
import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties;
Expand Down Expand Up @@ -48,6 +50,7 @@ public class DocumentIndexingFailureService {
private final EmbeddingJobAttemptConverter attemptConverter;
private final IndexingFailureTransitionService failureTransitionService;
private final Clock clock;
private final ApplicationEventPublisher applicationEventPublisher;

/**
* 운영 경로에서 공통 실패 전이 Service를 주입받는 생성자.
Expand All @@ -59,14 +62,16 @@ public DocumentIndexingFailureService(
EmbeddingJobOwnershipValidator ownershipValidator,
EmbeddingJobAttemptConverter attemptConverter,
IndexingFailureTransitionService failureTransitionService,
Clock clock
Clock clock,
ApplicationEventPublisher applicationEventPublisher
) {
this.embeddingJobRepository = embeddingJobRepository;
this.embeddingJobAttemptRepository = embeddingJobAttemptRepository;
this.ownershipValidator = ownershipValidator;
this.attemptConverter = attemptConverter;
this.failureTransitionService = failureTransitionService;
this.clock = clock;
this.applicationEventPublisher = applicationEventPublisher;
}

/**
Expand All @@ -82,7 +87,8 @@ public DocumentIndexingFailureService(
EmbeddingJobOwnershipValidator ownershipValidator,
EmbeddingJobAttemptConverter attemptConverter,
IndexingWorkerProperties workerProperties,
Clock clock
Clock clock,
ApplicationEventPublisher applicationEventPublisher
) {
this(
embeddingJobRepository,
Expand All @@ -96,7 +102,8 @@ public DocumentIndexingFailureService(
indexingEventRepository,
workerProperties
),
clock
clock,
applicationEventPublisher
);
}

Expand Down Expand Up @@ -144,6 +151,12 @@ public DocumentIndexingFailureResponse fail(
failedAt
);

// 4. 대시보드가 최신 집계를 다시 계산하도록 상태 전이를 알린다. transition()이 재시도 예약
// (PENDING)과 최종 실패(FAILED) 중 어느 쪽으로 끝났든 embeddingJob은 같은 영속 인스턴스라
// 최종 상태를 그대로 반영한다. AFTER_COMMIT 구독자만 반응하므로 이 Transaction이 실제로
// 커밋된 뒤에만 push로 이어진다.
applicationEventPublisher.publishEvent(new EmbeddingJobStatusChangedEvent(embeddingJob.getId()));

log.info(
"문서 인덱싱 실패 기록: jobId={}, attemptId={}, failureType={}, retryCount={}, terminal={}",
embeddingJob.getId(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,15 @@
import java.util.Set;
import java.util.UUID;

import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

import com.opensource.docgrid.domain.embedding.converter.EmbeddingJobConverter;
import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse;
import com.opensource.docgrid.domain.embedding.entity.EmbeddingJob;
import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus;
import com.opensource.docgrid.domain.embedding.event.EmbeddingJobStatusChangedEvent;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties;
import com.opensource.docgrid.domain.worker.entity.IndexingEvent;
Expand Down Expand Up @@ -50,6 +52,7 @@ public class EmbeddingJobClaimService {
private final EmbeddingJobConverter embeddingJobConverter;
private final IndexingWorkerProperties indexingWorkerProperties;
private final Clock clock;
private final ApplicationEventPublisher applicationEventPublisher;

/**
* Worker가 처리할 다음 PENDING Job을 Claim한다.
Expand Down Expand Up @@ -95,7 +98,11 @@ private ClaimedEmbeddingJobResponse claim(
.occurredAt(claimedAt)
.build());

// 4. Transaction 안에서 LAZY 연관 식별자를 읽어 Claim 결과 DTO를 완성한다.
// 4. 대시보드가 최신 집계를 다시 계산하도록 상태 전이를 알린다. AFTER_COMMIT 구독자만
// 반응하므로 이 Transaction이 실제로 커밋된 뒤에만 push로 이어진다.
applicationEventPublisher.publishEvent(new EmbeddingJobStatusChangedEvent(embeddingJob.getId()));

// 5. Transaction 안에서 LAZY 연관 식별자를 읽어 Claim 결과 DTO를 완성한다.
return embeddingJobConverter.toClaimedResponse(embeddingJob);
}

Expand Down
10 changes: 8 additions & 2 deletions src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,14 @@ spring:
task:
scheduling:
pool:
# DB Claim 지연이 Worker Heartbeat와 Lease 복구 실행을 막지 않도록 세 작업을 분리한다.
size: 3
# DB Claim 지연이 Worker Heartbeat·Lease 복구 실행을 막지 않도록, Dashboard debounce push까지
# 포함해 네 작업을 분리한다.
size: 4

dashboard:
push:
# 짧은 시간에 몰리는 Job 상태 전이 이벤트를 이 주기로 모아서 한 번만 집계·push한다(debounce).
debounce-interval-ms: ${DASHBOARD_PUSH_DEBOUNCE_INTERVAL_MS:300}

document:
upload:
Expand Down
Loading