diff --git a/build.gradle b/build.gradle
index 71c2236..f649ef9 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', 'worker-indexing-throughput'
+ excludeTags 'benchmark', 'minio-integration', 'claim-concurrency', 'local-e2e', 'vector-search-performance', 'worker-indexing-throughput', 'worker-horizontal-scaling'
}
}
@@ -174,6 +174,28 @@ tasks.register('workerIndexingThroughputTest', Test) {
outputs.upToDateWhen { false }
}
+tasks.register('workerHorizontalScalingTest', Test) {
+ group = 'verification'
+ description = '실제 PostgreSQL, MinIO와 BGE-M3에서 Worker 수·실행 Slot별 전체 인덱싱 확장성을 측정합니다.'
+ testClassesDirs = sourceSets.test.output.classesDirs
+ classpath = sourceSets.test.runtimeClasspath
+ useJUnitPlatform {
+ includeTags 'worker-horizontal-scaling'
+ }
+ maxParallelForks = 1
+ systemProperties System.properties.findAll { key, value ->
+ key.toString().startsWith('worker.horizontal.scaling.')
+ }
+ if (System.getProperty('worker.horizontal.scaling.output') == null) {
+ systemProperty(
+ 'worker.horizontal.scaling.output',
+ layout.buildDirectory.file('reports/worker-horizontal-scaling/worker-horizontal-scaling.json').get().asFile.absolutePath
+ )
+ }
+ // 여러 Spring Context가 같은 외부 Service와 Test Schema를 공유하므로 병렬 Fork와 Build Cache를 금지한다.
+ outputs.upToDateWhen { false }
+}
+
def configureOpenSqlDatabase = { Test task ->
task.maxParallelForks = 1
task.outputs.upToDateWhen { false }
diff --git a/docs/design/gimin-#138-worker-horizontal-scaling-benchmark.md b/docs/design/gimin-#138-worker-horizontal-scaling-benchmark.md
new file mode 100644
index 0000000..aa18086
--- /dev/null
+++ b/docs/design/gimin-#138-worker-horizontal-scaling-benchmark.md
@@ -0,0 +1,217 @@
+# 자동 Worker 수·실행 슬롯별 전체 인덱싱 수평 확장 Benchmark 설계
+
+## 1. 배경
+
+단일 자동 Worker와 실행 슬롯 2개에서 실제 PostgreSQL 17, MinIO와 `BAAI/bge-m3`를 통과하는 전체
+문서 인덱싱 처리량 기준선을 확보했다. 16개 문서는 분당 19.102개, 32개 문서는 분당 18.402개를
+처리했고 전체 지연 증가는 실제 처리보다 Queue 대기가 지배했다.
+
+이 기준선만으로 Worker 인스턴스를 늘렸을 때 처리량이 증가하는지, 실행 슬롯만 늘리는 것과 Worker 수를
+늘리는 것이 어떻게 다른지는 판단할 수 없다. 이번 작업은 같은 Job Queue를 여러 독립 Worker Context가
+경쟁해 처리하도록 구성하고 Worker 수와 Worker별 실행 슬롯 조합에 따른 처리량·지연·분산도를 측정한다.
+
+## 2. 목표
+
+1. Worker 1·2·4개가 실제 `FOR UPDATE SKIP LOCKED` Claim으로 Job을 나눠 처리하는지 검증한다.
+2. Worker별 실행 슬롯 1·2개가 문서·Chunk·Embedding 처리량과 Queue 지연에 주는 영향을 비교한다.
+3. `1 Worker × 1 Slot` 기준 Speedup과 전체 Slot 기준 Scaling Efficiency를 계산한다.
+4. 처리량이 증가하더라도 Attempt·Event·Chunk·Embedding·Vector 불변식이 유지되는지 확인한다.
+5. 단일 Host CPU BGE-M3가 수평 확장의 공통 병목이 되는 지점을 측정 결과로 설명한다.
+
+## 3. 범위
+
+### 3.1 포함
+
+- 실제 인증과 Multipart HTTP 문서 업로드
+- 실제 MinIO Object 저장·읽기
+- 실제 PostgreSQL Worker 등록·Heartbeat·Claim·Attempt·Lease
+- 독립 Worker별 Executor, Slot Pool, Scheduler와 Hikari Connection Pool
+- 실제 TXT Parsing·Chunk 저장·BGE-M3 Batch Embedding·`vector(1024)` 저장
+- Version `INDEXED`와 Document `current_version_id` 전환
+- Worker별 성공 Attempt 분포
+- 문서·Chunk·Embedding 처리량
+- Queue·처리·전체 p50·p95·p99·최댓값
+- Baseline 대비 Speedup과 Scaling Efficiency
+- 전용 Gradle Task와 Git 제외 JSON 원시 결과
+- 공개 가능한 실측 결과와 한계 문서
+
+### 3.2 제외
+
+- 여러 물리 Host·VM·Container 사이 Network 비용
+- Kubernetes, Service Discovery와 Load Balancer
+- 공식 Rocky Linux OpenSQL 원격 환경 재측정
+- Queue 상한과 요청 거부 같은 Backpressure 정책 구현
+- PDF·DOCX 형식별 Parser 성능 비교
+- GPU BGE-M3 확장성
+- 운영 SLO 확정
+
+## 4. 실행 구조
+
+Benchmark는 역할이 다른 두 종류의 Spring Context를 사용한다.
+
+```text
+Coordinator Context
+ ├─ RANDOM_PORT HTTP Server
+ ├─ 실제 인증·업로드 API
+ ├─ Benchmark 상태 조회·결과 검증
+ └─ indexing.worker.enabled=false
+
+Worker Context 1..N
+ ├─ WebApplicationType.SERVLET + server.port=0
+ ├─ 고유 Worker Name·Instance ID
+ ├─ 고유 Job Executor·Slot Pool
+ ├─ 고유 Hikari Pool
+ └─ 같은 PostgreSQL Schema·MinIO Bucket·BGE-M3 사용
+```
+
+Coordinator는 Job을 직접 Claim하지 않는다. Worker Context만 Worker Node로 등록되고 Production
+`WorkerJobPollingScheduler`를 실행한다. 따라서 Worker 수가 늘어날 때 각 Context가 실제 DB Lock과
+Claim Token 경계를 통과한다.
+
+Worker Context는 프로젝트의 실제 배포 애플리케이션과 같은 Servlet 자동 설정을 사용한다. 현재
+`SwaggerConfig`는 Web Application에서 제공되는 `SwaggerUiConfigProperties`를 주입받으므로, Worker만
+검증하려고 `WebApplicationType.NONE`을 사용하면 해당 Bean이 자기 자신을 주입하는 순환 생성이 발생한다.
+각 Context는 `server.port=0`으로 HTTP Port 충돌만 차단하고, 기동 시간은 처리량 측정에서 제외한다.
+
+Worker Context는 한 Profile 설정 동안 유지한다. 시작·Flyway 검증·종료 시간은 처리량 측정에서 제외한다.
+Profile이 끝난 뒤 모든 실행 슬롯이 반환된 것을 확인하고 Context를 역순으로 닫아 Worker를 `STOPPED`로
+전환한다.
+
+## 5. 격리 계약
+
+- Gradle 실행마다 UUID 기반 PostgreSQL Schema와 MinIO Bucket을 생성한다.
+- Schema 접미사는 PostgreSQL 식별자 63자 제한 안에서 24자로 제한한다.
+- 모든 Worker Context에 동일한 전용 Schema와 Bucket을 명시한다.
+- Worker Context마다 고유 `spring.application.name`, Worker Name과 Hikari Pool Name을 사용한다.
+- Profile 초기화는 Worker 실행 슬롯이 모두 반환된 뒤 Job·Attempt·Event·Document·Vector Data만 지운다.
+- Worker Context를 모두 닫은 뒤 이전 Profile의 `worker_nodes`를 정리한다.
+- 종료 시 Benchmark가 만든 Bucket과 Schema만 삭제한다.
+- DB·MinIO·JWT Credential과 접속 문자열은 JSON·Log·문서에 기록하지 않는다.
+
+## 6. Profile 계약
+
+기본 Profile은 Worker 수와 Worker별 Slot 수의 영향을 분리하면서 과도한 실행 조합을 피한다.
+
+| Profile | Worker 수 | Worker별 Slot | 전체 Slot | 비교 목적 |
+|---|---:|---:|---:|---|
+| `w1-s1` | 1 | 1 | 1 | Baseline |
+| `w1-s2` | 1 | 2 | 2 | 단일 Process 내부 Slot 확장 |
+| `w2-s1` | 2 | 1 | 2 | 같은 전체 Slot에서 Worker 수 효과 |
+| `w2-s2` | 2 | 2 | 4 | Worker와 Slot 동시 확장 |
+| `w4-s2` | 4 | 2 | 8 | Local 확장 상한과 BGE 병목 관찰 |
+
+기본 본 측정은 Profile마다 같은 16개 TXT 문서를 2회 처리하고, 각 Profile 시작 뒤 2개 문서로 예열한다.
+문서 본문은 6,400자로 고정해 문서당 Chunk·Embedding 수가 같게 한다.
+
+다음 System Property로 실행 범위를 조정할 수 있다.
+
+```text
+worker.horizontal.scaling.profiles
+worker.horizontal.scaling.document-characters
+worker.horizontal.scaling.document-count
+worker.horizontal.scaling.repetitions
+worker.horizontal.scaling.warm-up-documents
+worker.horizontal.scaling.uploader-threads
+worker.horizontal.scaling.profile-timeout-seconds
+worker.horizontal.scaling.status-polling-ms
+worker.horizontal.scaling.output
+```
+
+Profile 문자열은 `workerCount x slotsPerWorker` 형식의 쉼표 목록으로 받는다. Worker·Slot·문서·반복 값은
+모두 1 이상이어야 하고 중복 Profile은 거부한다. Speedup 기준선을 계산할 수 있도록 사용자 지정 Profile에도
+`1x1`을 반드시 포함해야 한다.
+
+## 7. 측정 경계와 지표
+
+### 7.1 측정 순서
+
+1. Profile Worker Context를 순서대로 시작하고 모든 Worker ID 등록을 확인한다.
+2. Warm-up 문서를 업로드하고 모두 `INDEXED`가 될 때까지 기다린다.
+3. Job Data를 초기화하고 Worker가 Idle인지 확인한다.
+4. 시작 시각을 기록하고 고정 Uploader Thread로 본 측정 문서를 병렬 접수한다.
+5. 마지막 Upload 응답 시각을 기록한다.
+6. 모든 대상 Job이 `INDEXED`이고 모든 Worker Slot이 반환될 때까지 기다린다.
+7. 결과 불변식과 Worker별 처리 분포를 검증한 뒤 지표를 계산한다.
+8. 반복 완료 뒤 중앙값을 기록하고 Worker Context를 닫는다.
+
+### 7.2 처리량과 지연
+
+| 지표 | 계산 |
+|---|---|
+| documents/s | 완료 문서 수 / 전체 측정 시간 |
+| documents/min | documents/s × 60 |
+| chunks/s | Chunk 수 / 전체 측정 시간 |
+| embeddings/s | Embedding 수 / 전체 측정 시간 |
+| Queue 대기 | Job `created_at` → 첫 `LOCKED` Event |
+| 실제 처리 | 첫 `LOCKED` → `INDEXED` Event |
+| 전체 지연 | Job `created_at` → `INDEXED` Event |
+| Speedup | Profile documents/s 중앙값 / `w1-s1` 중앙값 |
+| Scaling Efficiency | Speedup / Profile 전체 Slot 수 |
+
+지연은 선형 보간 p50·p95·p99와 max를 밀리초로 기록한다. Worker 분포는 성공 Attempt의
+`worker_node_id`별 Job 수를 저장한다.
+
+`Scaling Efficiency`는 전체 Slot 수 증가에 대한 단순 효율이다. Worker Context마다 별도 Connection Pool과
+Scheduler가 생기는 비용을 분리하지 않으므로 CPU·DB·BGE 자원 전체의 효율로 해석하지 않는다.
+
+## 8. 정합성 계약
+
+각 본 측정 반복은 다음 조건을 모두 만족해야 한다.
+
+- 대상 Job 전부 `INDEXED`
+- PENDING·PROCESSING 잔여 Job 없음
+- Job별 Retry 0회, `SUCCESS` Attempt 정확히 1개
+- `LOCKED → PARSE_STARTED → CHUNKED → EMBEDDING_STARTED → INDEXED` Event 순서
+- 실패·Lease 만료·Retry Event 없음
+- Document와 Version 전부 `INDEXED`
+- `documents.current_version_id`가 측정 Version을 가리킴
+- 문서별 Chunk 수와 Embedding 수 일치
+- 같은 Chunk의 Embedding 중복 없음
+- Vector 차원 1024, NaN·Infinity 없음
+- Profile의 모든 Worker가 등록 상태를 유지함
+- Worker 수가 2개 이상이고 문서 수가 Worker 수 이상이면 최소 2개 Worker가 성공 Attempt를 처리함
+- 측정 종료 시 모든 Worker 실행 Slot이 반환됨
+
+Worker 분배는 균등성을 강제하지 않는다. 단일 CPU BGE 응답 시간과 Scheduler Timing에 따라 분배 편차가
+생길 수 있으므로 참여 Worker 수와 실제 처리 건수만 관찰값으로 기록한다.
+
+## 9. 결과 판정
+
+- 처리량과 Speedup은 자동 합격선을 두지 않고 반복 중앙값으로 관찰한다.
+- 처리량이 감소해도 정합성을 만족하면 Benchmark는 통과하며 병목 근거로 기록한다.
+- 다중 Worker Profile에서 실제로 하나의 Worker만 모든 Job을 처리하면 수평 확장 검증 실패로 처리한다.
+- Slot 수가 늘어도 처리량이 증가하지 않으면 공유 BGE CPU, DB Connection 또는 Host Scheduling 병목 후보로
+ 기록한다.
+- Local 동일 JVM 결과를 다중 Host 운영 성능으로 표현하지 않는다.
+
+## 10. 실행 경계
+
+일반 `./gradlew test`는 이 Benchmark Tag를 제외한다. 실제 Infrastructure를 사용하는 전용 Task만 실행한다.
+
+```bash
+docker compose up -d postgres minio embedding-server
+DB_SSLMODE=disable ./gradlew workerHorizontalScalingTest
+```
+
+원시 JSON은 Git에 포함되지 않는 다음 경로에 생성한다.
+
+```text
+build/reports/worker-horizontal-scaling/worker-horizontal-scaling.json
+```
+
+## 11. 커밋 분할
+
+1. `docs: #138 Worker 수평 확장 Benchmark 설계 추가`
+2. `test: #138 다중 Worker 수평 확장 Benchmark 추가`
+3. `build: #138 Worker 수평 확장 전용 실행 경계 추가`
+4. `perf: #138 Worker 수평 확장 실측 결과 기록`
+
+리뷰 수정은 지적된 불변식과 실행 경계만 별도 커밋으로 반영한다.
+
+## 12. 완료 조건
+
+- 실제 Worker 1·2·4개와 Slot 조합이 같은 PENDING Queue를 나눠 처리한다.
+- 처리량·지연·Worker 분포·Speedup·Scaling Efficiency가 구조화된 JSON에 기록된다.
+- 모든 Profile에서 Job·Attempt·Event·Chunk·Embedding·Vector 불변식을 검증한다.
+- 일반 회귀와 실제 Infrastructure Smoke·정식 Benchmark가 통과한다.
+- 측정 환경, 병목 해석과 다중 Host로 일반화할 수 없는 한계를 결과 문서에 기록한다.
diff --git a/docs/test-results/gimin-#138-worker-horizontal-scaling-benchmark.md b/docs/test-results/gimin-#138-worker-horizontal-scaling-benchmark.md
new file mode 100644
index 0000000..77c9b05
--- /dev/null
+++ b/docs/test-results/gimin-#138-worker-horizontal-scaling-benchmark.md
@@ -0,0 +1,155 @@
+# Worker 수·실행 Slot별 전체 인덱싱 수평 확장 Benchmark 결과
+
+## 1. 목적
+
+[Worker 수평 확장 Benchmark 설계](../design/gimin-%23138-worker-horizontal-scaling-benchmark.md)에 따라
+실제 PostgreSQL, MinIO와 BGE-M3를 사용하는 자동 인덱싱 Pipeline에서 Worker 수와 Worker별 실행
+Slot 수의 영향을 비교했다.
+
+이번 결과는 다음 두 질문에 답한다.
+
+1. 여러 Worker Context가 같은 PENDING Job Queue를 실제로 나눠 처리하는가?
+2. Worker와 실행 Slot을 늘렸을 때 전체 처리량과 Queue·처리 지연이 어떻게 변하는가?
+
+## 2. 실행 환경
+
+| 항목 | 값 |
+|---|---|
+| 실행 일자 | 2026-08-10 |
+| OS | macOS aarch64 |
+| 사용 가능 CPU | 10 |
+| PostgreSQL | 17.8 |
+| pgvector | 0.8.1 |
+| Spring Boot | 3.5.16 |
+| Embedding Model | `BAAI/bge-m3` |
+| Embedding 차원 | 1024 |
+| Embedding Batch Size | 32 |
+| 문서 형식 | TXT |
+| 문서 크기 | 6,400자 |
+| Profile별 문서 수 | 16 |
+| 문서별 Chunk·Embedding 수 | 8·8 |
+| 반복 수 | 2 |
+| 예열 문서 수 | 2 |
+| HTTP Uploader Thread | 8 |
+
+모든 Worker는 같은 Mac, PostgreSQL Schema, MinIO Bucket과 CPU BGE-M3 Container를 공유했다.
+Worker Context 기동·Flyway 검증·종료 시간은 처리량에서 제외했다.
+
+## 3. 실행 명령
+
+```bash
+docker compose up -d postgres minio embedding-server
+DB_SSLMODE=disable ./gradlew workerHorizontalScalingTest
+```
+
+Gradle 결과:
+
+```text
+BUILD SUCCESSFUL in 8m 12s
+5 actionable tasks: 1 executed, 4 up-to-date
+```
+
+원시 JSON은 Git에 포함되지 않는 다음 경로에 생성했다.
+
+```text
+build/reports/worker-horizontal-scaling/worker-horizontal-scaling.json
+```
+
+## 4. 처리량·확장성 결과
+
+반복 2회의 중앙값이다.
+
+| Profile | Worker | Worker별 Slot | 전체 Slot | 문서/분 | Chunk·Embedding/초 | Speedup | Slot 효율 |
+|---|---:|---:|---:|---:|---:|---:|---:|
+| `w1-s1` | 1 | 1 | 1 | 19.475 | 2.597 | 1.000x | 1.000 |
+| `w1-s2` | 1 | 2 | 2 | 21.862 | 2.915 | 1.123x | 0.561 |
+| `w2-s1` | 2 | 1 | 2 | 21.576 | 2.877 | 1.108x | 0.554 |
+| `w2-s2` | 2 | 2 | 4 | 21.577 | 2.877 | 1.108x | 0.277 |
+| `w4-s2` | 4 | 2 | 8 | 21.677 | 2.890 | 1.113x | 0.139 |
+
+전체 Slot 1개에서 2개로 늘렸을 때 처리량은 약 11~12% 증가했다. 하지만 전체 Slot을 4개와 8개로
+늘려도 처리량은 약 21.6문서/분에서 더 증가하지 않았다. Slot 효율은 2 Slot 0.56, 4 Slot 0.28,
+8 Slot 0.14로 감소했다.
+
+## 5. 지연 결과
+
+반복 2회에서 각 p95 값을 구한 뒤 그 중앙값을 사용했다.
+
+| Profile | 전체 시간 | Queue p95 | 처리 p95 | E2E p95 |
+|---|---:|---:|---:|---:|
+| `w1-s1` | 49.298초 | 43.963초 | 3.391초 | 46.939초 |
+| `w1-s2` | 43.914초 | 38.302초 | 5.550초 | 43.671초 |
+| `w2-s1` | 44.494초 | 38.912초 | 5.647초 | 44.225초 |
+| `w2-s2` | 44.493초 | 33.627초 | 11.693초 | 44.294초 |
+| `w4-s2` | 44.300초 | 22.810초 | 22.741초 | 44.206초 |
+
+Slot이 늘수록 Job이 더 빨리 Claim되어 Queue p95는 감소했다. 반면 CPU BGE-M3 요청이 동시에 늘면서
+개별 Job 처리 p95가 증가했다. 두 효과가 상쇄되어 E2E p95와 전체 처리 시간은 2 Slot 이후 거의
+줄지 않았다.
+
+## 6. Worker 분배 결과
+
+| Profile | 반복 1 | 반복 2 | 참여 Worker 중앙값 |
+|---|---|---|---:|
+| `w1-s1` | 16 | 16 | 1 |
+| `w1-s2` | 16 | 16 | 1 |
+| `w2-s1` | 8 / 8 | 8 / 8 | 2 |
+| `w2-s2` | 8 / 8 | 8 / 8 | 2 |
+| `w4-s2` | 4 / 4 / 4 / 4 | 4 / 4 / 4 / 4 | 4 |
+
+모든 다중 Worker 본 측정에서 등록된 Worker가 실제 성공 Attempt를 나눠 처리했다. 이 결과는 같은
+Queue에 대한 `FOR UPDATE SKIP LOCKED`, Worker 소유권과 Claim Token 경계가 독립 Worker Context에서도
+유지됨을 보여준다.
+
+## 7. 정합성 검증
+
+모든 10개 본 측정 반복에서 다음 조건을 통과했다.
+
+- 16개 Job 전부 `INDEXED`, Retry 0회
+- Job별 `SUCCESS` Attempt 정확히 1개
+- `LOCKED → PARSE_STARTED → CHUNKED → EMBEDDING_STARTED → INDEXED` Event 순서
+- 실패·Lease 만료·Retry Event 없음
+- Document·Version `INDEXED`와 `current_version_id` 전환
+- 반복별 Chunk 128개, Embedding 128개
+- Chunk별 Embedding 중복 없음
+- 모든 Vector 1024차원, NaN·Infinity 없음
+- PENDING·PROCESSING 잔여 Job 없음
+- 다중 Worker 본 측정에서 최소 2개 Worker 참여
+- 완료 대기 경계에서 모든 실행 Slot 반환
+
+## 8. 실행 중 발견하고 보정한 경계
+
+### 8.1 Worker Context 유형
+
+`WebApplicationType.NONE`은 프로젝트의 `SwaggerConfig`가 웹 자동설정 Bean을 주입받는 현재 구조와
+맞지 않아 순환 Bean 생성이 발생했다. 제품 코드를 수정하지 않고 실제 배포 형태와 같은 Servlet
+Context를 사용하되 `server.port=0`으로 Port 충돌을 차단했다. Context 기동 시간은 측정에서 제외했다.
+
+### 8.2 JDBC Schema 전달
+
+`@DynamicPropertySource`는 별도 `SpringApplicationBuilder`에 자동 상속되지 않는다. Coordinator가 실제
+사용 중인 Schema 포함 JDBC URL을 Worker Context에 명시적으로 전달해 API 접수와 Worker Claim이 같은
+Queue를 보도록 했다. Credential은 결과와 Log에 기록하지 않았다.
+
+### 8.3 예열과 본 측정 분리
+
+2개 예열 문서는 Scheduler Timing에 따라 한 Worker가 모두 처리할 수 있다. 예열은 Cache·Model 준비와
+정합성만 검증하고, 다중 Worker 참여 조건은 16문서 본 측정에만 적용했다.
+
+## 9. 해석과 한계
+
+- 같은 전체 Slot 2개에서 `w1-s2`와 `w2-s1` 처리량이 거의 같으므로 Worker Context 자체의 추가
+ Overhead는 이번 규모에서 크지 않았다.
+- 2 Slot 이후 처리량이 포화되고 처리 p95만 늘어난 주된 후보는 모든 Worker가 공유한 단일 CPU
+ BGE-M3 Container다.
+- Worker 분배와 DB 정합성은 검증했지만 이 결과를 여러 물리 Host의 Network·Container 환경 성능으로
+ 일반화할 수 없다.
+- PostgreSQL과 MinIO도 같은 Host를 공유했으므로 BGE, DB, Storage 병목 기여도를 개별 분리하지 않았다.
+- 다중 Host 또는 GPU BGE 환경에서는 같은 Profile을 다시 실행해 확장 상한을 재측정해야 한다.
+
+## 10. 결론
+
+자동 Worker는 2개와 4개 Context에서 같은 Queue를 실제로 나눠 처리했고, 이번 실행에서는 2개 Worker가
+각 8건, 4개 Worker가 각 4건을 처리하는 분포가 관찰됐다. 모든 인덱싱·Vector 불변식도 유지했다. 현재
+단일 Host CPU 환경에서는 전체 Slot 2개가 실용적인 포화 지점이며, Worker·Slot을 그보다 늘리면 Queue
+대기는 줄지만 BGE 처리 대기가 증가해 최종 처리량은 약 21.6문서/분에 머물렀다.
diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingBenchmark.java b/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingBenchmark.java
new file mode 100644
index 0000000..d5b0dff
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingBenchmark.java
@@ -0,0 +1,1015 @@
+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.sql.Connection;
+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.Set;
+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 java.util.function.Supplier;
+
+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.WebApplicationType;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.boot.test.util.TestPropertyValues;
+import org.springframework.boot.test.web.client.TestRestTemplate;
+import org.springframework.context.ConfigurableApplicationContext;
+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.DocgridApplication;
+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.WorkerLifecycleManager;
+import com.opensource.docgrid.e2e.LocalE2eApiClient.UploadedDocument;
+import com.opensource.docgrid.e2e.LocalE2eDocumentFactory.DocumentPayload;
+import com.opensource.docgrid.e2e.WorkerHorizontalScalingStatistics.WorkerProfile;
+
+import io.minio.MinioClient;
+import lombok.extern.slf4j.Slf4j;
+
+/**
+ * 실제 PostgreSQL 17·pgvector·MinIO·BGE-M3에서 Worker 수와 Worker별 실행 Slot 확장성을 측정한다.
+ *
+ *
HTTP 접수만 담당하는 Coordinator와 Production Worker Lifecycle을 실행하는 독립 Spring Context를
+ * 분리한다. 모든 Worker는 같은 Job Queue를 DB Claim으로 경쟁하며 처리량, 지연, 실제 Worker 분포와
+ * 데이터 불변식을 함께 검증한다. 실제 외부 인프라를 점유하므로 전용 Gradle Task에서만 실행한다.
+ */
+@Slf4j
+@Tag("integration")
+@Tag("worker-horizontal-scaling")
+@ActiveProfiles({"test", "minio-integration"})
+@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
+@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@DisplayName("Worker 수·실행 Slot 수평 확장 Benchmark")
+class WorkerHorizontalScalingBenchmark {
+
+ private static final String EXECUTION_ID = UUID.randomUUID().toString().replace("-", "");
+ // PostgreSQL 식별자 63자 제한 안에서 별도 Gradle 실행이 Schema를 공유하지 않도록 격리한다.
+ private static final String TEST_SCHEMA = "docgrid_worker_horizontal_"
+ + EXECUTION_ID.substring(0, 24);
+ private static final String TEST_BUCKET = "docgrid-worker-horizontal-" + EXECUTION_ID;
+ private static final String TEST_JWT_SECRET =
+ "docgrid-worker-horizontal-scaling-test-secret-key-2026";
+ 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 List DEFAULT_PROFILES = List.of(
+ new WorkerProfile(1, 1),
+ new WorkerProfile(1, 2),
+ new WorkerProfile(2, 1),
+ new WorkerProfile(2, 2),
+ new WorkerProfile(4, 2)
+ );
+ private static final List PROFILES =
+ WorkerHorizontalScalingStatistics.parseProfiles(
+ System.getProperty("worker.horizontal.scaling.profiles"),
+ DEFAULT_PROFILES
+ );
+ private static final int DOCUMENT_CHARACTER_COUNT = positiveIntegerProperty(
+ "worker.horizontal.scaling.document-characters",
+ 6_400
+ );
+ private static final int DOCUMENT_COUNT = positiveIntegerProperty(
+ "worker.horizontal.scaling.document-count",
+ 16
+ );
+ private static final int REPETITIONS = positiveIntegerProperty(
+ "worker.horizontal.scaling.repetitions",
+ 2
+ );
+ private static final int WARM_UP_DOCUMENT_COUNT = positiveIntegerProperty(
+ "worker.horizontal.scaling.warm-up-documents",
+ 2
+ );
+ private static final int UPLOADER_THREADS = positiveIntegerProperty(
+ "worker.horizontal.scaling.uploader-threads",
+ 8
+ );
+ private static final long PROFILE_TIMEOUT_SECONDS = positiveLongProperty(
+ "worker.horizontal.scaling.profile-timeout-seconds",
+ 600L
+ );
+ private static final long POLLING_SLEEP_MILLIS = positiveLongProperty(
+ "worker.horizontal.scaling.status-polling-ms",
+ 100L
+ );
+ private static final Path OUTPUT_PATH = Path.of(System.getProperty(
+ "worker.horizontal.scaling.output",
+ "build/reports/worker-horizontal-scaling/worker-horizontal-scaling.json"
+ ));
+
+ @Autowired private TestRestTemplate restTemplate;
+ @Autowired private JdbcTemplate jdbcTemplate;
+ @Autowired private MinioClient minioClient;
+ @Autowired private ObjectMapper objectMapper;
+ @Autowired private EmbeddingBatchProperties embeddingBatchProperties;
+ @Autowired @Qualifier("embeddingRestClient") private RestClient embeddingRestClient;
+
+ private LocalE2eApiClient apiClient;
+ private LocalE2eMinioBucket minioBucket;
+ private String accessToken;
+ private String coordinatorJdbcUrl;
+ private WorkerCluster currentCluster;
+
+ @DynamicPropertySource
+ static void configureCoordinatorEnvironment(DynamicPropertyRegistry registry) {
+ registry.add("TEST_DB_SCHEMA", () -> TEST_SCHEMA);
+ registry.add("jwt.secret", () -> TEST_JWT_SECRET);
+ registry.add("minio.bucket", () -> TEST_BUCKET);
+ // Coordinator는 HTTP 접수와 검증만 담당하고 Job Claim 경쟁에는 참여하지 않는다.
+ registry.add("indexing.worker.enabled", () -> "false");
+ registry.add("embedding.server.read-timeout", () -> "2m");
+ registry.add("spring.datasource.hikari.maximum-pool-size", () -> "8");
+ }
+
+ @BeforeAll
+ void setUpInfrastructure() throws Exception {
+ apiClient = new LocalE2eApiClient(restTemplate);
+ minioBucket = new LocalE2eMinioBucket(minioClient, TEST_BUCKET);
+ minioBucket.create();
+ accessToken = apiClient.loginAdmin();
+ try (Connection connection = jdbcTemplate.getDataSource().getConnection()) {
+ // DynamicPropertySource로 완성된 Schema 포함 URL을 독립 Worker Context에도 그대로 전달한다.
+ coordinatorJdbcUrl = connection.getMetaData().getURL();
+ }
+ resetAllBenchmarkState();
+ }
+
+ @AfterAll
+ void cleanUpInfrastructure() throws Exception {
+ // 1. 실패 중단 경로에서도 독립 Worker Context를 먼저 닫아 DB 재접근을 차단한다.
+ if (currentCluster != null) {
+ currentCluster.close();
+ currentCluster = null;
+ }
+
+ // 2. Benchmark가 만든 Bucket과 Schema만 제거해 기존 로컬 개발 Data를 보존한다.
+ try {
+ if (minioBucket != null) {
+ minioBucket.close();
+ }
+ } finally {
+ jdbcTemplate.execute("DROP SCHEMA IF EXISTS " + TEST_SCHEMA + " CASCADE");
+ }
+ }
+
+ @Test
+ @Timeout(3_600)
+ @DisplayName("Worker 수와 Worker별 실행 Slot 조합의 전체 인덱싱 확장성을 반복 측정한다")
+ void measureWorkerAndExecutionSlotHorizontalScaling() throws Exception {
+ // 1. 환경과 Workload 계약을 기록한 뒤 각 Profile을 독립 Worker Cluster로 실행한다.
+ EnvironmentFingerprint environment = validateEnvironment();
+ List runs = new ArrayList<>();
+ writeReport(new BenchmarkReport(environment, runs, List.of()));
+ logJson("WORKER_HORIZONTAL_SCALING_ENV", environment);
+
+ for (WorkerProfile profile : PROFILES) {
+ resetAllBenchmarkState();
+ currentCluster = startWorkerCluster(profile);
+ Throwable primaryFailure = null;
+ try {
+ runWarmUp(profile, currentCluster);
+
+ // 2. 같은 문서 분포를 반복 처리해 Worker 수와 Slot 수 외의 입력 차이를 줄인다.
+ for (int repetition = 1; repetition <= REPETITIONS; repetition++) {
+ resetJobState(currentCluster);
+ ProfileRun run = runMeasuredProfile(profile, repetition, currentCluster);
+ runs.add(run);
+ logJson("WORKER_HORIZONTAL_SCALING_RESULT", run);
+ writeReport(new BenchmarkReport(environment, runs, medians(runs)));
+ }
+ } catch (Exception | Error failure) {
+ primaryFailure = failure;
+ throw failure;
+ } finally {
+ WorkerCluster completedCluster = currentCluster;
+ closeWorkerCluster(completedCluster, primaryFailure);
+ }
+ }
+
+ // 3. 1 Worker × 1 Slot 중앙값 대비 Speedup과 전체 Slot 기준 효율을 계산한다.
+ List medians = medians(runs);
+ for (ProfileMedian median : medians) {
+ logJson("WORKER_HORIZONTAL_SCALING_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");
+
+ return new EnvironmentFingerprint(
+ postgresVersion,
+ pgvectorVersion,
+ SpringBootVersion.getVersion(),
+ EXPECTED_MODEL,
+ EXPECTED_VECTOR_DIMENSION,
+ embeddingBatchProperties.getBatchSize(),
+ PROFILES,
+ DOCUMENT_COUNT,
+ REPETITIONS,
+ WARM_UP_DOCUMENT_COUNT,
+ DOCUMENT_CHARACTER_COUNT,
+ UPLOADER_THREADS,
+ System.getProperty("os.name"),
+ System.getProperty("os.arch"),
+ Runtime.getRuntime().availableProcessors()
+ );
+ }
+
+ private WorkerCluster startWorkerCluster(WorkerProfile profile) throws InterruptedException {
+ List contexts = new ArrayList<>(profile.workerCount());
+ List workers = new ArrayList<>(profile.workerCount());
+
+ try {
+ // Worker Context를 순차 시작해 같은 Schema의 Flyway 검증이 서로 경합하지 않게 한다.
+ for (int index = 1; index <= profile.workerCount(); index++) {
+ String workerName = profile.name() + "-worker-" + String.format(Locale.ROOT, "%02d", index);
+ String poolName = "worker-horizontal-" + profile.name() + "-" + index;
+ ConfigurableApplicationContext context = new SpringApplicationBuilder(DocgridApplication.class)
+ // 제품 Swagger 설정까지 포함한 실제 배포 형태를 유지하되 임의 포트로 충돌을 막는다.
+ .web(WebApplicationType.SERVLET)
+ .profiles("test", "minio-integration")
+ .properties(
+ "spring.main.banner-mode=off",
+ "spring.jmx.enabled=false"
+ )
+ .initializers(applicationContext -> TestPropertyValues.of(
+ "TEST_DB_SCHEMA=" + TEST_SCHEMA,
+ "spring.datasource.url=" + coordinatorJdbcUrl,
+ "jwt.secret=" + TEST_JWT_SECRET,
+ "minio.bucket=" + TEST_BUCKET,
+ "indexing.worker.enabled=true",
+ "indexing.worker.name=" + workerName,
+ "indexing.worker.polling-interval=50ms",
+ "indexing.worker.heartbeat-interval=1s",
+ "indexing.worker.dead-threshold=2m",
+ "indexing.worker.max-concurrency=" + profile.slotsPerWorker(),
+ "indexing.worker.lease-duration=2m",
+ "indexing.worker.lease-renewal-interval=10s",
+ "indexing.worker.lease-recovery-interval=10m",
+ "indexing.worker.shutdown-grace-period=30s",
+ "embedding.server.read-timeout=2m",
+ "server.port=0",
+ "spring.application.name=" + workerName,
+ "spring.datasource.hikari.pool-name=" + poolName,
+ "spring.datasource.hikari.maximum-pool-size="
+ + Math.max(4, profile.slotsPerWorker() + 2)
+ ).applyTo(applicationContext))
+ .run();
+ contexts.add(context);
+
+ WorkerLifecycleManager lifecycleManager = context.getBean(WorkerLifecycleManager.class);
+ WorkerExecutionSlotPool slotPool = context.getBean(WorkerExecutionSlotPool.class);
+ IndexingWorkerProperties properties = context.getBean(IndexingWorkerProperties.class);
+ awaitCondition(
+ () -> workerName + "가 Application Ready 뒤 등록되지 않았습니다.",
+ () -> lifecycleManager.getWorkerId().isPresent()
+ );
+ assertThat(properties.getMaxConcurrency()).isEqualTo(profile.slotsPerWorker());
+ assertThat(slotPool.getCapacity()).isEqualTo(profile.slotsPerWorker());
+ workers.add(new WorkerHandle(
+ workerName,
+ lifecycleManager.getWorkerId().orElseThrow(),
+ slotPool
+ ));
+ }
+
+ WorkerCluster cluster = new WorkerCluster(contexts, workers);
+ assertRegisteredWorkers(cluster);
+ return cluster;
+ } catch (RuntimeException | InterruptedException | AssertionError exception) {
+ closeContexts(contexts);
+ throw exception;
+ }
+ }
+
+ private void runWarmUp(WorkerProfile profile, WorkerCluster cluster) throws Exception {
+ resetJobState(cluster);
+ List uploads = uploadDocuments(
+ profile.name() + "-warm-up",
+ WARM_UP_DOCUMENT_COUNT
+ );
+ awaitIndexedAndIdle(uploads, profile.name() + " Warm-up", cluster);
+ // 예열은 Cache와 Model 준비가 목적이므로 Worker 분산 자체는 본 측정에서만 강제한다.
+ assertProfileInvariants(uploads, cluster, false);
+ log.info(
+ "Worker 수평 확장 Benchmark 예열을 완료했습니다. profile={}, documentCount={}",
+ profile.name(),
+ WARM_UP_DOCUMENT_COUNT
+ );
+ }
+
+ private ProfileRun runMeasuredProfile(
+ WorkerProfile profile,
+ int repetition,
+ WorkerCluster cluster
+ ) throws Exception {
+ String runName = profile.name() + "-run-" + repetition;
+ long profileStartedAt = System.nanoTime();
+ List uploads = uploadDocuments(runName, DOCUMENT_COUNT);
+ long uploadCompletedAt = System.nanoTime();
+
+ awaitIndexedAndIdle(uploads, runName, cluster);
+ long profileCompletedAt = System.nanoTime();
+ ProfileData profileData = assertProfileInvariants(uploads, cluster, true);
+
+ double elapsedSeconds = seconds(profileCompletedAt - profileStartedAt);
+ double uploadSeconds = seconds(uploadCompletedAt - profileStartedAt);
+ double queueDrainSeconds = seconds(profileCompletedAt - uploadCompletedAt);
+ return new ProfileRun(
+ profile.name(),
+ profile.workerCount(),
+ profile.slotsPerWorker(),
+ profile.totalSlots(),
+ repetition,
+ DOCUMENT_COUNT,
+ profileData.chunkCount(),
+ profileData.embeddingCount(),
+ elapsedSeconds,
+ uploadSeconds,
+ queueDrainSeconds,
+ DOCUMENT_COUNT / elapsedSeconds,
+ DOCUMENT_COUNT / elapsedSeconds * 60.0,
+ profileData.chunkCount() / elapsedSeconds,
+ profileData.embeddingCount() / elapsedSeconds,
+ profileData.workerDistribution(),
+ profileData.workerDistribution().size(),
+ 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 runName, 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,
+ scalingDocument(runName, 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 scalingDocument(String runName, int documentIndex) {
+ String marker = String.format(
+ Locale.ROOT,
+ "DocGrid horizontal scaling document %04d. ",
+ documentIndex
+ );
+ String sentence = "Distributed 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 = runName + "-" + String.format(Locale.ROOT, "%04d", documentIndex) + ".txt";
+ return LocalE2eDocumentFactory.text(fileName, "Horizontal Scaling " + fileName, body.toString());
+ }
+
+ private void awaitIndexedAndIdle(
+ List uploads,
+ String runName,
+ WorkerCluster cluster
+ ) throws InterruptedException {
+ awaitCondition(
+ () -> runName + " Profile이 완료되지 않았습니다. jobs=" + jobSnapshot(uploads)
+ + ", workers=" + workerSnapshot(cluster),
+ () -> indexedJobCount(uploads) == uploads.size() && cluster.allSlotsReturned()
+ );
+ }
+
+ private ProfileData assertProfileInvariants(
+ List uploads,
+ WorkerCluster cluster,
+ boolean requireMultipleWorkerParticipation
+ ) {
+ List timings = new ArrayList<>(uploads.size());
+ List chunkCounts = new ArrayList<>(uploads.size());
+ Map workerDistribution = new LinkedHashMap<>();
+ Set clusterWorkerIds = cluster.workerIds();
+ 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();
+ SuccessfulAttempt attempt = readSuccessfulAttempt(upload.embeddingJobId());
+ assertThat(clusterWorkerIds).contains(attempt.workerId());
+ workerDistribution.merge(attempt.workerName(), 1, Integer::sum);
+ 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 차원과 유한값 계약을 지키는지 확인한다.
+ 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);
+ assertThat(count(
+ "SELECT COUNT(*) FROM embeddings WHERE document_version_id = ? "
+ + "AND (LOWER(vector::text) LIKE '%nan%' "
+ + "OR LOWER(vector::text) LIKE '%infinity%')",
+ upload.documentVersionId()
+ )).isZero();
+ 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 전체 Queue와 실행 슬롯이 비었으며 다중 Worker가 실제 처리에 참여했는지 확인한다.
+ assertThat(chunkCounts).allMatch(count -> count.equals(chunkCounts.get(0)));
+ assertThat(count(
+ "SELECT COUNT(*) FROM embedding_jobs WHERE status IN ('PENDING', 'PROCESSING')"
+ )).isZero();
+ assertThat(workerDistribution.values().stream().mapToInt(Integer::intValue).sum())
+ .isEqualTo(uploads.size());
+ if (requireMultipleWorkerParticipation && cluster.workerCount() > 1) {
+ assertThat(workerDistribution.size()).isGreaterThanOrEqualTo(2);
+ }
+ assertRegisteredWorkers(cluster);
+ return new ProfileData(
+ chunkCounts.stream().mapToInt(Integer::intValue).sum(),
+ totalEmbeddings,
+ Map.copyOf(workerDistribution),
+ List.copyOf(timings)
+ );
+ }
+
+ private SuccessfulAttempt readSuccessfulAttempt(Long jobId) {
+ return jdbcTemplate.queryForObject(
+ "SELECT attempt.worker_node_id, worker.worker_name "
+ + "FROM embedding_job_attempts attempt "
+ + "JOIN worker_nodes worker ON worker.id = attempt.worker_node_id "
+ + "WHERE attempt.embedding_job_id = ? AND attempt.status = 'SUCCESS'",
+ (resultSet, rowNumber) -> new SuccessfulAttempt(
+ resultSet.getLong("worker_node_id"),
+ resultSet.getString("worker_name")
+ ),
+ jobId
+ );
+ }
+
+ 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 assertRegisteredWorkers(WorkerCluster cluster) {
+ for (WorkerHandle worker : cluster.workers()) {
+ assertThat(queryString("SELECT worker_name FROM worker_nodes WHERE id = ?", worker.id()))
+ .isEqualTo(worker.name());
+ assertThat(queryString("SELECT status FROM worker_nodes WHERE id = ?", worker.id()))
+ .isIn("ACTIVE", "IDLE");
+ assertThat(worker.slotPool().getCapacity()).isPositive();
+ }
+ }
+
+ private void awaitWorkersStopped(WorkerCluster cluster) throws InterruptedException {
+ String placeholders = String.join(",", cluster.workers().stream().map(worker -> "?").toList());
+ Object[] workerIds = cluster.workers().stream().map(WorkerHandle::id).toArray();
+ awaitCondition(
+ () -> "종료한 Worker가 STOPPED 상태로 전환되지 않았습니다: " + workerSnapshot(cluster),
+ () -> count(
+ "SELECT COUNT(*) FROM worker_nodes WHERE status = 'STOPPED' AND id IN ("
+ + placeholders + ")",
+ workerIds
+ ) == cluster.workerCount()
+ );
+ }
+
+ private void resetJobState(WorkerCluster cluster) throws Exception {
+ awaitCondition(
+ () -> "이전 Profile의 Worker 실행 Slot이 반환되지 않았습니다: " + workerSnapshot(cluster),
+ cluster::allSlotsReturned
+ );
+ 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 void resetAllBenchmarkState() throws Exception {
+ minioBucket.clear();
+ jdbcTemplate.execute("""
+ TRUNCATE TABLE
+ indexing_events,
+ embedding_job_attempts,
+ embeddings,
+ document_chunks,
+ embedding_jobs,
+ document_versions,
+ documents,
+ file_objects,
+ worker_nodes
+ 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 String workerSnapshot(WorkerCluster cluster) {
+ String placeholders = String.join(",", cluster.workers().stream().map(worker -> "?").toList());
+ Object[] workerIds = cluster.workers().stream().map(WorkerHandle::id).toArray();
+ return jdbcTemplate.queryForList(
+ "SELECT id || ':' || worker_name || ':' || status FROM worker_nodes WHERE id IN ("
+ + placeholders + ") ORDER BY id",
+ String.class,
+ workerIds
+ ).toString();
+ }
+
+ private List medians(List runs) {
+ Map throughputMedians = new LinkedHashMap<>();
+ for (WorkerProfile profile : PROFILES) {
+ List profileRuns = runs.stream()
+ .filter(run -> run.profileName().equals(profile.name()))
+ .toList();
+ if (!profileRuns.isEmpty()) {
+ throughputMedians.put(
+ profile.name(),
+ median(profileRuns, ProfileRun::documentsPerSecond)
+ );
+ }
+ }
+
+ Double baseline = throughputMedians.get(new WorkerProfile(1, 1).name());
+ if (baseline == null) {
+ return List.of();
+ }
+
+ List result = new ArrayList<>();
+ for (WorkerProfile profile : PROFILES) {
+ List profileRuns = runs.stream()
+ .filter(run -> run.profileName().equals(profile.name()))
+ .toList();
+ if (profileRuns.isEmpty()) {
+ continue;
+ }
+ double documentsPerSecond = throughputMedians.get(profile.name());
+ double speedup = WorkerHorizontalScalingStatistics.speedup(baseline, documentsPerSecond);
+ result.add(new ProfileMedian(
+ profile.name(),
+ profile.workerCount(),
+ profile.slotsPerWorker(),
+ profile.totalSlots(),
+ profileRuns.size(),
+ 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, ProfileRun::participatingWorkers),
+ median(profileRuns, run -> run.queueWait().p95Millis()),
+ median(profileRuns, run -> run.processing().p95Millis()),
+ median(profileRuns, run -> run.endToEnd().p95Millis()),
+ speedup,
+ WorkerHorizontalScalingStatistics.scalingEfficiency(speedup, profile.totalSlots())
+ ));
+ }
+ return List.copyOf(result);
+ }
+
+ 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(Supplier 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.get());
+ }
+
+ private void closeWorkerCluster(WorkerCluster cluster, Throwable primaryFailure)
+ throws InterruptedException {
+ try {
+ // 1. 독립 Context를 닫고 DB의 Worker 종착 상태까지 확인한다.
+ cluster.close();
+ awaitWorkersStopped(cluster);
+ } catch (RuntimeException | Error | InterruptedException cleanupFailure) {
+ // 2. 본 측정 실패가 있으면 정리 실패는 보조 근거로 남겨 최초 원인을 보존한다.
+ if (primaryFailure != null) {
+ primaryFailure.addSuppressed(cleanupFailure);
+ return;
+ }
+ throw cleanupFailure;
+ } finally {
+ currentCluster = null;
+ }
+ }
+
+ 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 void closeContexts(List contexts) {
+ for (int index = contexts.size() - 1; index >= 0; index--) {
+ contexts.get(index).close();
+ }
+ }
+
+ 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;
+ }
+
+ /** Secret을 제외하고 실행 환경과 Workload를 재현하는 지문이다. */
+ private record EnvironmentFingerprint(
+ String postgresVersion,
+ String pgvectorVersion,
+ String springBootVersion,
+ String embeddingModel,
+ int vectorDimension,
+ int embeddingBatchSize,
+ List profiles,
+ int documentCount,
+ int repetitions,
+ int warmUpDocumentCount,
+ int documentCharacterCount,
+ int uploaderThreads,
+ String osName,
+ String osArchitecture,
+ int availableProcessors
+ ) {
+ }
+
+ /** 한 Worker·Slot Profile 반복의 처리량, 지연과 실제 처리 분포 원시 결과다. */
+ private record ProfileRun(
+ String profileName,
+ int workerCount,
+ int slotsPerWorker,
+ int totalSlots,
+ int repetition,
+ int documentCount,
+ int chunkCount,
+ int embeddingCount,
+ double totalElapsedSeconds,
+ double uploadElapsedSeconds,
+ double queueDrainSeconds,
+ double documentsPerSecond,
+ double documentsPerMinute,
+ double chunksPerSecond,
+ double embeddingsPerSecond,
+ Map workerDistribution,
+ int participatingWorkers,
+ LatencySummary queueWait,
+ LatencySummary processing,
+ LatencySummary endToEnd
+ ) {
+ }
+
+ /** 같은 Profile 반복 중앙값과 1×1 Baseline 대비 수평 확장 지표다. */
+ private record ProfileMedian(
+ String profileName,
+ int workerCount,
+ int slotsPerWorker,
+ int totalSlots,
+ int completedRepetitions,
+ double documentsPerSecond,
+ double documentsPerMinute,
+ double chunksPerSecond,
+ double embeddingsPerSecond,
+ double totalElapsedSeconds,
+ double uploadElapsedSeconds,
+ double queueDrainSeconds,
+ double participatingWorkersMedian,
+ double queueWaitP95Millis,
+ double processingP95Millis,
+ double endToEndP95Millis,
+ double speedup,
+ double scalingEfficiency
+ ) {
+ }
+
+ /** 한 지연 분포의 중앙값, 꼬리 지연과 최댓값을 밀리초로 보존한다. */
+ 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 수, Worker 분포와 Job 단계별 시각이다. */
+ private record ProfileData(
+ int chunkCount,
+ int embeddingCount,
+ Map workerDistribution,
+ List timings
+ ) {
+ }
+
+ /** 한 Job의 Queue, 실제 처리와 전체 지연을 밀리초로 표현한다. */
+ private record JobTiming(double queueWaitMillis, double processingMillis, double endToEndMillis) {
+ }
+
+ /** 성공한 Attempt가 실제 어느 Worker에서 처리됐는지 보존한다. */
+ private record SuccessfulAttempt(Long workerId, String workerName) {
+ }
+
+ /** 한 Worker Context의 DB 식별자와 실행 Slot Pool을 묶는다. */
+ private record WorkerHandle(String name, Long id, WorkerExecutionSlotPool slotPool) {
+ }
+
+ /**
+ * 한 Profile 동안 유지되는 독립 Worker Context 집합의 시작·실행 Slot·종료 경계를 관리한다.
+ */
+ private static final class WorkerCluster implements AutoCloseable {
+
+ private final List contexts;
+ private final List workers;
+
+ private WorkerCluster(
+ List contexts,
+ List workers
+ ) {
+ this.contexts = List.copyOf(contexts);
+ this.workers = List.copyOf(workers);
+ }
+
+ private List workers() {
+ return workers;
+ }
+
+ private int workerCount() {
+ return workers.size();
+ }
+
+ private Set workerIds() {
+ return workers.stream().map(WorkerHandle::id).collect(java.util.stream.Collectors.toSet());
+ }
+
+ private boolean allSlotsReturned() {
+ return workers.stream().allMatch(worker ->
+ worker.slotPool().getActiveSlots() == 0
+ && worker.slotPool().getAvailableSlots() == worker.slotPool().getCapacity()
+ );
+ }
+
+ @Override
+ public void close() {
+ closeContexts(contexts);
+ }
+ }
+
+ /** 환경, 원시 반복과 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/WorkerHorizontalScalingStatistics.java b/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingStatistics.java
new file mode 100644
index 0000000..a4c9aa7
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingStatistics.java
@@ -0,0 +1,97 @@
+package com.opensource.docgrid.e2e;
+
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Locale;
+import java.util.Set;
+
+/**
+ * Worker 수평 확장 Benchmark의 Profile 입력과 Baseline 대비 비교 지표를 계산한다.
+ *
+ * 실제 인프라 실행과 분리된 순수 계산만 담당해 Profile 오입력, Speedup과 Slot 기준 효율 계산을
+ * 일반 단위 테스트에서 빠르게 검증할 수 있게 한다.
+ */
+final class WorkerHorizontalScalingStatistics {
+
+ private static final WorkerProfile BASELINE = new WorkerProfile(1, 1);
+
+ private WorkerHorizontalScalingStatistics() {
+ }
+
+ static List parseProfiles(String rawValue, List defaults) {
+ List profiles = rawValue == null || rawValue.isBlank()
+ ? List.copyOf(defaults)
+ : List.of(rawValue.split(",")).stream()
+ .map(String::trim)
+ .map(WorkerHorizontalScalingStatistics::parseProfile)
+ .toList();
+
+ if (profiles.isEmpty()) {
+ throw new IllegalArgumentException("Worker 수평 확장 Profile은 하나 이상이어야 합니다.");
+ }
+ Set uniqueProfiles = new LinkedHashSet<>(profiles);
+ if (uniqueProfiles.size() != profiles.size()) {
+ throw new IllegalArgumentException("Worker 수평 확장 Profile은 중복될 수 없습니다.");
+ }
+ if (!uniqueProfiles.contains(BASELINE)) {
+ throw new IllegalArgumentException("Speedup 기준 Profile 1x1이 필요합니다.");
+ }
+ return List.copyOf(uniqueProfiles);
+ }
+
+ static double speedup(double baselineDocumentsPerSecond, double documentsPerSecond) {
+ if (baselineDocumentsPerSecond <= 0.0) {
+ throw new IllegalArgumentException("Baseline 처리량은 0보다 커야 합니다.");
+ }
+ if (documentsPerSecond < 0.0) {
+ throw new IllegalArgumentException("Profile 처리량은 0보다 작을 수 없습니다.");
+ }
+ return documentsPerSecond / baselineDocumentsPerSecond;
+ }
+
+ static double scalingEfficiency(double speedup, int totalSlots) {
+ if (speedup < 0.0) {
+ throw new IllegalArgumentException("Speedup은 0보다 작을 수 없습니다.");
+ }
+ if (totalSlots < 1) {
+ throw new IllegalArgumentException("전체 실행 Slot 수는 1 이상이어야 합니다.");
+ }
+ return speedup / totalSlots;
+ }
+
+ private static WorkerProfile parseProfile(String value) {
+ String normalized = value.toLowerCase(Locale.ROOT);
+ String[] parts = normalized.split("x", -1);
+ if (parts.length != 2) {
+ throw new IllegalArgumentException(
+ "Worker 수평 확장 Profile은 workerCountxslotsPerWorker 형식이어야 합니다: " + value
+ );
+ }
+ try {
+ return new WorkerProfile(Integer.parseInt(parts[0]), Integer.parseInt(parts[1]));
+ } catch (NumberFormatException exception) {
+ throw new IllegalArgumentException(
+ "Worker 수평 확장 Profile에는 양의 정수만 사용할 수 있습니다: " + value,
+ exception
+ );
+ }
+ }
+
+ /** Worker 수와 Worker별 실행 Slot 수를 하나의 측정 Profile로 표현한다. */
+ record WorkerProfile(int workerCount, int slotsPerWorker) {
+
+ WorkerProfile {
+ if (workerCount < 1 || slotsPerWorker < 1) {
+ throw new IllegalArgumentException("Worker 수와 Worker별 실행 Slot 수는 1 이상이어야 합니다.");
+ }
+ }
+
+ int totalSlots() {
+ return Math.multiplyExact(workerCount, slotsPerWorker);
+ }
+
+ String name() {
+ return "w" + workerCount + "-s" + slotsPerWorker;
+ }
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingStatisticsTest.java b/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingStatisticsTest.java
new file mode 100644
index 0000000..18f1aeb
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/e2e/WorkerHorizontalScalingStatisticsTest.java
@@ -0,0 +1,60 @@
+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;
+
+import com.opensource.docgrid.e2e.WorkerHorizontalScalingStatistics.WorkerProfile;
+
+/** Worker 수평 확장 Profile 파싱과 비교 지표의 순수 계산 계약을 검증한다. */
+@DisplayName("Worker 수평 확장 통계")
+class WorkerHorizontalScalingStatisticsTest {
+
+ @Test
+ @DisplayName("Worker 수와 Slot 수 Profile을 입력 순서대로 파싱한다")
+ void parsesWorkerAndSlotProfilesInOrder() {
+ List profiles = WorkerHorizontalScalingStatistics.parseProfiles(
+ "1x1, 1x2,2x1,4x2",
+ List.of()
+ );
+
+ assertThat(profiles).containsExactly(
+ new WorkerProfile(1, 1),
+ new WorkerProfile(1, 2),
+ new WorkerProfile(2, 1),
+ new WorkerProfile(4, 2)
+ );
+ assertThat(profiles.get(3).name()).isEqualTo("w4-s2");
+ assertThat(profiles.get(3).totalSlots()).isEqualTo(8);
+ }
+
+ @Test
+ @DisplayName("중복 Profile과 1x1 Baseline 누락을 거부한다")
+ void rejectsDuplicateOrMissingBaselineProfiles() {
+ assertThatThrownBy(() -> WorkerHorizontalScalingStatistics.parseProfiles(
+ "1x1,1x1",
+ List.of()
+ )).isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("중복");
+
+ assertThatThrownBy(() -> WorkerHorizontalScalingStatistics.parseProfiles(
+ "2x1,2x2",
+ List.of()
+ )).isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("1x1");
+ }
+
+ @Test
+ @DisplayName("Baseline 처리량 대비 Speedup과 전체 Slot 기준 효율을 계산한다")
+ void calculatesSpeedupAndScalingEfficiency() {
+ double speedup = WorkerHorizontalScalingStatistics.speedup(2.0, 6.0);
+
+ assertThat(speedup).isEqualTo(3.0);
+ assertThat(WorkerHorizontalScalingStatistics.scalingEfficiency(speedup, 4))
+ .isEqualTo(0.75);
+ }
+}