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); + } +}