diff --git a/build.gradle b/build.gradle
index a2830ab..71c2236 100644
--- a/build.gradle
+++ b/build.gradle
@@ -52,7 +52,7 @@ dependencies {
tasks.named('test') {
useJUnitPlatform {
- excludeTags 'benchmark', 'minio-integration', 'claim-concurrency', 'local-e2e', 'vector-search-performance'
+ excludeTags 'benchmark', 'minio-integration', 'claim-concurrency', 'local-e2e', 'vector-search-performance', 'worker-indexing-throughput'
}
}
@@ -152,6 +152,28 @@ tasks.register('vectorSearchPerformanceTest', Test) {
outputs.upToDateWhen { false }
}
+tasks.register('workerIndexingThroughputTest', Test) {
+ group = 'verification'
+ description = '실제 PostgreSQL, MinIO와 BGE-M3에서 자동 Worker 전체 문서 인덱싱 처리량을 측정합니다.'
+ testClassesDirs = sourceSets.test.output.classesDirs
+ classpath = sourceSets.test.runtimeClasspath
+ useJUnitPlatform {
+ includeTags 'worker-indexing-throughput'
+ }
+ maxParallelForks = 1
+ systemProperties System.properties.findAll { key, value ->
+ key.toString().startsWith('worker.indexing.throughput.')
+ }
+ if (System.getProperty('worker.indexing.throughput.output') == null) {
+ systemProperty(
+ 'worker.indexing.throughput.output',
+ layout.buildDirectory.file('reports/worker-indexing-throughput/worker-indexing-throughput.json').get().asFile.absolutePath
+ )
+ }
+ // 실제 외부 Service를 점유하고 테스트 데이터를 초기화하는 Benchmark이므로 일반 Test와 Build Cache에서 분리한다.
+ outputs.upToDateWhen { false }
+}
+
def configureOpenSqlDatabase = { Test task ->
task.maxParallelForks = 1
task.outputs.upToDateWhen { false }
diff --git a/docs/design/gimin-#133-worker-indexing-throughput-benchmark.md b/docs/design/gimin-#133-worker-indexing-throughput-benchmark.md
new file mode 100644
index 0000000..040b948
--- /dev/null
+++ b/docs/design/gimin-#133-worker-indexing-throughput-benchmark.md
@@ -0,0 +1,176 @@
+# 자동 Worker 전체 문서 인덱싱 처리량 Benchmark 설계
+
+## 1. 배경
+
+자동 Worker는 Job 실행 슬롯을 먼저 확보한 뒤 `PENDING` Job을 Claim하고, Parsing, Chunk 저장,
+Batch Embedding, Vector 저장과 검색 Version 전환까지 수행한다. 기존 통합·E2E 테스트는 상태 전이와
+소유권 정합성을 검증하지만 여러 문서를 연속 처리할 때의 처리량과 지연 기준선은 제공하지 않는다.
+
+이번 작업은 실제 PostgreSQL 17, pgvector, MinIO와 `BAAI/bge-m3`를 연결한 자동 Worker 전체
+Pipeline을 반복 측정해 다음 질문에 답한다.
+
+1. 기본 Worker 동시성에서 문서와 Chunk를 초당 얼마나 인덱싱하는가?
+2. 업로드 접수, Queue 대기와 실제 처리 시간 중 어느 구간이 전체 지연을 지배하는가?
+3. Queue를 모두 소진할 때 처리 누락, 중복 Attempt 또는 중복 Vector가 발생하지 않는가?
+
+## 2. 범위
+
+### 2.1 포함
+
+- 실제 인증과 Multipart HTTP 문서 업로드
+- 실제 MinIO Object 저장과 읽기
+- 실제 자동 Polling, Claim, Attempt, Lease 갱신과 전체 인덱싱 Pipeline
+- 실제 BGE-M3 Batch Embedding과 PostgreSQL `vector(1024)` 저장
+- 결정적 TXT Corpus와 같은 Chunk 분포를 사용한 반복 측정
+- Warm-up과 본 측정 분리
+- 문서·Chunk·Embedding 처리량과 Queue·처리·전체 지연 수집
+- Profile별 상태, Attempt, Event, Chunk와 Vector 정합성 검증
+- PostgreSQL, pgvector, BGE, Batch Size와 Worker 설정 환경 지문 수집
+- Git 제외 JSON 원본 결과와 실행 결과 Markdown 기록
+- 일반 테스트와 분리된 전용 Gradle Task
+
+### 2.2 제외
+
+- Worker 프로세스 수에 따른 수평 확장 비교
+- Queue 허용 한계와 Backpressure 붕괴 지점 측정
+- Lease 갱신 주기별 DB 부하 비교
+- BGE-M3 Batch Size 재선정
+- PDF·DOCX Parser 형식별 성능 비교
+- 공식 OpenSQL 공급사 환경의 최종 성능 수치
+- 운영 SLO 확정
+
+## 3. Workload 계약
+
+### 3.1 결정적 문서
+
+- 모든 Profile은 같은 UTF-8 TXT Corpus Template을 사용한다.
+- 문서마다 고유한 식별 문구만 바꾸고 본문 길이와 문단 구조는 동일하게 유지한다.
+- 본문은 기본 Chunk Size와 Overlap에서 같은 수의 Chunk가 생성되도록 고정한다.
+- 파일 이름과 Document 제목은 Profile, 반복과 문서 순번을 포함해 충돌을 방지한다.
+- 전체 Profile의 실제 Chunk 수가 같지 않으면 처리량 비교를 실패로 처리한다.
+
+### 3.2 Profile
+
+기본값은 짧은 로컬 실행과 Queue가 유지되는 본 측정을 함께 제공한다.
+
+| 구분 | 문서 수 | 반복 | 통계 포함 |
+|---|---:|---:|---|
+| Warm-up | 4 | 1 | 제외 |
+| 작은 Queue | 16 | 3 | 포함 |
+| 지속 Queue | 32 | 3 | 포함 |
+
+문서 수와 반복은 `worker.indexing.throughput.*` System Property로 변경할 수 있다. Worker
+`max-concurrency`는 제품 기본값인 `2`로 고정하고 결과 환경 지문에 기록한다. Worker 수평 확장은
+별도 Benchmark에서 다룬다.
+
+### 3.3 측정 구간
+
+1. Warm-up 문서를 업로드하고 모두 `INDEXED`가 될 때까지 기다린다.
+2. Profile 시작 시각을 기록하고 여러 Uploader Thread로 측정 문서를 접수한다.
+3. 마지막 업로드 완료 시각을 기록한다.
+4. 자동 Worker가 Profile의 모든 Job을 `INDEXED`로 전환할 때까지 기다린다.
+5. 상태와 결과 정합성을 검증한 뒤 통계를 계산한다.
+
+Upload와 Worker 실행이 겹치는 실제 동작을 유지한다. 대신 업로드 시간과 마지막 업로드 뒤 Queue
+소진 시간을 분리해 HTTP·MinIO 접수가 Worker 처리량을 가리는지 확인한다.
+
+## 4. 지표 계약
+
+| 지표 | 계산 |
+|---|---|
+| documents/s | 완료 문서 수 / Profile 전체 경과 시간 |
+| documents/min | documents/s × 60 |
+| chunks/s | 저장 Chunk 수 / Profile 전체 경과 시간 |
+| embeddings/s | 저장 Embedding 수 / Profile 전체 경과 시간 |
+| Upload 경과 | 첫 Upload 시작부터 마지막 Upload 응답까지 |
+| Queue 소진 | 마지막 Upload 응답부터 마지막 `INDEXED`까지 |
+| Queue 대기 | Job `created_at`부터 `LOCKED` Event까지 |
+| 실제 처리 | `LOCKED`부터 `INDEXED` Event까지 |
+| 전체 Job 지연 | Job `created_at`부터 `INDEXED` Event까지 |
+
+지연 분포는 선형 보간 p50·p95·p99와 max를 밀리초로 기록한다. 모든 Profile 원본과 중앙값 요약은
+JSON으로 기록한다.
+
+## 5. 정합성 계약
+
+각 Profile은 측정 뒤 다음 조건을 모두 검증한다.
+
+- 대상 Job 전부 `INDEXED`
+- 전체 Schema에 `PENDING` 또는 `PROCESSING` 잔여 Job 없음
+- Job별 `SUCCESS` Attempt 정확히 1개
+- Job별 `LOCKED`, `PARSE_STARTED`, `CHUNKED`, `EMBEDDING_STARTED`, `INDEXED` 순서 유지
+- Document와 Document Version 모두 `INDEXED`
+- `documents.current_version_id`가 측정 Version을 가리킴
+- 문서별 Chunk 수와 Embedding 수 일치
+- 같은 `(chunk_id, embedding_model_id)` 중복 없음
+- 모든 Vector 차원 1024
+- 실패 Attempt, Retry와 최종 실패 Event 없음
+
+정합성 실패가 하나라도 발생하면 성능 숫자를 유효한 결과로 취급하지 않고 Test를 실패시킨다.
+
+## 6. 환경 검증과 결과 보존
+
+Benchmark 시작 전에 다음을 확인한다.
+
+- PostgreSQL Server Version `17.x`
+- pgvector Extension `0.8.1`
+- 전용 Test Schema와 MinIO Bucket
+- BGE Health와 Model명 `BAAI/bge-m3`
+- Document Batch Size와 Worker 최대 동시성
+
+접속 Secret, JWT와 Object Storage Credential은 결과에 기록하지 않는다. 원본 JSON은 다음 Git 제외
+경로에 생성한다.
+
+```text
+build/reports/worker-indexing-throughput/worker-indexing-throughput.json
+```
+
+실행 환경, 중앙값, 해석과 한계만 `docs/test-results/`에 기록한다.
+
+## 7. 실행 경계
+
+일반 `./gradlew test`는 Docker와 실제 BGE-M3에 의존하지 않는다. 전용 Task만 실제 Infrastructure를
+요구한다.
+
+```bash
+docker compose up -d postgres minio embedding-server
+./gradlew workerIndexingThroughputTest
+```
+
+확장 측정 예시는 다음과 같다.
+
+```bash
+./gradlew workerIndexingThroughputTest \
+ -Dworker.indexing.throughput.document-counts=32,64 \
+ -Dworker.indexing.throughput.repetitions=5
+```
+
+## 8. 실패 정책
+
+- Infrastructure Health, DB Version, pgvector Version 또는 BGE Model 계약이 다르면 즉시 실패한다.
+- Upload, Worker 처리, Embedding 또는 상태 Polling이 제한 시간을 넘으면 Job Snapshot을 포함해 실패한다.
+- 완료된 Profile 결과는 후속 진단을 위해 JSON에 보존하되 실패 Profile은 중앙값 판단에서 제외한다.
+- Profile 간 데이터는 같은 전용 Schema에서 고유 식별자로 격리하고 마지막에 Schema와 Bucket을 정리한다.
+
+## 9. 검증
+
+- Percentile, 중앙값과 처리량 계산 단위 테스트
+- 실제 Infrastructure를 사용하는 작은 Smoke Benchmark
+- 기본 Profile 전체 Benchmark
+- 기존 로컬 문서 인덱싱 E2E 회귀
+- 전체 일반 Java 테스트
+
+## 10. 커밋 분할
+
+1. `docs: #133 자동 Worker 처리량 Benchmark 설계 추가`
+2. `test: #133 자동 Worker 전체 인덱싱 처리량 Benchmark 추가`
+3. `build: #133 Worker 처리량 전용 실행 경계 추가`
+4. `perf: #133 자동 Worker 인덱싱 처리량 실측 결과 기록`
+
+## 11. 완료 조건
+
+- 실제 자동 Worker 전체 Pipeline을 한 명령으로 반복 측정할 수 있다.
+- 문서·Chunk·Embedding 처리량과 Queue·처리·전체 지연이 구조화돼 기록된다.
+- 성능 Profile마다 완료·소유권·Attempt·Event·Vector 불변식을 검증한다.
+- 일반 테스트는 실제 Infrastructure 없이 계속 실행된다.
+- 측정 환경의 한계와 운영 수치로 해석하면 안 되는 범위를 결과 문서에 명시한다.
diff --git a/docs/test-results/gimin-#133-worker-indexing-throughput-benchmark.md b/docs/test-results/gimin-#133-worker-indexing-throughput-benchmark.md
new file mode 100644
index 0000000..fee0bae
--- /dev/null
+++ b/docs/test-results/gimin-#133-worker-indexing-throughput-benchmark.md
@@ -0,0 +1,164 @@
+# Issue #133 자동 Worker 전체 문서 인덱싱 처리량 Benchmark 결과
+
+## 1. 결과 요약
+
+PostgreSQL 17.8, pgvector 0.8.1, MinIO와 실제 `BAAI/bge-m3`를 연결하고 자동 Polling Worker가
+문서 업로드부터 `INDEXED` 전환까지 처리하는 전체 경로를 측정했다. 4개 문서 예열 후 16개와 32개
+문서 Profile을 각각 3회 실행했다.
+
+| 문서 수 | 반복 | 중앙 총 시간 | 문서/초 | 문서/분 | Chunk·Embedding/초 | 전체 P95 |
+|---:|---:|---:|---:|---:|---:|---:|
+| 16 | 3회 | 50.255초 | 0.318 | 19.102 | 2.547 | 50.127초 |
+| 32 | 3회 | 104.335초 | 0.307 | 18.402 | 2.454 | 100.665초 |
+
+Queue를 16개에서 32개로 두 배 늘렸을 때 분당 처리량은 약 3.7% 감소했다. Worker 실행 슬롯을
+2개로 고정했기 때문에 처리 시간 P95는 6.768초에서 7.285초로 비슷하게 유지됐고, 전체 지연 증가는
+주로 Queue 대기 P95가 43.794초에서 93.554초로 늘어난 데서 발생했다.
+
+## 2. 공개 가능한 실행 환경
+
+| 항목 | 값 |
+|---|---|
+| Database | PostgreSQL 17.8, Local Docker |
+| pgvector | 0.8.1 |
+| Object Storage | MinIO, Local Docker |
+| Embedding Provider | `BAAI/bge-m3`, Local Docker CPU 추론 |
+| Vector 차원 | 1024 |
+| Embedding Batch Size | 32 |
+| Worker 동시 실행 슬롯 | 2 |
+| Worker Polling 주기 | 50 ms |
+| 업로더 Thread | 4 |
+| 문서 크기 | TXT 6,400자 |
+| 문서당 Chunk·Embedding | 각각 8개 |
+| Warm-up | 4개 문서 |
+| 본 측정 | 16 / 32개 문서, Profile당 3회 |
+| Application | Spring Boot 3.5.16, Java 17 |
+| 실행 장비 | macOS `aarch64`, 가용 Processor 10개 |
+| 실행 일자 | 2026-08-10 KST |
+
+DB Host·Database 이름·Username·Password, MinIO Credential과 JWT 값은 결과에 기록하지 않았다.
+이번 결과는 로컬 개발 장비의 기준선이며 공식 OpenSQL 원격 Server 성능이나 운영 SLO가 아니다.
+
+## 3. 측정 범위
+
+각 문서는 다음 전체 경로를 통과했다.
+
+```text
+TXT 업로드
+→ MinIO 저장
+→ Embedding Job PENDING
+→ 자동 Worker Polling·Claim
+→ Attempt 시작
+→ 텍스트 Parsing·Chunk 저장
+→ 실제 BGE-M3 Batch 호출
+→ vector(1024) 저장
+→ Version INDEXED
+→ Document current_version 전환
+→ Worker 실행 슬롯 반환
+```
+
+Profile마다 Benchmark 전용 Schema와 Bucket의 데이터를 초기화해 이전 실행의 Job·Chunk·Embedding이
+다음 실행의 수치에 포함되지 않게 했다. Worker Node는 Application 수명주기를 유지하기 위해 Profile
+사이에서 재사용했다.
+
+## 4. 실행 방법
+
+PostgreSQL, MinIO와 Embedding Server가 모두 건강 상태인 로컬 환경에서 실행했다.
+
+```bash
+docker compose up -d postgres minio embedding-server
+DB_SSLMODE=disable ./gradlew workerIndexingThroughputTest
+```
+
+구조화 원시 결과는 Git에 포함하지 않는 다음 경로에 생성된다.
+
+```text
+build/reports/worker-indexing-throughput/worker-indexing-throughput.json
+```
+
+실환경 연결과 결과 계약을 빠르게 확인할 때는 다음 Smoke Profile을 사용했다.
+
+```bash
+DB_SSLMODE=disable ./gradlew workerIndexingThroughputTest \
+ -Dworker.indexing.throughput.warm-up-documents=2 \
+ -Dworker.indexing.throughput.document-counts=4 \
+ -Dworker.indexing.throughput.repetitions=1 \
+ -Dworker.indexing.throughput.output=build/reports/worker-indexing-throughput/smoke.json
+```
+
+Smoke 결과는 4개 문서, 32개 Chunk, 32개 Embedding을 12.461초에 처리했고 분당 19.260문서를
+기록했다.
+
+## 5. 반복별 원시 결과
+
+| 문서 수 | 회차 | 총 시간 | 문서/분 | Chunk·Embedding/초 | Queue P95 | 처리 P95 | 전체 P95 |
+|---:|---:|---:|---:|---:|---:|---:|---:|
+| 16 | 1 | 49.379초 | 19.442 | 2.592 | 42.914초 | 7.121초 | 49.134초 |
+| 16 | 2 | 51.898초 | 18.498 | 2.466 | 44.887초 | 6.768초 | 51.742초 |
+| 16 | 3 | 50.255초 | 19.102 | 2.547 | 43.794초 | 6.468초 | 50.127초 |
+| 32 | 1 | 102.558초 | 18.721 | 2.496 | 92.526초 | 6.974초 | 98.769초 |
+| 32 | 2 | 104.335초 | 18.402 | 2.454 | 93.554초 | 7.285초 | 100.665초 |
+| 32 | 3 | 105.928초 | 18.125 | 2.417 | 95.968초 | 7.431초 | 102.253초 |
+
+같은 Profile 세 번의 분당 처리량 범위는 16문서에서 18.498~19.442, 32문서에서
+18.125~18.721이었다. 한 번의 최고값이 아니라 Profile별 중앙값을 비교 기준으로 사용했다.
+
+## 6. 데이터 완전성 검증
+
+| 문서 수 | 회차 | Chunk 수 | Embedding 수 | 결과 |
+|---:|---:|---:|---:|---|
+| 16 | 1~3 | 매회 128 | 매회 128 | PASS |
+| 32 | 1~3 | 매회 256 | 매회 256 | PASS |
+
+각 Profile 완료 시 다음 불변식을 함께 확인했다.
+
+- 모든 Embedding Job이 `INDEXED`다.
+- 각 Job에 성공한 Attempt가 정확히 하나 존재한다.
+- 모든 Document와 Version이 `INDEXED`이고 `current_version_id`가 측정 Version을 가리킨다.
+- Chunk 수와 Embedding 수가 일치하고 중복 Chunk Embedding이 없다.
+- 저장된 모든 Vector의 차원은 1024이며 NaN·Infinity가 없다.
+- 자동 Worker 실행 슬롯이 0으로 반환되고 Worker가 살아 있다.
+
+따라서 이 결과는 HTTP 접수 시간만 측정한 값이 아니라 Vector 저장과 검색 Version 전환이 완료된 시점까지의
+전체 처리량이다.
+
+## 7. 실행 중 발견하고 해결한 환경 문제
+
+### 7.1 로컬 PostgreSQL SSL 설정
+
+첫 Smoke 시도는 외부 Shell의 SSL 설정이 비-SSL 로컬 PostgreSQL에 적용돼 Flyway 연결 전에 실패했다.
+`DB_SSLMODE=disable`을 명시한 뒤 Migration과 전체 Pipeline이 정상 실행됐다.
+
+### 7.2 CPU BGE-M3 응답 제한
+
+기존 Embedding HTTP 응답 제한은 5초로 고정돼 있었다. 로컬 CPU BGE-M3는 정상적인 Batch 응답에도
+5초를 넘겨 Worker가 `EMBEDDING_PROVIDER_UNAVAILABLE`로 재시도했다. 연결·응답 제한을 환경 설정으로
+분리하고 Benchmark에만 2분 응답 제한을 적용했다. 일반 실행의 기본값은 기존과 같은 5초다.
+
+### 7.3 전체 회귀의 JWT 환경 값
+
+전체 회귀 첫 시도는 `JWT_SECRET` 미설정으로 Spring Context 12건이 연쇄 실패했다. Test 전용 JWT와
+`DB_SSLMODE=disable`을 명시한 재실행에서 720개 Test가 모두 통과했다. 첫 실패는 Source 결함이나
+Benchmark 실패가 아니며 통과 결과로 계산하지 않았다.
+
+## 8. 검증 결과
+
+| 검증 | 결과 |
+|---|---|
+| 통계 계약 단위 테스트 | PASS |
+| 전용 Gradle Task 노출 | PASS |
+| 4문서 실제 BGE-M3 Smoke | PASS, 23초 |
+| 16·32문서 각 3회 본 측정 | PASS, 8분 1초 |
+| 전체 Java 회귀 | PASS, 720 tests, failure/error/skipped 0 |
+| `git diff --check` | PASS |
+
+## 9. 결론과 남은 한계
+
+- 구현됨: 자동 Worker의 실제 문서 인덱싱 처리량과 Queue·처리·전체 지연 분포를 반복 측정할 수 있다.
+- 검증됨: 16문서에서 32문서로 Queue가 두 배가 되어도 분당 처리량 감소는 약 3.7%였다.
+- 검증됨: 여섯 실행 모두 Chunk·Embedding·1024차원 Vector 완전성을 만족했다.
+- 관찰됨: 동시 실행 슬롯 2개에서는 Queue 증가가 처리 시간보다 전체 P95를 지배했다.
+- 한계: 단일 Local Apple Silicon CPU 장비의 결과로, 공식 OpenSQL Server나 GPU BGE-M3 결과가 아니다.
+- 한계: TXT 6,400자 고정 입력이라 PDF·DOCX Parser 비용이나 다양한 문서 길이 분포를 대표하지 않는다.
+- 한계: Worker 수·동시 실행 슬롯·Embedding Batch Size를 바꾼 수평 확장 비교는 이번 범위에 포함하지 않았다.
+- 후속: 같은 Harness로 Worker 수·동시성 변화, Queue 적체와 Backpressure, 장애 주입 Profile을 비교해야 한다.
diff --git a/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java b/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java
index 343c499..6638c41 100644
--- a/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java
+++ b/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java
@@ -9,20 +9,37 @@
import org.springframework.http.client.JdkClientHttpRequestFactory;
import org.springframework.web.client.RestClient;
+/**
+ * Embedding Provider HTTP 연결 설정을 구성한다.
+ *
+ *
Provider 호출 계약과 오류 분류는 Domain Client가 담당하고, 이 설정은 연결·응답 제한 시간과
+ * Base URL 같은 Transport 경계만 책임진다.
+ */
@Configuration
public class EmbeddingServerConfig {
@Value("${embedding.server.base-url}")
private String baseUrl;
+ @Value("${embedding.server.connect-timeout:5s}")
+ private Duration connectTimeout;
+
+ @Value("${embedding.server.read-timeout:5s}")
+ private Duration readTimeout;
+
+ /**
+ * 실행 환경별 제한 시간을 적용한 Embedding Provider 전용 Client를 만든다.
+ */
@Bean("embeddingRestClient")
public RestClient embeddingRestClient() {
+ // 1. 연결과 추론 응답 제한을 분리해 느린 CPU 추론 환경에서도 Timeout을 독립적으로 조정한다.
HttpClient httpClient = HttpClient.newBuilder()
- .connectTimeout(Duration.ofSeconds(5))
+ .connectTimeout(connectTimeout)
.build();
JdkClientHttpRequestFactory requestFactory = new JdkClientHttpRequestFactory(httpClient);
- requestFactory.setReadTimeout(Duration.ofSeconds(5));
+ requestFactory.setReadTimeout(readTimeout);
+ // 2. Domain Client가 같은 Transport 정책을 공유하도록 단일 이름의 RestClient를 제공한다.
return RestClient.builder()
.baseUrl(baseUrl)
.requestFactory(requestFactory)
diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml
index 3821d5a..4e983c2 100644
--- a/src/main/resources/application.yml
+++ b/src/main/resources/application.yml
@@ -61,6 +61,8 @@ jwt:
embedding:
server:
base-url: ${EMBEDDING_SERVER_URL:http://localhost:8000}
+ connect-timeout: ${EMBEDDING_SERVER_CONNECT_TIMEOUT:5s}
+ read-timeout: ${EMBEDDING_SERVER_READ_TIMEOUT:5s}
document:
# 실제 BGE-M3 CPU Benchmark에서 최고 처리량의 98.3%를 유지하면서 Batch 64보다 p95가 47% 낮은 균형점이다.
batch-size: ${EMBEDDING_DOCUMENT_BATCH_SIZE:32}
diff --git a/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java b/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java
index 9a8db8d..3b5ccfa 100644
--- a/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java
+++ b/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java
@@ -56,6 +56,17 @@ List objectKeys() throws Exception {
return List.copyOf(keys);
}
+ /**
+ * Bucket은 유지하고 현재 Test가 저장한 Object만 제거해 반복 Profile의 입력 상태를 초기화한다.
+ */
+ void clear() throws Exception {
+ for (String objectKey : objectKeys()) {
+ minioClient.removeObject(
+ RemoveObjectArgs.builder().bucket(bucket).object(objectKey).build()
+ );
+ }
+ }
+
@Override
public void close() throws Exception {
boolean exists = minioClient.bucketExists(
@@ -66,11 +77,7 @@ public void close() throws Exception {
}
// 1. MinIO는 비어 있지 않은 Bucket 삭제를 거부하므로 실제 Object를 모두 먼저 제거한다.
- for (String objectKey : objectKeys()) {
- minioClient.removeObject(
- RemoveObjectArgs.builder().bucket(bucket).object(objectKey).build()
- );
- }
+ clear();
// 2. Test가 생성한 격리 Bucket만 삭제하고 다른 개발 Bucket은 조회하거나 변경하지 않는다.
minioClient.removeBucket(RemoveBucketArgs.builder().bucket(bucket).build());
diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputBenchmark.java b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputBenchmark.java
new file mode 100644
index 0000000..c7f5455
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputBenchmark.java
@@ -0,0 +1,721 @@
+package com.opensource.docgrid.e2e;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Duration;
+import java.time.LocalDateTime;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.Timeout;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.SpringBootVersion;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.boot.test.web.client.TestRestTemplate;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.test.annotation.DirtiesContext;
+import org.springframework.test.context.ActiveProfiles;
+import org.springframework.test.context.DynamicPropertyRegistry;
+import org.springframework.test.context.DynamicPropertySource;
+import org.springframework.web.client.RestClient;
+
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.opensource.docgrid.domain.embedding.config.EmbeddingBatchProperties;
+import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties;
+import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool;
+import com.opensource.docgrid.domain.worker.lifecycle.WorkerExecutionLifecycleManager;
+import com.opensource.docgrid.domain.worker.lifecycle.WorkerJobPollingScheduler;
+import com.opensource.docgrid.domain.worker.lifecycle.WorkerLifecycleManager;
+import com.opensource.docgrid.e2e.LocalE2eApiClient.UploadedDocument;
+import com.opensource.docgrid.e2e.LocalE2eDocumentFactory.DocumentPayload;
+
+import io.minio.MinioClient;
+import lombok.extern.slf4j.Slf4j;
+
+/**
+ * 실제 PostgreSQL 17·pgvector·MinIO·BGE-M3에서 자동 Worker 전체 인덱싱 처리량을 측정한다.
+ *
+ * 결정적 TXT 문서를 실제 HTTP로 동시에 접수하고 Production Scheduler가 Queue를 소진하게 한다.
+ * 각 Profile은 처리량과 단계별 지연을 기록한 뒤 Job, Attempt, Event, Chunk와 Vector 정합성을 함께
+ * 검증한다. 실제 Infrastructure를 요구하므로 일반 회귀 테스트와 분리된 Tag에서만 실행한다.
+ */
+@Slf4j
+@Tag("worker-indexing-throughput")
+@ActiveProfiles({"test", "minio-integration"})
+@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
+@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@DisplayName("자동 Worker 전체 문서 인덱싱 처리량 Benchmark")
+class WorkerIndexingThroughputBenchmark {
+
+ private static final String EXECUTION_ID = UUID.randomUUID().toString().replace("-", "");
+ // PostgreSQL 식별자 63자 제한 안에서 별도 Gradle 실행이 Schema를 공유하지 않도록 격리한다.
+ private static final String TEST_SCHEMA = "docgrid_worker_indexing_throughput_"
+ + EXECUTION_ID.substring(0, 24);
+ private static final String TEST_BUCKET = "docgrid-worker-throughput-" + EXECUTION_ID;
+ private static final String EXPECTED_POSTGRES_VERSION_PREFIX = "17.";
+ private static final String EXPECTED_PGVECTOR_VERSION = "0.8.1";
+ private static final String EXPECTED_MODEL = "BAAI/bge-m3";
+ private static final int EXPECTED_VECTOR_DIMENSION = 1024;
+ private static final int DOCUMENT_CHARACTER_COUNT = positiveIntegerProperty(
+ "worker.indexing.throughput.document-characters",
+ 6_400
+ );
+ private static final int WARM_UP_DOCUMENT_COUNT = positiveIntegerProperty(
+ "worker.indexing.throughput.warm-up-documents",
+ 4
+ );
+ private static final List DOCUMENT_COUNTS = positiveIntegerListProperty(
+ "worker.indexing.throughput.document-counts",
+ List.of(16, 32)
+ );
+ private static final int REPETITIONS = positiveIntegerProperty(
+ "worker.indexing.throughput.repetitions",
+ 3
+ );
+ private static final int UPLOADER_THREADS = positiveIntegerProperty(
+ "worker.indexing.throughput.uploader-threads",
+ 4
+ );
+ private static final long PROFILE_TIMEOUT_SECONDS = positiveLongProperty(
+ "worker.indexing.throughput.profile-timeout-seconds",
+ 600L
+ );
+ private static final long POLLING_SLEEP_MILLIS = positiveLongProperty(
+ "worker.indexing.throughput.status-polling-ms",
+ 100L
+ );
+ private static final Path OUTPUT_PATH = Path.of(System.getProperty(
+ "worker.indexing.throughput.output",
+ "build/reports/worker-indexing-throughput/worker-indexing-throughput.json"
+ ));
+
+ @Autowired private TestRestTemplate restTemplate;
+ @Autowired private JdbcTemplate jdbcTemplate;
+ @Autowired private MinioClient minioClient;
+ @Autowired private ObjectMapper objectMapper;
+ @Autowired private WorkerLifecycleManager workerLifecycleManager;
+ @Autowired private WorkerJobPollingScheduler pollingScheduler;
+ @Autowired private WorkerExecutionLifecycleManager executionLifecycleManager;
+ @Autowired private WorkerExecutionSlotPool executionSlotPool;
+ @Autowired private IndexingWorkerProperties workerProperties;
+ @Autowired private EmbeddingBatchProperties embeddingBatchProperties;
+ @Autowired @Qualifier("embeddingRestClient") private RestClient embeddingRestClient;
+
+ private LocalE2eApiClient apiClient;
+ private LocalE2eMinioBucket minioBucket;
+ private String accessToken;
+
+ @DynamicPropertySource
+ static void configureEnvironment(DynamicPropertyRegistry registry) {
+ registry.add("TEST_DB_SCHEMA", () -> TEST_SCHEMA);
+ registry.add("jwt.secret", () -> "docgrid-worker-indexing-throughput-test-secret-key-2026");
+ registry.add("minio.bucket", () -> TEST_BUCKET);
+ registry.add("indexing.worker.enabled", () -> "true");
+ registry.add("indexing.worker.name", () -> "worker-indexing-throughput-benchmark");
+ registry.add("indexing.worker.polling-interval", () -> "50ms");
+ registry.add("indexing.worker.heartbeat-interval", () -> "1s");
+ registry.add("indexing.worker.dead-threshold", () -> "2m");
+ registry.add("indexing.worker.max-concurrency", () -> "2");
+ registry.add("indexing.worker.lease-duration", () -> "2m");
+ registry.add("indexing.worker.lease-renewal-interval", () -> "10s");
+ registry.add("indexing.worker.lease-recovery-interval", () -> "10m");
+ registry.add("indexing.worker.shutdown-grace-period", () -> "30s");
+ // CPU 기반 BGE-M3의 실측 추론 시간을 5초 기본 운영 Timeout과 분리한다.
+ registry.add("embedding.server.read-timeout", () -> "2m");
+ }
+
+ @BeforeAll
+ void setUpInfrastructure() throws Exception {
+ apiClient = new LocalE2eApiClient(restTemplate);
+ minioBucket = new LocalE2eMinioBucket(minioClient, TEST_BUCKET);
+ minioBucket.create();
+ accessToken = apiClient.loginAdmin();
+ awaitCondition(
+ "Worker가 Application Ready 뒤 등록되지 않았습니다.",
+ () -> workerLifecycleManager.getWorkerId().isPresent()
+ );
+ }
+
+ @AfterAll
+ void cleanUpInfrastructure() throws Exception {
+ // 1. Scheduler와 실행 Thread를 먼저 닫아 Profile 자원 정리 뒤 DB 접근이 재개되지 않게 한다.
+ pollingScheduler.stopPolling();
+ executionLifecycleManager.shutdown();
+ workerLifecycleManager.stopWorker();
+
+ // 2. Benchmark 전용 Bucket과 Schema만 제거해 기존 로컬 개발 Data를 보존한다.
+ minioBucket.close();
+ jdbcTemplate.execute("DROP SCHEMA IF EXISTS " + TEST_SCHEMA + " CASCADE");
+ }
+
+ @Test
+ @Timeout(1_800)
+ @DisplayName("실제 자동 Worker의 처리량과 Queue·처리·전체 지연을 반복 측정한다")
+ void measureAutomaticWorkerIndexingThroughput() throws Exception {
+ // 1. 잘못된 DB나 외부 Service에서 얻은 수치를 기록하지 않도록 환경 계약을 먼저 확인한다.
+ EnvironmentFingerprint environment = validateEnvironment();
+ List runs = new ArrayList<>();
+ writeReport(new BenchmarkReport(environment, runs, List.of()));
+ logJson("WORKER_INDEXING_THROUGHPUT_ENV", environment);
+
+ // 2. Migration, HTTP, MinIO, Model과 주요 DB 경로를 예열하되 측정 결과에는 포함하지 않는다.
+ runWarmUp();
+
+ // 3. 같은 문서 분포로 Queue 크기와 반복만 바꿔 Profile 원시 결과를 수집한다.
+ for (int documentCount : DOCUMENT_COUNTS) {
+ for (int repetition = 1; repetition <= REPETITIONS; repetition++) {
+ resetProfileState();
+ ProfileRun run = runMeasuredProfile(documentCount, repetition);
+ runs.add(run);
+ logJson("WORKER_INDEXING_THROUGHPUT_RESULT", run);
+ writeReport(new BenchmarkReport(environment, runs, medians(runs)));
+ }
+ }
+
+ // 4. 장비 내 변동을 줄인 Profile별 중앙값을 최종 결과와 Log에 별도로 남긴다.
+ List medians = medians(runs);
+ for (ProfileMedian median : medians) {
+ logJson("WORKER_INDEXING_THROUGHPUT_MEDIAN", median);
+ }
+ writeReport(new BenchmarkReport(environment, runs, medians));
+ }
+
+ private EnvironmentFingerprint validateEnvironment() {
+ String postgresVersion = jdbcTemplate.queryForObject("SHOW server_version", String.class);
+ String pgvectorVersion = jdbcTemplate.queryForObject(
+ "SELECT extversion FROM pg_extension WHERE extname = 'vector'",
+ String.class
+ );
+ Map model = jdbcTemplate.queryForMap(
+ "SELECT model_name, dimension FROM embedding_models "
+ + "WHERE is_active = TRUE AND is_searchable = TRUE"
+ );
+ JsonNode health = embeddingRestClient.get()
+ .uri("/health")
+ .retrieve()
+ .body(JsonNode.class);
+
+ assertThat(jdbcTemplate.queryForObject("SELECT current_schema()", String.class))
+ .isEqualTo(TEST_SCHEMA);
+ assertThat(postgresVersion).startsWith(EXPECTED_POSTGRES_VERSION_PREFIX);
+ assertThat(pgvectorVersion).isEqualTo(EXPECTED_PGVECTOR_VERSION);
+ assertThat(model.get("model_name").toString()).isEqualTo(EXPECTED_MODEL);
+ assertThat(((Number) model.get("dimension")).intValue()).isEqualTo(EXPECTED_VECTOR_DIMENSION);
+ assertThat(health).isNotNull();
+ assertThat(health.path("status").asText()).isEqualTo("ok");
+ assertThat(workerProperties.getMaxConcurrency()).isEqualTo(2);
+
+ return new EnvironmentFingerprint(
+ postgresVersion,
+ pgvectorVersion,
+ SpringBootVersion.getVersion(),
+ EXPECTED_MODEL,
+ EXPECTED_VECTOR_DIMENSION,
+ embeddingBatchProperties.getBatchSize(),
+ workerProperties.getMaxConcurrency(),
+ workerProperties.getPollingInterval().toString(),
+ UPLOADER_THREADS,
+ WARM_UP_DOCUMENT_COUNT,
+ DOCUMENT_COUNTS,
+ REPETITIONS,
+ DOCUMENT_CHARACTER_COUNT,
+ System.getProperty("os.name"),
+ System.getProperty("os.arch"),
+ Runtime.getRuntime().availableProcessors()
+ );
+ }
+
+ private void runWarmUp() throws Exception {
+ resetProfileState();
+ List uploads = uploadDocuments("warm-up", WARM_UP_DOCUMENT_COUNT);
+ awaitIndexedAndIdle(uploads, "Warm-up");
+ assertProfileInvariants(uploads);
+ log.info(
+ "자동 Worker 처리량 Benchmark 예열을 완료했습니다. documentCount={}",
+ WARM_UP_DOCUMENT_COUNT
+ );
+ }
+
+ private ProfileRun runMeasuredProfile(int documentCount, int repetition) throws Exception {
+ String profileName = "documents-" + documentCount + "-run-" + repetition;
+ long profileStartedAt = System.nanoTime();
+ List uploads = uploadDocuments(profileName, documentCount);
+ long uploadCompletedAt = System.nanoTime();
+
+ awaitIndexedAndIdle(uploads, profileName);
+ long profileCompletedAt = System.nanoTime();
+ ProfileData profileData = assertProfileInvariants(uploads);
+
+ double elapsedSeconds = seconds(profileCompletedAt - profileStartedAt);
+ double uploadSeconds = seconds(uploadCompletedAt - profileStartedAt);
+ double queueDrainSeconds = seconds(profileCompletedAt - uploadCompletedAt);
+ return new ProfileRun(
+ documentCount,
+ repetition,
+ profileData.chunkCount(),
+ profileData.embeddingCount(),
+ elapsedSeconds,
+ uploadSeconds,
+ queueDrainSeconds,
+ documentCount / elapsedSeconds,
+ documentCount / elapsedSeconds * 60.0,
+ profileData.chunkCount() / elapsedSeconds,
+ profileData.embeddingCount() / elapsedSeconds,
+ LatencySummary.from(profileData.timings().stream().map(JobTiming::queueWaitMillis).toList()),
+ LatencySummary.from(profileData.timings().stream().map(JobTiming::processingMillis).toList()),
+ LatencySummary.from(profileData.timings().stream().map(JobTiming::endToEndMillis).toList())
+ );
+ }
+
+ private List uploadDocuments(String profileName, int documentCount) throws Exception {
+ int threadCount = Math.min(UPLOADER_THREADS, documentCount);
+ ExecutorService uploader = Executors.newFixedThreadPool(threadCount);
+ List> futures = new ArrayList<>(documentCount);
+
+ try {
+ // 1. Upload를 병렬 제출해 Worker 슬롯이 유지될 만큼 빠르게 PENDING Queue를 만든다.
+ for (int index = 0; index < documentCount; index++) {
+ int documentIndex = index;
+ futures.add(uploader.submit(() -> apiClient.upload(
+ accessToken,
+ throughputDocument(profileName, documentIndex)
+ )));
+ }
+ uploader.shutdown();
+
+ // 2. 각 HTTP 응답을 제한 시간 안에서 수집하고 원래 문서 순서를 보존한다.
+ List uploads = new ArrayList<>(documentCount);
+ for (Future future : futures) {
+ uploads.add(future.get(PROFILE_TIMEOUT_SECONDS, TimeUnit.SECONDS));
+ }
+ return List.copyOf(uploads);
+ } finally {
+ uploader.shutdownNow();
+ uploader.awaitTermination(10, TimeUnit.SECONDS);
+ }
+ }
+
+ private DocumentPayload throughputDocument(String profileName, int documentIndex) {
+ String marker = String.format(
+ Locale.ROOT,
+ "DocGrid automatic worker throughput document %04d. ",
+ documentIndex
+ );
+ String sentence = "Automatic workers claim pending jobs, parse deterministic text, "
+ + "store chunks, generate BGE-M3 vectors, and switch the searchable version atomically. ";
+ StringBuilder body = new StringBuilder(DOCUMENT_CHARACTER_COUNT);
+ body.append(marker);
+ while (body.length() < DOCUMENT_CHARACTER_COUNT) {
+ body.append(sentence);
+ }
+ body.setLength(DOCUMENT_CHARACTER_COUNT);
+
+ String fileName = profileName + "-" + String.format(Locale.ROOT, "%04d", documentIndex) + ".txt";
+ return LocalE2eDocumentFactory.text(fileName, "Throughput " + fileName, body.toString());
+ }
+
+ private void awaitIndexedAndIdle(List uploads, String profileName)
+ throws InterruptedException {
+ awaitCondition(
+ profileName + " Profile이 완료되지 않았습니다. " + jobSnapshot(uploads),
+ () -> indexedJobCount(uploads) == uploads.size() && executionSlotPool.getActiveSlots() == 0
+ );
+ }
+
+ private ProfileData assertProfileInvariants(List uploads) {
+ List timings = new ArrayList<>(uploads.size());
+ List chunkCounts = new ArrayList<>(uploads.size());
+ int totalEmbeddings = 0;
+
+ for (UploadedDocument upload : uploads) {
+ // 1. Job과 검색 Version 전이가 끝났고 재시도 없이 한 Attempt만 성공했는지 확인한다.
+ assertThat(queryString("SELECT status FROM embedding_jobs WHERE id = ?", upload.embeddingJobId()))
+ .isEqualTo("INDEXED");
+ assertThat(queryInteger(
+ "SELECT retry_count FROM embedding_jobs WHERE id = ?",
+ upload.embeddingJobId()
+ )).isZero();
+ assertThat(count(
+ "SELECT COUNT(*) FROM embedding_job_attempts WHERE embedding_job_id = ?",
+ upload.embeddingJobId()
+ )).isOne();
+ assertThat(count(
+ "SELECT COUNT(*) FROM embedding_job_attempts "
+ + "WHERE embedding_job_id = ? AND status = 'SUCCESS'",
+ upload.embeddingJobId()
+ )).isOne();
+ assertThat(queryString("SELECT status FROM documents WHERE id = ?", upload.documentId()))
+ .isEqualTo("INDEXED");
+ assertThat(queryLong("SELECT current_version_id FROM documents WHERE id = ?", upload.documentId()))
+ .isEqualTo(upload.documentVersionId());
+ assertThat(queryString(
+ "SELECT status FROM document_versions WHERE id = ?",
+ upload.documentVersionId()
+ )).isEqualTo("INDEXED");
+
+ // 2. Chunk와 Embedding Set이 일대일이며 모든 Vector가 1024차원인지 확인한다.
+ int chunkCount = count(
+ "SELECT COUNT(*) FROM document_chunks WHERE document_version_id = ?",
+ upload.documentVersionId()
+ );
+ int embeddingCount = count(
+ "SELECT COUNT(*) FROM embeddings WHERE document_version_id = ?",
+ upload.documentVersionId()
+ );
+ assertThat(chunkCount).isPositive();
+ assertThat(embeddingCount).isEqualTo(chunkCount);
+ assertThat(count(
+ "SELECT COUNT(DISTINCT chunk_id) FROM embeddings WHERE document_version_id = ?",
+ upload.documentVersionId()
+ )).isEqualTo(embeddingCount);
+ assertThat(queryInteger(
+ "SELECT MIN(vector_dims(vector)) FROM embeddings WHERE document_version_id = ?",
+ upload.documentVersionId()
+ )).isEqualTo(EXPECTED_VECTOR_DIMENSION);
+ assertThat(queryInteger(
+ "SELECT MAX(vector_dims(vector)) FROM embeddings WHERE document_version_id = ?",
+ upload.documentVersionId()
+ )).isEqualTo(EXPECTED_VECTOR_DIMENSION);
+ chunkCounts.add(chunkCount);
+ totalEmbeddings += embeddingCount;
+
+ // 3. 정상 Event 순서와 실패·Retry 부재를 확인한 뒤 단계별 시각을 수집한다.
+ List events = jdbcTemplate.queryForList(
+ "SELECT event_type FROM indexing_events WHERE embedding_job_id = ? ORDER BY id",
+ String.class,
+ upload.embeddingJobId()
+ );
+ assertThat(events).containsSubsequence(
+ "LOCKED",
+ "PARSE_STARTED",
+ "CHUNKED",
+ "EMBEDDING_STARTED",
+ "INDEXED"
+ );
+ assertThat(events).doesNotContain(
+ "PARSE_FAILED",
+ "EMBEDDING_FAILED",
+ "LEASE_EXPIRED",
+ "FAILED",
+ "RETRY",
+ "MANUAL_RETRY"
+ );
+ timings.add(readJobTiming(upload.embeddingJobId()));
+ }
+
+ // 4. 같은 Profile의 문서별 Chunk 분포와 전체 Queue 소진 상태를 확인한다.
+ assertThat(chunkCounts).allMatch(count -> count.equals(chunkCounts.get(0)));
+ assertThat(count(
+ "SELECT COUNT(*) FROM embedding_jobs WHERE status IN ('PENDING', 'PROCESSING')"
+ )).isZero();
+ assertThat(workerLifecycleManager.getWorkerId()).isPresent();
+ return new ProfileData(
+ chunkCounts.stream().mapToInt(Integer::intValue).sum(),
+ totalEmbeddings,
+ List.copyOf(timings)
+ );
+ }
+
+ private JobTiming readJobTiming(Long jobId) {
+ return jdbcTemplate.queryForObject(
+ """
+ SELECT job.created_at AS job_created_at,
+ MIN(event.occurred_at) FILTER (WHERE event.event_type = 'LOCKED') AS first_locked_at,
+ MIN(event.occurred_at) FILTER (WHERE event.event_type = 'INDEXED') AS indexed_at
+ FROM embedding_jobs job
+ JOIN indexing_events event ON event.embedding_job_id = job.id
+ WHERE job.id = ?
+ GROUP BY job.id, job.created_at
+ """,
+ (resultSet, rowNumber) -> {
+ LocalDateTime createdAt = resultSet.getTimestamp("job_created_at").toLocalDateTime();
+ LocalDateTime lockedAt = resultSet.getTimestamp("first_locked_at").toLocalDateTime();
+ LocalDateTime indexedAt = resultSet.getTimestamp("indexed_at").toLocalDateTime();
+ double queueWaitMillis = millis(Duration.between(createdAt, lockedAt));
+ double processingMillis = millis(Duration.between(lockedAt, indexedAt));
+ double endToEndMillis = millis(Duration.between(createdAt, indexedAt));
+ assertThat(queueWaitMillis).isGreaterThanOrEqualTo(0.0);
+ assertThat(processingMillis).isGreaterThanOrEqualTo(0.0);
+ assertThat(endToEndMillis).isGreaterThanOrEqualTo(processingMillis);
+ return new JobTiming(queueWaitMillis, processingMillis, endToEndMillis);
+ },
+ jobId
+ );
+ }
+
+ private void resetProfileState() throws Exception {
+ awaitCondition(
+ "이전 Profile의 Worker 실행 슬롯이 반환되지 않았습니다.",
+ () -> executionSlotPool.getActiveSlots() == 0
+ );
+ minioBucket.clear();
+ jdbcTemplate.execute("""
+ TRUNCATE TABLE
+ indexing_events,
+ embedding_job_attempts,
+ embeddings,
+ document_chunks,
+ embedding_jobs,
+ document_versions,
+ documents,
+ file_objects
+ RESTART IDENTITY CASCADE
+ """);
+ }
+
+ private int indexedJobCount(List uploads) {
+ String placeholders = String.join(",", uploads.stream().map(upload -> "?").toList());
+ Object[] jobIds = uploads.stream().map(UploadedDocument::embeddingJobId).toArray();
+ return count(
+ "SELECT COUNT(*) FROM embedding_jobs WHERE status = 'INDEXED' AND id IN (" + placeholders + ")",
+ jobIds
+ );
+ }
+
+ private String jobSnapshot(List uploads) {
+ String placeholders = String.join(",", uploads.stream().map(upload -> "?").toList());
+ Object[] jobIds = uploads.stream().map(UploadedDocument::embeddingJobId).toArray();
+ return jdbcTemplate.queryForList(
+ "SELECT id || ':' || status || ':retry=' || retry_count || ':error=' "
+ + "|| COALESCE(error_code, 'none') FROM embedding_jobs WHERE id IN (" + placeholders + ") "
+ + "ORDER BY id",
+ String.class,
+ jobIds
+ ).toString();
+ }
+
+ private List medians(List runs) {
+ List medians = new ArrayList<>();
+ for (int documentCount : DOCUMENT_COUNTS) {
+ List profileRuns = runs.stream()
+ .filter(run -> run.documentCount() == documentCount)
+ .toList();
+ if (profileRuns.isEmpty()) {
+ continue;
+ }
+ medians.add(new ProfileMedian(
+ documentCount,
+ profileRuns.size(),
+ median(profileRuns, ProfileRun::documentsPerSecond),
+ median(profileRuns, ProfileRun::documentsPerMinute),
+ median(profileRuns, ProfileRun::chunksPerSecond),
+ median(profileRuns, ProfileRun::embeddingsPerSecond),
+ median(profileRuns, ProfileRun::totalElapsedSeconds),
+ median(profileRuns, ProfileRun::uploadElapsedSeconds),
+ median(profileRuns, ProfileRun::queueDrainSeconds),
+ median(profileRuns, run -> run.queueWait().p95Millis()),
+ median(profileRuns, run -> run.processing().p95Millis()),
+ median(profileRuns, run -> run.endToEnd().p95Millis())
+ ));
+ }
+ return List.copyOf(medians);
+ }
+
+ private double median(List runs, ProfileMetric metric) {
+ return WorkerIndexingThroughputStatistics.median(runs.stream().map(metric::value).toList());
+ }
+
+ private void writeReport(BenchmarkReport report) throws IOException {
+ Path parent = OUTPUT_PATH.toAbsolutePath().getParent();
+ if (parent != null) {
+ Files.createDirectories(parent);
+ }
+ objectMapper.writerWithDefaultPrettyPrinter().writeValue(OUTPUT_PATH.toFile(), report);
+ }
+
+ private void logJson(String prefix, Object value) throws JsonProcessingException {
+ log.info("{} {}", prefix, objectMapper.writeValueAsString(value));
+ }
+
+ private void awaitCondition(String failureMessage, CheckedCondition condition) throws InterruptedException {
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(PROFILE_TIMEOUT_SECONDS);
+ while (System.nanoTime() < deadline) {
+ if (condition.evaluate()) {
+ return;
+ }
+ Thread.sleep(POLLING_SLEEP_MILLIS);
+ }
+ throw new AssertionError(failureMessage);
+ }
+
+ private int count(String sql, Object... arguments) {
+ return jdbcTemplate.queryForObject(sql, Integer.class, arguments);
+ }
+
+ private int queryInteger(String sql, Object... arguments) {
+ return jdbcTemplate.queryForObject(sql, Integer.class, arguments);
+ }
+
+ private Long queryLong(String sql, Object... arguments) {
+ return jdbcTemplate.queryForObject(sql, Long.class, arguments);
+ }
+
+ private String queryString(String sql, Object... arguments) {
+ return jdbcTemplate.queryForObject(sql, String.class, arguments);
+ }
+
+ private static double seconds(long nanoseconds) {
+ return nanoseconds / 1_000_000_000.0;
+ }
+
+ private static double millis(Duration duration) {
+ return duration.toNanos() / 1_000_000.0;
+ }
+
+ private static int positiveIntegerProperty(String name, int defaultValue) {
+ String value = System.getProperty(name);
+ int parsed = value == null ? defaultValue : Integer.parseInt(value.trim());
+ if (parsed < 1) {
+ throw new IllegalArgumentException(name + "은 1 이상이어야 합니다.");
+ }
+ return parsed;
+ }
+
+ private static long positiveLongProperty(String name, long defaultValue) {
+ String value = System.getProperty(name);
+ long parsed = value == null ? defaultValue : Long.parseLong(value.trim());
+ if (parsed < 1L) {
+ throw new IllegalArgumentException(name + "은 1 이상이어야 합니다.");
+ }
+ return parsed;
+ }
+
+ private static List positiveIntegerListProperty(String name, List defaults) {
+ String value = System.getProperty(name);
+ if (value == null || value.isBlank()) {
+ return defaults;
+ }
+ List parsed = List.of(value.split(",")).stream()
+ .map(String::trim)
+ .map(Integer::parseInt)
+ .distinct()
+ .toList();
+ if (parsed.isEmpty() || parsed.stream().anyMatch(item -> item < 1)) {
+ throw new IllegalArgumentException(name + "에는 1 이상의 정수만 사용할 수 있습니다.");
+ }
+ return parsed;
+ }
+
+ /** 실행 환경과 Workload 설정을 Secret 없이 재현할 수 있는 지문이다. */
+ private record EnvironmentFingerprint(
+ String postgresVersion,
+ String pgvectorVersion,
+ String springBootVersion,
+ String embeddingModel,
+ int vectorDimension,
+ int embeddingBatchSize,
+ int workerMaxConcurrency,
+ String workerPollingInterval,
+ int uploaderThreads,
+ int warmUpDocumentCount,
+ List documentCounts,
+ int repetitions,
+ int documentCharacterCount,
+ String osName,
+ String osArchitecture,
+ int availableProcessors
+ ) {
+ }
+
+ /** 한 Profile 반복의 처리량과 단계별 지연 원시 결과다. */
+ private record ProfileRun(
+ int documentCount,
+ int repetition,
+ int chunkCount,
+ int embeddingCount,
+ double totalElapsedSeconds,
+ double uploadElapsedSeconds,
+ double queueDrainSeconds,
+ double documentsPerSecond,
+ double documentsPerMinute,
+ double chunksPerSecond,
+ double embeddingsPerSecond,
+ LatencySummary queueWait,
+ LatencySummary processing,
+ LatencySummary endToEnd
+ ) {
+ }
+
+ /** 같은 문서 수 Profile의 반복 결과를 중앙값으로 요약한 비교 단위다. */
+ private record ProfileMedian(
+ int documentCount,
+ int completedRepetitions,
+ double documentsPerSecond,
+ double documentsPerMinute,
+ double chunksPerSecond,
+ double embeddingsPerSecond,
+ double totalElapsedSeconds,
+ double uploadElapsedSeconds,
+ double queueDrainSeconds,
+ double queueWaitP95Millis,
+ double processingP95Millis,
+ double endToEndP95Millis
+ ) {
+ }
+
+ /** 한 지연 분포의 중앙값, 꼬리 지연과 최댓값을 밀리초로 보존한다. */
+ private record LatencySummary(
+ double p50Millis,
+ double p95Millis,
+ double p99Millis,
+ double maxMillis
+ ) {
+ private static LatencySummary from(List samples) {
+ return new LatencySummary(
+ WorkerIndexingThroughputStatistics.percentile(samples, 50.0),
+ WorkerIndexingThroughputStatistics.percentile(samples, 95.0),
+ WorkerIndexingThroughputStatistics.percentile(samples, 99.0),
+ samples.stream().mapToDouble(Double::doubleValue).max().orElseThrow()
+ );
+ }
+ }
+
+ /** Profile 검증에서 수집한 전체 Chunk·Embedding 수와 Job 단계별 시각이다. */
+ private record ProfileData(int chunkCount, int embeddingCount, List timings) {
+ }
+
+ /** 한 Job의 Queue, 실제 처리와 전체 지연을 밀리초로 표현한다. */
+ private record JobTiming(double queueWaitMillis, double processingMillis, double endToEndMillis) {
+ }
+
+ /** 완료된 Profile과 중앙값을 환경 지문과 함께 JSON으로 저장하는 최상위 결과다. */
+ private record BenchmarkReport(
+ EnvironmentFingerprint environment,
+ List runs,
+ List medians
+ ) {
+ }
+
+ /** Profile 중앙값을 계산할 실수 지표를 선택한다. */
+ @FunctionalInterface
+ private interface ProfileMetric {
+ double value(ProfileRun run);
+ }
+
+ /** 제한 시간 동안 반복 평가할 DB·Worker 완료 조건이다. */
+ @FunctionalInterface
+ private interface CheckedCondition {
+ boolean evaluate();
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatistics.java b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatistics.java
new file mode 100644
index 0000000..1f2b0cf
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatistics.java
@@ -0,0 +1,47 @@
+package com.opensource.docgrid.e2e;
+
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * Worker 처리량 Benchmark의 Percentile과 중앙값을 결정적으로 계산하는 Test 전용 통계 경계다.
+ *
+ * 측정 I/O와 분리된 순수 계산만 제공해 실제 Infrastructure 없이도 통계 계약을 회귀 검증한다.
+ */
+final class WorkerIndexingThroughputStatistics {
+
+ private WorkerIndexingThroughputStatistics() {
+ }
+
+ static double median(List values) {
+ return percentile(values, 50.0);
+ }
+
+ static double percentile(List values, double percentile) {
+ if (values == null || values.isEmpty()) {
+ throw new IllegalArgumentException("Percentile 입력은 비어 있을 수 없습니다.");
+ }
+ if (!Double.isFinite(percentile) || percentile < 0.0 || percentile > 100.0) {
+ throw new IllegalArgumentException("Percentile은 0 이상 100 이하여야 합니다.");
+ }
+
+ List sorted = new ArrayList<>(values.size());
+ for (Double value : values) {
+ if (value == null || !Double.isFinite(value)) {
+ throw new IllegalArgumentException("Percentile 입력은 유한한 수여야 합니다.");
+ }
+ sorted.add(value);
+ }
+ sorted.sort(Double::compareTo);
+
+ double rank = percentile / 100.0 * (sorted.size() - 1);
+ int lowerIndex = (int) Math.floor(rank);
+ int upperIndex = (int) Math.ceil(rank);
+ if (lowerIndex == upperIndex) {
+ return sorted.get(lowerIndex);
+ }
+
+ double weight = rank - lowerIndex;
+ return sorted.get(lowerIndex) * (1.0 - weight) + sorted.get(upperIndex) * weight;
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatisticsTest.java b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatisticsTest.java
new file mode 100644
index 0000000..30e5b04
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatisticsTest.java
@@ -0,0 +1,41 @@
+package com.opensource.docgrid.e2e;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+import java.util.List;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Worker 처리량 Benchmark가 사용하는 선형 보간 Percentile과 입력 검증 계약을 확인한다.
+ */
+@DisplayName("Worker 인덱싱 처리량 통계 단위 테스트")
+class WorkerIndexingThroughputStatisticsTest {
+
+ @Test
+ @DisplayName("정렬되지 않은 짝수 표본의 중앙값과 p95를 선형 보간한다")
+ void percentile_interpolatesUnsortedEvenSamples() {
+ List samples = List.of(40.0, 10.0, 30.0, 20.0);
+
+ assertThat(WorkerIndexingThroughputStatistics.median(samples)).isEqualTo(25.0);
+ assertThat(WorkerIndexingThroughputStatistics.percentile(samples, 95.0)).isEqualTo(38.5);
+ }
+
+ @Test
+ @DisplayName("단일 표본은 모든 Percentile에서 같은 값을 반환한다")
+ void percentile_returnsOnlySample() {
+ assertThat(WorkerIndexingThroughputStatistics.percentile(List.of(17.0), 99.0))
+ .isEqualTo(17.0);
+ }
+
+ @Test
+ @DisplayName("빈 표본이나 유효 범위 밖 Percentile은 거부한다")
+ void percentile_rejectsInvalidInput() {
+ assertThatThrownBy(() -> WorkerIndexingThroughputStatistics.percentile(List.of(), 50.0))
+ .isInstanceOf(IllegalArgumentException.class);
+ assertThatThrownBy(() -> WorkerIndexingThroughputStatistics.percentile(List.of(1.0), 101.0))
+ .isInstanceOf(IllegalArgumentException.class);
+ }
+}