diff --git a/build.gradle b/build.gradle index a2830ab..71c2236 100644 --- a/build.gradle +++ b/build.gradle @@ -52,7 +52,7 @@ dependencies { tasks.named('test') { useJUnitPlatform { - excludeTags 'benchmark', 'minio-integration', 'claim-concurrency', 'local-e2e', 'vector-search-performance' + excludeTags 'benchmark', 'minio-integration', 'claim-concurrency', 'local-e2e', 'vector-search-performance', 'worker-indexing-throughput' } } @@ -152,6 +152,28 @@ tasks.register('vectorSearchPerformanceTest', Test) { outputs.upToDateWhen { false } } +tasks.register('workerIndexingThroughputTest', Test) { + group = 'verification' + description = '실제 PostgreSQL, MinIO와 BGE-M3에서 자동 Worker 전체 문서 인덱싱 처리량을 측정합니다.' + testClassesDirs = sourceSets.test.output.classesDirs + classpath = sourceSets.test.runtimeClasspath + useJUnitPlatform { + includeTags 'worker-indexing-throughput' + } + maxParallelForks = 1 + systemProperties System.properties.findAll { key, value -> + key.toString().startsWith('worker.indexing.throughput.') + } + if (System.getProperty('worker.indexing.throughput.output') == null) { + systemProperty( + 'worker.indexing.throughput.output', + layout.buildDirectory.file('reports/worker-indexing-throughput/worker-indexing-throughput.json').get().asFile.absolutePath + ) + } + // 실제 외부 Service를 점유하고 테스트 데이터를 초기화하는 Benchmark이므로 일반 Test와 Build Cache에서 분리한다. + outputs.upToDateWhen { false } +} + def configureOpenSqlDatabase = { Test task -> task.maxParallelForks = 1 task.outputs.upToDateWhen { false } diff --git a/docs/design/gimin-#133-worker-indexing-throughput-benchmark.md b/docs/design/gimin-#133-worker-indexing-throughput-benchmark.md new file mode 100644 index 0000000..040b948 --- /dev/null +++ b/docs/design/gimin-#133-worker-indexing-throughput-benchmark.md @@ -0,0 +1,176 @@ +# 자동 Worker 전체 문서 인덱싱 처리량 Benchmark 설계 + +## 1. 배경 + +자동 Worker는 Job 실행 슬롯을 먼저 확보한 뒤 `PENDING` Job을 Claim하고, Parsing, Chunk 저장, +Batch Embedding, Vector 저장과 검색 Version 전환까지 수행한다. 기존 통합·E2E 테스트는 상태 전이와 +소유권 정합성을 검증하지만 여러 문서를 연속 처리할 때의 처리량과 지연 기준선은 제공하지 않는다. + +이번 작업은 실제 PostgreSQL 17, pgvector, MinIO와 `BAAI/bge-m3`를 연결한 자동 Worker 전체 +Pipeline을 반복 측정해 다음 질문에 답한다. + +1. 기본 Worker 동시성에서 문서와 Chunk를 초당 얼마나 인덱싱하는가? +2. 업로드 접수, Queue 대기와 실제 처리 시간 중 어느 구간이 전체 지연을 지배하는가? +3. Queue를 모두 소진할 때 처리 누락, 중복 Attempt 또는 중복 Vector가 발생하지 않는가? + +## 2. 범위 + +### 2.1 포함 + +- 실제 인증과 Multipart HTTP 문서 업로드 +- 실제 MinIO Object 저장과 읽기 +- 실제 자동 Polling, Claim, Attempt, Lease 갱신과 전체 인덱싱 Pipeline +- 실제 BGE-M3 Batch Embedding과 PostgreSQL `vector(1024)` 저장 +- 결정적 TXT Corpus와 같은 Chunk 분포를 사용한 반복 측정 +- Warm-up과 본 측정 분리 +- 문서·Chunk·Embedding 처리량과 Queue·처리·전체 지연 수집 +- Profile별 상태, Attempt, Event, Chunk와 Vector 정합성 검증 +- PostgreSQL, pgvector, BGE, Batch Size와 Worker 설정 환경 지문 수집 +- Git 제외 JSON 원본 결과와 실행 결과 Markdown 기록 +- 일반 테스트와 분리된 전용 Gradle Task + +### 2.2 제외 + +- Worker 프로세스 수에 따른 수평 확장 비교 +- Queue 허용 한계와 Backpressure 붕괴 지점 측정 +- Lease 갱신 주기별 DB 부하 비교 +- BGE-M3 Batch Size 재선정 +- PDF·DOCX Parser 형식별 성능 비교 +- 공식 OpenSQL 공급사 환경의 최종 성능 수치 +- 운영 SLO 확정 + +## 3. Workload 계약 + +### 3.1 결정적 문서 + +- 모든 Profile은 같은 UTF-8 TXT Corpus Template을 사용한다. +- 문서마다 고유한 식별 문구만 바꾸고 본문 길이와 문단 구조는 동일하게 유지한다. +- 본문은 기본 Chunk Size와 Overlap에서 같은 수의 Chunk가 생성되도록 고정한다. +- 파일 이름과 Document 제목은 Profile, 반복과 문서 순번을 포함해 충돌을 방지한다. +- 전체 Profile의 실제 Chunk 수가 같지 않으면 처리량 비교를 실패로 처리한다. + +### 3.2 Profile + +기본값은 짧은 로컬 실행과 Queue가 유지되는 본 측정을 함께 제공한다. + +| 구분 | 문서 수 | 반복 | 통계 포함 | +|---|---:|---:|---| +| Warm-up | 4 | 1 | 제외 | +| 작은 Queue | 16 | 3 | 포함 | +| 지속 Queue | 32 | 3 | 포함 | + +문서 수와 반복은 `worker.indexing.throughput.*` System Property로 변경할 수 있다. Worker +`max-concurrency`는 제품 기본값인 `2`로 고정하고 결과 환경 지문에 기록한다. Worker 수평 확장은 +별도 Benchmark에서 다룬다. + +### 3.3 측정 구간 + +1. Warm-up 문서를 업로드하고 모두 `INDEXED`가 될 때까지 기다린다. +2. Profile 시작 시각을 기록하고 여러 Uploader Thread로 측정 문서를 접수한다. +3. 마지막 업로드 완료 시각을 기록한다. +4. 자동 Worker가 Profile의 모든 Job을 `INDEXED`로 전환할 때까지 기다린다. +5. 상태와 결과 정합성을 검증한 뒤 통계를 계산한다. + +Upload와 Worker 실행이 겹치는 실제 동작을 유지한다. 대신 업로드 시간과 마지막 업로드 뒤 Queue +소진 시간을 분리해 HTTP·MinIO 접수가 Worker 처리량을 가리는지 확인한다. + +## 4. 지표 계약 + +| 지표 | 계산 | +|---|---| +| documents/s | 완료 문서 수 / Profile 전체 경과 시간 | +| documents/min | documents/s × 60 | +| chunks/s | 저장 Chunk 수 / Profile 전체 경과 시간 | +| embeddings/s | 저장 Embedding 수 / Profile 전체 경과 시간 | +| Upload 경과 | 첫 Upload 시작부터 마지막 Upload 응답까지 | +| Queue 소진 | 마지막 Upload 응답부터 마지막 `INDEXED`까지 | +| Queue 대기 | Job `created_at`부터 `LOCKED` Event까지 | +| 실제 처리 | `LOCKED`부터 `INDEXED` Event까지 | +| 전체 Job 지연 | Job `created_at`부터 `INDEXED` Event까지 | + +지연 분포는 선형 보간 p50·p95·p99와 max를 밀리초로 기록한다. 모든 Profile 원본과 중앙값 요약은 +JSON으로 기록한다. + +## 5. 정합성 계약 + +각 Profile은 측정 뒤 다음 조건을 모두 검증한다. + +- 대상 Job 전부 `INDEXED` +- 전체 Schema에 `PENDING` 또는 `PROCESSING` 잔여 Job 없음 +- Job별 `SUCCESS` Attempt 정확히 1개 +- Job별 `LOCKED`, `PARSE_STARTED`, `CHUNKED`, `EMBEDDING_STARTED`, `INDEXED` 순서 유지 +- Document와 Document Version 모두 `INDEXED` +- `documents.current_version_id`가 측정 Version을 가리킴 +- 문서별 Chunk 수와 Embedding 수 일치 +- 같은 `(chunk_id, embedding_model_id)` 중복 없음 +- 모든 Vector 차원 1024 +- 실패 Attempt, Retry와 최종 실패 Event 없음 + +정합성 실패가 하나라도 발생하면 성능 숫자를 유효한 결과로 취급하지 않고 Test를 실패시킨다. + +## 6. 환경 검증과 결과 보존 + +Benchmark 시작 전에 다음을 확인한다. + +- PostgreSQL Server Version `17.x` +- pgvector Extension `0.8.1` +- 전용 Test Schema와 MinIO Bucket +- BGE Health와 Model명 `BAAI/bge-m3` +- Document Batch Size와 Worker 최대 동시성 + +접속 Secret, JWT와 Object Storage Credential은 결과에 기록하지 않는다. 원본 JSON은 다음 Git 제외 +경로에 생성한다. + +```text +build/reports/worker-indexing-throughput/worker-indexing-throughput.json +``` + +실행 환경, 중앙값, 해석과 한계만 `docs/test-results/`에 기록한다. + +## 7. 실행 경계 + +일반 `./gradlew test`는 Docker와 실제 BGE-M3에 의존하지 않는다. 전용 Task만 실제 Infrastructure를 +요구한다. + +```bash +docker compose up -d postgres minio embedding-server +./gradlew workerIndexingThroughputTest +``` + +확장 측정 예시는 다음과 같다. + +```bash +./gradlew workerIndexingThroughputTest \ + -Dworker.indexing.throughput.document-counts=32,64 \ + -Dworker.indexing.throughput.repetitions=5 +``` + +## 8. 실패 정책 + +- Infrastructure Health, DB Version, pgvector Version 또는 BGE Model 계약이 다르면 즉시 실패한다. +- Upload, Worker 처리, Embedding 또는 상태 Polling이 제한 시간을 넘으면 Job Snapshot을 포함해 실패한다. +- 완료된 Profile 결과는 후속 진단을 위해 JSON에 보존하되 실패 Profile은 중앙값 판단에서 제외한다. +- Profile 간 데이터는 같은 전용 Schema에서 고유 식별자로 격리하고 마지막에 Schema와 Bucket을 정리한다. + +## 9. 검증 + +- Percentile, 중앙값과 처리량 계산 단위 테스트 +- 실제 Infrastructure를 사용하는 작은 Smoke Benchmark +- 기본 Profile 전체 Benchmark +- 기존 로컬 문서 인덱싱 E2E 회귀 +- 전체 일반 Java 테스트 + +## 10. 커밋 분할 + +1. `docs: #133 자동 Worker 처리량 Benchmark 설계 추가` +2. `test: #133 자동 Worker 전체 인덱싱 처리량 Benchmark 추가` +3. `build: #133 Worker 처리량 전용 실행 경계 추가` +4. `perf: #133 자동 Worker 인덱싱 처리량 실측 결과 기록` + +## 11. 완료 조건 + +- 실제 자동 Worker 전체 Pipeline을 한 명령으로 반복 측정할 수 있다. +- 문서·Chunk·Embedding 처리량과 Queue·처리·전체 지연이 구조화돼 기록된다. +- 성능 Profile마다 완료·소유권·Attempt·Event·Vector 불변식을 검증한다. +- 일반 테스트는 실제 Infrastructure 없이 계속 실행된다. +- 측정 환경의 한계와 운영 수치로 해석하면 안 되는 범위를 결과 문서에 명시한다. diff --git a/docs/test-results/gimin-#133-worker-indexing-throughput-benchmark.md b/docs/test-results/gimin-#133-worker-indexing-throughput-benchmark.md new file mode 100644 index 0000000..fee0bae --- /dev/null +++ b/docs/test-results/gimin-#133-worker-indexing-throughput-benchmark.md @@ -0,0 +1,164 @@ +# Issue #133 자동 Worker 전체 문서 인덱싱 처리량 Benchmark 결과 + +## 1. 결과 요약 + +PostgreSQL 17.8, pgvector 0.8.1, MinIO와 실제 `BAAI/bge-m3`를 연결하고 자동 Polling Worker가 +문서 업로드부터 `INDEXED` 전환까지 처리하는 전체 경로를 측정했다. 4개 문서 예열 후 16개와 32개 +문서 Profile을 각각 3회 실행했다. + +| 문서 수 | 반복 | 중앙 총 시간 | 문서/초 | 문서/분 | Chunk·Embedding/초 | 전체 P95 | +|---:|---:|---:|---:|---:|---:|---:| +| 16 | 3회 | 50.255초 | 0.318 | 19.102 | 2.547 | 50.127초 | +| 32 | 3회 | 104.335초 | 0.307 | 18.402 | 2.454 | 100.665초 | + +Queue를 16개에서 32개로 두 배 늘렸을 때 분당 처리량은 약 3.7% 감소했다. Worker 실행 슬롯을 +2개로 고정했기 때문에 처리 시간 P95는 6.768초에서 7.285초로 비슷하게 유지됐고, 전체 지연 증가는 +주로 Queue 대기 P95가 43.794초에서 93.554초로 늘어난 데서 발생했다. + +## 2. 공개 가능한 실행 환경 + +| 항목 | 값 | +|---|---| +| Database | PostgreSQL 17.8, Local Docker | +| pgvector | 0.8.1 | +| Object Storage | MinIO, Local Docker | +| Embedding Provider | `BAAI/bge-m3`, Local Docker CPU 추론 | +| Vector 차원 | 1024 | +| Embedding Batch Size | 32 | +| Worker 동시 실행 슬롯 | 2 | +| Worker Polling 주기 | 50 ms | +| 업로더 Thread | 4 | +| 문서 크기 | TXT 6,400자 | +| 문서당 Chunk·Embedding | 각각 8개 | +| Warm-up | 4개 문서 | +| 본 측정 | 16 / 32개 문서, Profile당 3회 | +| Application | Spring Boot 3.5.16, Java 17 | +| 실행 장비 | macOS `aarch64`, 가용 Processor 10개 | +| 실행 일자 | 2026-08-10 KST | + +DB Host·Database 이름·Username·Password, MinIO Credential과 JWT 값은 결과에 기록하지 않았다. +이번 결과는 로컬 개발 장비의 기준선이며 공식 OpenSQL 원격 Server 성능이나 운영 SLO가 아니다. + +## 3. 측정 범위 + +각 문서는 다음 전체 경로를 통과했다. + +```text +TXT 업로드 +→ MinIO 저장 +→ Embedding Job PENDING +→ 자동 Worker Polling·Claim +→ Attempt 시작 +→ 텍스트 Parsing·Chunk 저장 +→ 실제 BGE-M3 Batch 호출 +→ vector(1024) 저장 +→ Version INDEXED +→ Document current_version 전환 +→ Worker 실행 슬롯 반환 +``` + +Profile마다 Benchmark 전용 Schema와 Bucket의 데이터를 초기화해 이전 실행의 Job·Chunk·Embedding이 +다음 실행의 수치에 포함되지 않게 했다. Worker Node는 Application 수명주기를 유지하기 위해 Profile +사이에서 재사용했다. + +## 4. 실행 방법 + +PostgreSQL, MinIO와 Embedding Server가 모두 건강 상태인 로컬 환경에서 실행했다. + +```bash +docker compose up -d postgres minio embedding-server +DB_SSLMODE=disable ./gradlew workerIndexingThroughputTest +``` + +구조화 원시 결과는 Git에 포함하지 않는 다음 경로에 생성된다. + +```text +build/reports/worker-indexing-throughput/worker-indexing-throughput.json +``` + +실환경 연결과 결과 계약을 빠르게 확인할 때는 다음 Smoke Profile을 사용했다. + +```bash +DB_SSLMODE=disable ./gradlew workerIndexingThroughputTest \ + -Dworker.indexing.throughput.warm-up-documents=2 \ + -Dworker.indexing.throughput.document-counts=4 \ + -Dworker.indexing.throughput.repetitions=1 \ + -Dworker.indexing.throughput.output=build/reports/worker-indexing-throughput/smoke.json +``` + +Smoke 결과는 4개 문서, 32개 Chunk, 32개 Embedding을 12.461초에 처리했고 분당 19.260문서를 +기록했다. + +## 5. 반복별 원시 결과 + +| 문서 수 | 회차 | 총 시간 | 문서/분 | Chunk·Embedding/초 | Queue P95 | 처리 P95 | 전체 P95 | +|---:|---:|---:|---:|---:|---:|---:|---:| +| 16 | 1 | 49.379초 | 19.442 | 2.592 | 42.914초 | 7.121초 | 49.134초 | +| 16 | 2 | 51.898초 | 18.498 | 2.466 | 44.887초 | 6.768초 | 51.742초 | +| 16 | 3 | 50.255초 | 19.102 | 2.547 | 43.794초 | 6.468초 | 50.127초 | +| 32 | 1 | 102.558초 | 18.721 | 2.496 | 92.526초 | 6.974초 | 98.769초 | +| 32 | 2 | 104.335초 | 18.402 | 2.454 | 93.554초 | 7.285초 | 100.665초 | +| 32 | 3 | 105.928초 | 18.125 | 2.417 | 95.968초 | 7.431초 | 102.253초 | + +같은 Profile 세 번의 분당 처리량 범위는 16문서에서 18.498~19.442, 32문서에서 +18.125~18.721이었다. 한 번의 최고값이 아니라 Profile별 중앙값을 비교 기준으로 사용했다. + +## 6. 데이터 완전성 검증 + +| 문서 수 | 회차 | Chunk 수 | Embedding 수 | 결과 | +|---:|---:|---:|---:|---| +| 16 | 1~3 | 매회 128 | 매회 128 | PASS | +| 32 | 1~3 | 매회 256 | 매회 256 | PASS | + +각 Profile 완료 시 다음 불변식을 함께 확인했다. + +- 모든 Embedding Job이 `INDEXED`다. +- 각 Job에 성공한 Attempt가 정확히 하나 존재한다. +- 모든 Document와 Version이 `INDEXED`이고 `current_version_id`가 측정 Version을 가리킨다. +- Chunk 수와 Embedding 수가 일치하고 중복 Chunk Embedding이 없다. +- 저장된 모든 Vector의 차원은 1024이며 NaN·Infinity가 없다. +- 자동 Worker 실행 슬롯이 0으로 반환되고 Worker가 살아 있다. + +따라서 이 결과는 HTTP 접수 시간만 측정한 값이 아니라 Vector 저장과 검색 Version 전환이 완료된 시점까지의 +전체 처리량이다. + +## 7. 실행 중 발견하고 해결한 환경 문제 + +### 7.1 로컬 PostgreSQL SSL 설정 + +첫 Smoke 시도는 외부 Shell의 SSL 설정이 비-SSL 로컬 PostgreSQL에 적용돼 Flyway 연결 전에 실패했다. +`DB_SSLMODE=disable`을 명시한 뒤 Migration과 전체 Pipeline이 정상 실행됐다. + +### 7.2 CPU BGE-M3 응답 제한 + +기존 Embedding HTTP 응답 제한은 5초로 고정돼 있었다. 로컬 CPU BGE-M3는 정상적인 Batch 응답에도 +5초를 넘겨 Worker가 `EMBEDDING_PROVIDER_UNAVAILABLE`로 재시도했다. 연결·응답 제한을 환경 설정으로 +분리하고 Benchmark에만 2분 응답 제한을 적용했다. 일반 실행의 기본값은 기존과 같은 5초다. + +### 7.3 전체 회귀의 JWT 환경 값 + +전체 회귀 첫 시도는 `JWT_SECRET` 미설정으로 Spring Context 12건이 연쇄 실패했다. Test 전용 JWT와 +`DB_SSLMODE=disable`을 명시한 재실행에서 720개 Test가 모두 통과했다. 첫 실패는 Source 결함이나 +Benchmark 실패가 아니며 통과 결과로 계산하지 않았다. + +## 8. 검증 결과 + +| 검증 | 결과 | +|---|---| +| 통계 계약 단위 테스트 | PASS | +| 전용 Gradle Task 노출 | PASS | +| 4문서 실제 BGE-M3 Smoke | PASS, 23초 | +| 16·32문서 각 3회 본 측정 | PASS, 8분 1초 | +| 전체 Java 회귀 | PASS, 720 tests, failure/error/skipped 0 | +| `git diff --check` | PASS | + +## 9. 결론과 남은 한계 + +- 구현됨: 자동 Worker의 실제 문서 인덱싱 처리량과 Queue·처리·전체 지연 분포를 반복 측정할 수 있다. +- 검증됨: 16문서에서 32문서로 Queue가 두 배가 되어도 분당 처리량 감소는 약 3.7%였다. +- 검증됨: 여섯 실행 모두 Chunk·Embedding·1024차원 Vector 완전성을 만족했다. +- 관찰됨: 동시 실행 슬롯 2개에서는 Queue 증가가 처리 시간보다 전체 P95를 지배했다. +- 한계: 단일 Local Apple Silicon CPU 장비의 결과로, 공식 OpenSQL Server나 GPU BGE-M3 결과가 아니다. +- 한계: TXT 6,400자 고정 입력이라 PDF·DOCX Parser 비용이나 다양한 문서 길이 분포를 대표하지 않는다. +- 한계: Worker 수·동시 실행 슬롯·Embedding Batch Size를 바꾼 수평 확장 비교는 이번 범위에 포함하지 않았다. +- 후속: 같은 Harness로 Worker 수·동시성 변화, Queue 적체와 Backpressure, 장애 주입 Profile을 비교해야 한다. diff --git a/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java b/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java index 343c499..6638c41 100644 --- a/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java +++ b/src/main/java/com/opensource/docgrid/global/config/EmbeddingServerConfig.java @@ -9,20 +9,37 @@ import org.springframework.http.client.JdkClientHttpRequestFactory; import org.springframework.web.client.RestClient; +/** + * Embedding Provider HTTP 연결 설정을 구성한다. + * + *

Provider 호출 계약과 오류 분류는 Domain Client가 담당하고, 이 설정은 연결·응답 제한 시간과 + * Base URL 같은 Transport 경계만 책임진다.

+ */ @Configuration public class EmbeddingServerConfig { @Value("${embedding.server.base-url}") private String baseUrl; + @Value("${embedding.server.connect-timeout:5s}") + private Duration connectTimeout; + + @Value("${embedding.server.read-timeout:5s}") + private Duration readTimeout; + + /** + * 실행 환경별 제한 시간을 적용한 Embedding Provider 전용 Client를 만든다. + */ @Bean("embeddingRestClient") public RestClient embeddingRestClient() { + // 1. 연결과 추론 응답 제한을 분리해 느린 CPU 추론 환경에서도 Timeout을 독립적으로 조정한다. HttpClient httpClient = HttpClient.newBuilder() - .connectTimeout(Duration.ofSeconds(5)) + .connectTimeout(connectTimeout) .build(); JdkClientHttpRequestFactory requestFactory = new JdkClientHttpRequestFactory(httpClient); - requestFactory.setReadTimeout(Duration.ofSeconds(5)); + requestFactory.setReadTimeout(readTimeout); + // 2. Domain Client가 같은 Transport 정책을 공유하도록 단일 이름의 RestClient를 제공한다. return RestClient.builder() .baseUrl(baseUrl) .requestFactory(requestFactory) diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 3821d5a..4e983c2 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -61,6 +61,8 @@ jwt: embedding: server: base-url: ${EMBEDDING_SERVER_URL:http://localhost:8000} + connect-timeout: ${EMBEDDING_SERVER_CONNECT_TIMEOUT:5s} + read-timeout: ${EMBEDDING_SERVER_READ_TIMEOUT:5s} document: # 실제 BGE-M3 CPU Benchmark에서 최고 처리량의 98.3%를 유지하면서 Batch 64보다 p95가 47% 낮은 균형점이다. batch-size: ${EMBEDDING_DOCUMENT_BATCH_SIZE:32} diff --git a/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java b/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java index 9a8db8d..3b5ccfa 100644 --- a/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java +++ b/src/test/java/com/opensource/docgrid/e2e/LocalE2eMinioBucket.java @@ -56,6 +56,17 @@ List objectKeys() throws Exception { return List.copyOf(keys); } + /** + * Bucket은 유지하고 현재 Test가 저장한 Object만 제거해 반복 Profile의 입력 상태를 초기화한다. + */ + void clear() throws Exception { + for (String objectKey : objectKeys()) { + minioClient.removeObject( + RemoveObjectArgs.builder().bucket(bucket).object(objectKey).build() + ); + } + } + @Override public void close() throws Exception { boolean exists = minioClient.bucketExists( @@ -66,11 +77,7 @@ public void close() throws Exception { } // 1. MinIO는 비어 있지 않은 Bucket 삭제를 거부하므로 실제 Object를 모두 먼저 제거한다. - for (String objectKey : objectKeys()) { - minioClient.removeObject( - RemoveObjectArgs.builder().bucket(bucket).object(objectKey).build() - ); - } + clear(); // 2. Test가 생성한 격리 Bucket만 삭제하고 다른 개발 Bucket은 조회하거나 변경하지 않는다. minioClient.removeBucket(RemoveBucketArgs.builder().bucket(bucket).build()); diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputBenchmark.java b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputBenchmark.java new file mode 100644 index 0000000..c7f5455 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputBenchmark.java @@ -0,0 +1,721 @@ +package com.opensource.docgrid.e2e; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.Duration; +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.Timeout; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.SpringBootVersion; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.web.client.TestRestTemplate; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; +import org.springframework.web.client.RestClient; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.opensource.docgrid.domain.embedding.config.EmbeddingBatchProperties; +import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool; +import com.opensource.docgrid.domain.worker.lifecycle.WorkerExecutionLifecycleManager; +import com.opensource.docgrid.domain.worker.lifecycle.WorkerJobPollingScheduler; +import com.opensource.docgrid.domain.worker.lifecycle.WorkerLifecycleManager; +import com.opensource.docgrid.e2e.LocalE2eApiClient.UploadedDocument; +import com.opensource.docgrid.e2e.LocalE2eDocumentFactory.DocumentPayload; + +import io.minio.MinioClient; +import lombok.extern.slf4j.Slf4j; + +/** + * 실제 PostgreSQL 17·pgvector·MinIO·BGE-M3에서 자동 Worker 전체 인덱싱 처리량을 측정한다. + * + *

결정적 TXT 문서를 실제 HTTP로 동시에 접수하고 Production Scheduler가 Queue를 소진하게 한다. + * 각 Profile은 처리량과 단계별 지연을 기록한 뒤 Job, Attempt, Event, Chunk와 Vector 정합성을 함께 + * 검증한다. 실제 Infrastructure를 요구하므로 일반 회귀 테스트와 분리된 Tag에서만 실행한다. + */ +@Slf4j +@Tag("worker-indexing-throughput") +@ActiveProfiles({"test", "minio-integration"}) +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) +@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@DisplayName("자동 Worker 전체 문서 인덱싱 처리량 Benchmark") +class WorkerIndexingThroughputBenchmark { + + private static final String EXECUTION_ID = UUID.randomUUID().toString().replace("-", ""); + // PostgreSQL 식별자 63자 제한 안에서 별도 Gradle 실행이 Schema를 공유하지 않도록 격리한다. + private static final String TEST_SCHEMA = "docgrid_worker_indexing_throughput_" + + EXECUTION_ID.substring(0, 24); + private static final String TEST_BUCKET = "docgrid-worker-throughput-" + EXECUTION_ID; + private static final String EXPECTED_POSTGRES_VERSION_PREFIX = "17."; + private static final String EXPECTED_PGVECTOR_VERSION = "0.8.1"; + private static final String EXPECTED_MODEL = "BAAI/bge-m3"; + private static final int EXPECTED_VECTOR_DIMENSION = 1024; + private static final int DOCUMENT_CHARACTER_COUNT = positiveIntegerProperty( + "worker.indexing.throughput.document-characters", + 6_400 + ); + private static final int WARM_UP_DOCUMENT_COUNT = positiveIntegerProperty( + "worker.indexing.throughput.warm-up-documents", + 4 + ); + private static final List DOCUMENT_COUNTS = positiveIntegerListProperty( + "worker.indexing.throughput.document-counts", + List.of(16, 32) + ); + private static final int REPETITIONS = positiveIntegerProperty( + "worker.indexing.throughput.repetitions", + 3 + ); + private static final int UPLOADER_THREADS = positiveIntegerProperty( + "worker.indexing.throughput.uploader-threads", + 4 + ); + private static final long PROFILE_TIMEOUT_SECONDS = positiveLongProperty( + "worker.indexing.throughput.profile-timeout-seconds", + 600L + ); + private static final long POLLING_SLEEP_MILLIS = positiveLongProperty( + "worker.indexing.throughput.status-polling-ms", + 100L + ); + private static final Path OUTPUT_PATH = Path.of(System.getProperty( + "worker.indexing.throughput.output", + "build/reports/worker-indexing-throughput/worker-indexing-throughput.json" + )); + + @Autowired private TestRestTemplate restTemplate; + @Autowired private JdbcTemplate jdbcTemplate; + @Autowired private MinioClient minioClient; + @Autowired private ObjectMapper objectMapper; + @Autowired private WorkerLifecycleManager workerLifecycleManager; + @Autowired private WorkerJobPollingScheduler pollingScheduler; + @Autowired private WorkerExecutionLifecycleManager executionLifecycleManager; + @Autowired private WorkerExecutionSlotPool executionSlotPool; + @Autowired private IndexingWorkerProperties workerProperties; + @Autowired private EmbeddingBatchProperties embeddingBatchProperties; + @Autowired @Qualifier("embeddingRestClient") private RestClient embeddingRestClient; + + private LocalE2eApiClient apiClient; + private LocalE2eMinioBucket minioBucket; + private String accessToken; + + @DynamicPropertySource + static void configureEnvironment(DynamicPropertyRegistry registry) { + registry.add("TEST_DB_SCHEMA", () -> TEST_SCHEMA); + registry.add("jwt.secret", () -> "docgrid-worker-indexing-throughput-test-secret-key-2026"); + registry.add("minio.bucket", () -> TEST_BUCKET); + registry.add("indexing.worker.enabled", () -> "true"); + registry.add("indexing.worker.name", () -> "worker-indexing-throughput-benchmark"); + registry.add("indexing.worker.polling-interval", () -> "50ms"); + registry.add("indexing.worker.heartbeat-interval", () -> "1s"); + registry.add("indexing.worker.dead-threshold", () -> "2m"); + registry.add("indexing.worker.max-concurrency", () -> "2"); + registry.add("indexing.worker.lease-duration", () -> "2m"); + registry.add("indexing.worker.lease-renewal-interval", () -> "10s"); + registry.add("indexing.worker.lease-recovery-interval", () -> "10m"); + registry.add("indexing.worker.shutdown-grace-period", () -> "30s"); + // CPU 기반 BGE-M3의 실측 추론 시간을 5초 기본 운영 Timeout과 분리한다. + registry.add("embedding.server.read-timeout", () -> "2m"); + } + + @BeforeAll + void setUpInfrastructure() throws Exception { + apiClient = new LocalE2eApiClient(restTemplate); + minioBucket = new LocalE2eMinioBucket(minioClient, TEST_BUCKET); + minioBucket.create(); + accessToken = apiClient.loginAdmin(); + awaitCondition( + "Worker가 Application Ready 뒤 등록되지 않았습니다.", + () -> workerLifecycleManager.getWorkerId().isPresent() + ); + } + + @AfterAll + void cleanUpInfrastructure() throws Exception { + // 1. Scheduler와 실행 Thread를 먼저 닫아 Profile 자원 정리 뒤 DB 접근이 재개되지 않게 한다. + pollingScheduler.stopPolling(); + executionLifecycleManager.shutdown(); + workerLifecycleManager.stopWorker(); + + // 2. Benchmark 전용 Bucket과 Schema만 제거해 기존 로컬 개발 Data를 보존한다. + minioBucket.close(); + jdbcTemplate.execute("DROP SCHEMA IF EXISTS " + TEST_SCHEMA + " CASCADE"); + } + + @Test + @Timeout(1_800) + @DisplayName("실제 자동 Worker의 처리량과 Queue·처리·전체 지연을 반복 측정한다") + void measureAutomaticWorkerIndexingThroughput() throws Exception { + // 1. 잘못된 DB나 외부 Service에서 얻은 수치를 기록하지 않도록 환경 계약을 먼저 확인한다. + EnvironmentFingerprint environment = validateEnvironment(); + List runs = new ArrayList<>(); + writeReport(new BenchmarkReport(environment, runs, List.of())); + logJson("WORKER_INDEXING_THROUGHPUT_ENV", environment); + + // 2. Migration, HTTP, MinIO, Model과 주요 DB 경로를 예열하되 측정 결과에는 포함하지 않는다. + runWarmUp(); + + // 3. 같은 문서 분포로 Queue 크기와 반복만 바꿔 Profile 원시 결과를 수집한다. + for (int documentCount : DOCUMENT_COUNTS) { + for (int repetition = 1; repetition <= REPETITIONS; repetition++) { + resetProfileState(); + ProfileRun run = runMeasuredProfile(documentCount, repetition); + runs.add(run); + logJson("WORKER_INDEXING_THROUGHPUT_RESULT", run); + writeReport(new BenchmarkReport(environment, runs, medians(runs))); + } + } + + // 4. 장비 내 변동을 줄인 Profile별 중앙값을 최종 결과와 Log에 별도로 남긴다. + List medians = medians(runs); + for (ProfileMedian median : medians) { + logJson("WORKER_INDEXING_THROUGHPUT_MEDIAN", median); + } + writeReport(new BenchmarkReport(environment, runs, medians)); + } + + private EnvironmentFingerprint validateEnvironment() { + String postgresVersion = jdbcTemplate.queryForObject("SHOW server_version", String.class); + String pgvectorVersion = jdbcTemplate.queryForObject( + "SELECT extversion FROM pg_extension WHERE extname = 'vector'", + String.class + ); + Map model = jdbcTemplate.queryForMap( + "SELECT model_name, dimension FROM embedding_models " + + "WHERE is_active = TRUE AND is_searchable = TRUE" + ); + JsonNode health = embeddingRestClient.get() + .uri("/health") + .retrieve() + .body(JsonNode.class); + + assertThat(jdbcTemplate.queryForObject("SELECT current_schema()", String.class)) + .isEqualTo(TEST_SCHEMA); + assertThat(postgresVersion).startsWith(EXPECTED_POSTGRES_VERSION_PREFIX); + assertThat(pgvectorVersion).isEqualTo(EXPECTED_PGVECTOR_VERSION); + assertThat(model.get("model_name").toString()).isEqualTo(EXPECTED_MODEL); + assertThat(((Number) model.get("dimension")).intValue()).isEqualTo(EXPECTED_VECTOR_DIMENSION); + assertThat(health).isNotNull(); + assertThat(health.path("status").asText()).isEqualTo("ok"); + assertThat(workerProperties.getMaxConcurrency()).isEqualTo(2); + + return new EnvironmentFingerprint( + postgresVersion, + pgvectorVersion, + SpringBootVersion.getVersion(), + EXPECTED_MODEL, + EXPECTED_VECTOR_DIMENSION, + embeddingBatchProperties.getBatchSize(), + workerProperties.getMaxConcurrency(), + workerProperties.getPollingInterval().toString(), + UPLOADER_THREADS, + WARM_UP_DOCUMENT_COUNT, + DOCUMENT_COUNTS, + REPETITIONS, + DOCUMENT_CHARACTER_COUNT, + System.getProperty("os.name"), + System.getProperty("os.arch"), + Runtime.getRuntime().availableProcessors() + ); + } + + private void runWarmUp() throws Exception { + resetProfileState(); + List uploads = uploadDocuments("warm-up", WARM_UP_DOCUMENT_COUNT); + awaitIndexedAndIdle(uploads, "Warm-up"); + assertProfileInvariants(uploads); + log.info( + "자동 Worker 처리량 Benchmark 예열을 완료했습니다. documentCount={}", + WARM_UP_DOCUMENT_COUNT + ); + } + + private ProfileRun runMeasuredProfile(int documentCount, int repetition) throws Exception { + String profileName = "documents-" + documentCount + "-run-" + repetition; + long profileStartedAt = System.nanoTime(); + List uploads = uploadDocuments(profileName, documentCount); + long uploadCompletedAt = System.nanoTime(); + + awaitIndexedAndIdle(uploads, profileName); + long profileCompletedAt = System.nanoTime(); + ProfileData profileData = assertProfileInvariants(uploads); + + double elapsedSeconds = seconds(profileCompletedAt - profileStartedAt); + double uploadSeconds = seconds(uploadCompletedAt - profileStartedAt); + double queueDrainSeconds = seconds(profileCompletedAt - uploadCompletedAt); + return new ProfileRun( + documentCount, + repetition, + profileData.chunkCount(), + profileData.embeddingCount(), + elapsedSeconds, + uploadSeconds, + queueDrainSeconds, + documentCount / elapsedSeconds, + documentCount / elapsedSeconds * 60.0, + profileData.chunkCount() / elapsedSeconds, + profileData.embeddingCount() / elapsedSeconds, + LatencySummary.from(profileData.timings().stream().map(JobTiming::queueWaitMillis).toList()), + LatencySummary.from(profileData.timings().stream().map(JobTiming::processingMillis).toList()), + LatencySummary.from(profileData.timings().stream().map(JobTiming::endToEndMillis).toList()) + ); + } + + private List uploadDocuments(String profileName, int documentCount) throws Exception { + int threadCount = Math.min(UPLOADER_THREADS, documentCount); + ExecutorService uploader = Executors.newFixedThreadPool(threadCount); + List> futures = new ArrayList<>(documentCount); + + try { + // 1. Upload를 병렬 제출해 Worker 슬롯이 유지될 만큼 빠르게 PENDING Queue를 만든다. + for (int index = 0; index < documentCount; index++) { + int documentIndex = index; + futures.add(uploader.submit(() -> apiClient.upload( + accessToken, + throughputDocument(profileName, documentIndex) + ))); + } + uploader.shutdown(); + + // 2. 각 HTTP 응답을 제한 시간 안에서 수집하고 원래 문서 순서를 보존한다. + List uploads = new ArrayList<>(documentCount); + for (Future future : futures) { + uploads.add(future.get(PROFILE_TIMEOUT_SECONDS, TimeUnit.SECONDS)); + } + return List.copyOf(uploads); + } finally { + uploader.shutdownNow(); + uploader.awaitTermination(10, TimeUnit.SECONDS); + } + } + + private DocumentPayload throughputDocument(String profileName, int documentIndex) { + String marker = String.format( + Locale.ROOT, + "DocGrid automatic worker throughput document %04d. ", + documentIndex + ); + String sentence = "Automatic workers claim pending jobs, parse deterministic text, " + + "store chunks, generate BGE-M3 vectors, and switch the searchable version atomically. "; + StringBuilder body = new StringBuilder(DOCUMENT_CHARACTER_COUNT); + body.append(marker); + while (body.length() < DOCUMENT_CHARACTER_COUNT) { + body.append(sentence); + } + body.setLength(DOCUMENT_CHARACTER_COUNT); + + String fileName = profileName + "-" + String.format(Locale.ROOT, "%04d", documentIndex) + ".txt"; + return LocalE2eDocumentFactory.text(fileName, "Throughput " + fileName, body.toString()); + } + + private void awaitIndexedAndIdle(List uploads, String profileName) + throws InterruptedException { + awaitCondition( + profileName + " Profile이 완료되지 않았습니다. " + jobSnapshot(uploads), + () -> indexedJobCount(uploads) == uploads.size() && executionSlotPool.getActiveSlots() == 0 + ); + } + + private ProfileData assertProfileInvariants(List uploads) { + List timings = new ArrayList<>(uploads.size()); + List chunkCounts = new ArrayList<>(uploads.size()); + int totalEmbeddings = 0; + + for (UploadedDocument upload : uploads) { + // 1. Job과 검색 Version 전이가 끝났고 재시도 없이 한 Attempt만 성공했는지 확인한다. + assertThat(queryString("SELECT status FROM embedding_jobs WHERE id = ?", upload.embeddingJobId())) + .isEqualTo("INDEXED"); + assertThat(queryInteger( + "SELECT retry_count FROM embedding_jobs WHERE id = ?", + upload.embeddingJobId() + )).isZero(); + assertThat(count( + "SELECT COUNT(*) FROM embedding_job_attempts WHERE embedding_job_id = ?", + upload.embeddingJobId() + )).isOne(); + assertThat(count( + "SELECT COUNT(*) FROM embedding_job_attempts " + + "WHERE embedding_job_id = ? AND status = 'SUCCESS'", + upload.embeddingJobId() + )).isOne(); + assertThat(queryString("SELECT status FROM documents WHERE id = ?", upload.documentId())) + .isEqualTo("INDEXED"); + assertThat(queryLong("SELECT current_version_id FROM documents WHERE id = ?", upload.documentId())) + .isEqualTo(upload.documentVersionId()); + assertThat(queryString( + "SELECT status FROM document_versions WHERE id = ?", + upload.documentVersionId() + )).isEqualTo("INDEXED"); + + // 2. Chunk와 Embedding Set이 일대일이며 모든 Vector가 1024차원인지 확인한다. + int chunkCount = count( + "SELECT COUNT(*) FROM document_chunks WHERE document_version_id = ?", + upload.documentVersionId() + ); + int embeddingCount = count( + "SELECT COUNT(*) FROM embeddings WHERE document_version_id = ?", + upload.documentVersionId() + ); + assertThat(chunkCount).isPositive(); + assertThat(embeddingCount).isEqualTo(chunkCount); + assertThat(count( + "SELECT COUNT(DISTINCT chunk_id) FROM embeddings WHERE document_version_id = ?", + upload.documentVersionId() + )).isEqualTo(embeddingCount); + assertThat(queryInteger( + "SELECT MIN(vector_dims(vector)) FROM embeddings WHERE document_version_id = ?", + upload.documentVersionId() + )).isEqualTo(EXPECTED_VECTOR_DIMENSION); + assertThat(queryInteger( + "SELECT MAX(vector_dims(vector)) FROM embeddings WHERE document_version_id = ?", + upload.documentVersionId() + )).isEqualTo(EXPECTED_VECTOR_DIMENSION); + chunkCounts.add(chunkCount); + totalEmbeddings += embeddingCount; + + // 3. 정상 Event 순서와 실패·Retry 부재를 확인한 뒤 단계별 시각을 수집한다. + List events = jdbcTemplate.queryForList( + "SELECT event_type FROM indexing_events WHERE embedding_job_id = ? ORDER BY id", + String.class, + upload.embeddingJobId() + ); + assertThat(events).containsSubsequence( + "LOCKED", + "PARSE_STARTED", + "CHUNKED", + "EMBEDDING_STARTED", + "INDEXED" + ); + assertThat(events).doesNotContain( + "PARSE_FAILED", + "EMBEDDING_FAILED", + "LEASE_EXPIRED", + "FAILED", + "RETRY", + "MANUAL_RETRY" + ); + timings.add(readJobTiming(upload.embeddingJobId())); + } + + // 4. 같은 Profile의 문서별 Chunk 분포와 전체 Queue 소진 상태를 확인한다. + assertThat(chunkCounts).allMatch(count -> count.equals(chunkCounts.get(0))); + assertThat(count( + "SELECT COUNT(*) FROM embedding_jobs WHERE status IN ('PENDING', 'PROCESSING')" + )).isZero(); + assertThat(workerLifecycleManager.getWorkerId()).isPresent(); + return new ProfileData( + chunkCounts.stream().mapToInt(Integer::intValue).sum(), + totalEmbeddings, + List.copyOf(timings) + ); + } + + private JobTiming readJobTiming(Long jobId) { + return jdbcTemplate.queryForObject( + """ + SELECT job.created_at AS job_created_at, + MIN(event.occurred_at) FILTER (WHERE event.event_type = 'LOCKED') AS first_locked_at, + MIN(event.occurred_at) FILTER (WHERE event.event_type = 'INDEXED') AS indexed_at + FROM embedding_jobs job + JOIN indexing_events event ON event.embedding_job_id = job.id + WHERE job.id = ? + GROUP BY job.id, job.created_at + """, + (resultSet, rowNumber) -> { + LocalDateTime createdAt = resultSet.getTimestamp("job_created_at").toLocalDateTime(); + LocalDateTime lockedAt = resultSet.getTimestamp("first_locked_at").toLocalDateTime(); + LocalDateTime indexedAt = resultSet.getTimestamp("indexed_at").toLocalDateTime(); + double queueWaitMillis = millis(Duration.between(createdAt, lockedAt)); + double processingMillis = millis(Duration.between(lockedAt, indexedAt)); + double endToEndMillis = millis(Duration.between(createdAt, indexedAt)); + assertThat(queueWaitMillis).isGreaterThanOrEqualTo(0.0); + assertThat(processingMillis).isGreaterThanOrEqualTo(0.0); + assertThat(endToEndMillis).isGreaterThanOrEqualTo(processingMillis); + return new JobTiming(queueWaitMillis, processingMillis, endToEndMillis); + }, + jobId + ); + } + + private void resetProfileState() throws Exception { + awaitCondition( + "이전 Profile의 Worker 실행 슬롯이 반환되지 않았습니다.", + () -> executionSlotPool.getActiveSlots() == 0 + ); + minioBucket.clear(); + jdbcTemplate.execute(""" + TRUNCATE TABLE + indexing_events, + embedding_job_attempts, + embeddings, + document_chunks, + embedding_jobs, + document_versions, + documents, + file_objects + RESTART IDENTITY CASCADE + """); + } + + private int indexedJobCount(List uploads) { + String placeholders = String.join(",", uploads.stream().map(upload -> "?").toList()); + Object[] jobIds = uploads.stream().map(UploadedDocument::embeddingJobId).toArray(); + return count( + "SELECT COUNT(*) FROM embedding_jobs WHERE status = 'INDEXED' AND id IN (" + placeholders + ")", + jobIds + ); + } + + private String jobSnapshot(List uploads) { + String placeholders = String.join(",", uploads.stream().map(upload -> "?").toList()); + Object[] jobIds = uploads.stream().map(UploadedDocument::embeddingJobId).toArray(); + return jdbcTemplate.queryForList( + "SELECT id || ':' || status || ':retry=' || retry_count || ':error=' " + + "|| COALESCE(error_code, 'none') FROM embedding_jobs WHERE id IN (" + placeholders + ") " + + "ORDER BY id", + String.class, + jobIds + ).toString(); + } + + private List medians(List runs) { + List medians = new ArrayList<>(); + for (int documentCount : DOCUMENT_COUNTS) { + List profileRuns = runs.stream() + .filter(run -> run.documentCount() == documentCount) + .toList(); + if (profileRuns.isEmpty()) { + continue; + } + medians.add(new ProfileMedian( + documentCount, + profileRuns.size(), + median(profileRuns, ProfileRun::documentsPerSecond), + median(profileRuns, ProfileRun::documentsPerMinute), + median(profileRuns, ProfileRun::chunksPerSecond), + median(profileRuns, ProfileRun::embeddingsPerSecond), + median(profileRuns, ProfileRun::totalElapsedSeconds), + median(profileRuns, ProfileRun::uploadElapsedSeconds), + median(profileRuns, ProfileRun::queueDrainSeconds), + median(profileRuns, run -> run.queueWait().p95Millis()), + median(profileRuns, run -> run.processing().p95Millis()), + median(profileRuns, run -> run.endToEnd().p95Millis()) + )); + } + return List.copyOf(medians); + } + + private double median(List runs, ProfileMetric metric) { + return WorkerIndexingThroughputStatistics.median(runs.stream().map(metric::value).toList()); + } + + private void writeReport(BenchmarkReport report) throws IOException { + Path parent = OUTPUT_PATH.toAbsolutePath().getParent(); + if (parent != null) { + Files.createDirectories(parent); + } + objectMapper.writerWithDefaultPrettyPrinter().writeValue(OUTPUT_PATH.toFile(), report); + } + + private void logJson(String prefix, Object value) throws JsonProcessingException { + log.info("{} {}", prefix, objectMapper.writeValueAsString(value)); + } + + private void awaitCondition(String failureMessage, CheckedCondition condition) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(PROFILE_TIMEOUT_SECONDS); + while (System.nanoTime() < deadline) { + if (condition.evaluate()) { + return; + } + Thread.sleep(POLLING_SLEEP_MILLIS); + } + throw new AssertionError(failureMessage); + } + + private int count(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, Integer.class, arguments); + } + + private int queryInteger(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, Integer.class, arguments); + } + + private Long queryLong(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, Long.class, arguments); + } + + private String queryString(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, String.class, arguments); + } + + private static double seconds(long nanoseconds) { + return nanoseconds / 1_000_000_000.0; + } + + private static double millis(Duration duration) { + return duration.toNanos() / 1_000_000.0; + } + + private static int positiveIntegerProperty(String name, int defaultValue) { + String value = System.getProperty(name); + int parsed = value == null ? defaultValue : Integer.parseInt(value.trim()); + if (parsed < 1) { + throw new IllegalArgumentException(name + "은 1 이상이어야 합니다."); + } + return parsed; + } + + private static long positiveLongProperty(String name, long defaultValue) { + String value = System.getProperty(name); + long parsed = value == null ? defaultValue : Long.parseLong(value.trim()); + if (parsed < 1L) { + throw new IllegalArgumentException(name + "은 1 이상이어야 합니다."); + } + return parsed; + } + + private static List positiveIntegerListProperty(String name, List defaults) { + String value = System.getProperty(name); + if (value == null || value.isBlank()) { + return defaults; + } + List parsed = List.of(value.split(",")).stream() + .map(String::trim) + .map(Integer::parseInt) + .distinct() + .toList(); + if (parsed.isEmpty() || parsed.stream().anyMatch(item -> item < 1)) { + throw new IllegalArgumentException(name + "에는 1 이상의 정수만 사용할 수 있습니다."); + } + return parsed; + } + + /** 실행 환경과 Workload 설정을 Secret 없이 재현할 수 있는 지문이다. */ + private record EnvironmentFingerprint( + String postgresVersion, + String pgvectorVersion, + String springBootVersion, + String embeddingModel, + int vectorDimension, + int embeddingBatchSize, + int workerMaxConcurrency, + String workerPollingInterval, + int uploaderThreads, + int warmUpDocumentCount, + List documentCounts, + int repetitions, + int documentCharacterCount, + String osName, + String osArchitecture, + int availableProcessors + ) { + } + + /** 한 Profile 반복의 처리량과 단계별 지연 원시 결과다. */ + private record ProfileRun( + int documentCount, + int repetition, + int chunkCount, + int embeddingCount, + double totalElapsedSeconds, + double uploadElapsedSeconds, + double queueDrainSeconds, + double documentsPerSecond, + double documentsPerMinute, + double chunksPerSecond, + double embeddingsPerSecond, + LatencySummary queueWait, + LatencySummary processing, + LatencySummary endToEnd + ) { + } + + /** 같은 문서 수 Profile의 반복 결과를 중앙값으로 요약한 비교 단위다. */ + private record ProfileMedian( + int documentCount, + int completedRepetitions, + double documentsPerSecond, + double documentsPerMinute, + double chunksPerSecond, + double embeddingsPerSecond, + double totalElapsedSeconds, + double uploadElapsedSeconds, + double queueDrainSeconds, + double queueWaitP95Millis, + double processingP95Millis, + double endToEndP95Millis + ) { + } + + /** 한 지연 분포의 중앙값, 꼬리 지연과 최댓값을 밀리초로 보존한다. */ + private record LatencySummary( + double p50Millis, + double p95Millis, + double p99Millis, + double maxMillis + ) { + private static LatencySummary from(List samples) { + return new LatencySummary( + WorkerIndexingThroughputStatistics.percentile(samples, 50.0), + WorkerIndexingThroughputStatistics.percentile(samples, 95.0), + WorkerIndexingThroughputStatistics.percentile(samples, 99.0), + samples.stream().mapToDouble(Double::doubleValue).max().orElseThrow() + ); + } + } + + /** Profile 검증에서 수집한 전체 Chunk·Embedding 수와 Job 단계별 시각이다. */ + private record ProfileData(int chunkCount, int embeddingCount, List timings) { + } + + /** 한 Job의 Queue, 실제 처리와 전체 지연을 밀리초로 표현한다. */ + private record JobTiming(double queueWaitMillis, double processingMillis, double endToEndMillis) { + } + + /** 완료된 Profile과 중앙값을 환경 지문과 함께 JSON으로 저장하는 최상위 결과다. */ + private record BenchmarkReport( + EnvironmentFingerprint environment, + List runs, + List medians + ) { + } + + /** Profile 중앙값을 계산할 실수 지표를 선택한다. */ + @FunctionalInterface + private interface ProfileMetric { + double value(ProfileRun run); + } + + /** 제한 시간 동안 반복 평가할 DB·Worker 완료 조건이다. */ + @FunctionalInterface + private interface CheckedCondition { + boolean evaluate(); + } +} diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatistics.java b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatistics.java new file mode 100644 index 0000000..1f2b0cf --- /dev/null +++ b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatistics.java @@ -0,0 +1,47 @@ +package com.opensource.docgrid.e2e; + +import java.util.ArrayList; +import java.util.List; + +/** + * Worker 처리량 Benchmark의 Percentile과 중앙값을 결정적으로 계산하는 Test 전용 통계 경계다. + * + *

측정 I/O와 분리된 순수 계산만 제공해 실제 Infrastructure 없이도 통계 계약을 회귀 검증한다. + */ +final class WorkerIndexingThroughputStatistics { + + private WorkerIndexingThroughputStatistics() { + } + + static double median(List values) { + return percentile(values, 50.0); + } + + static double percentile(List values, double percentile) { + if (values == null || values.isEmpty()) { + throw new IllegalArgumentException("Percentile 입력은 비어 있을 수 없습니다."); + } + if (!Double.isFinite(percentile) || percentile < 0.0 || percentile > 100.0) { + throw new IllegalArgumentException("Percentile은 0 이상 100 이하여야 합니다."); + } + + List sorted = new ArrayList<>(values.size()); + for (Double value : values) { + if (value == null || !Double.isFinite(value)) { + throw new IllegalArgumentException("Percentile 입력은 유한한 수여야 합니다."); + } + sorted.add(value); + } + sorted.sort(Double::compareTo); + + double rank = percentile / 100.0 * (sorted.size() - 1); + int lowerIndex = (int) Math.floor(rank); + int upperIndex = (int) Math.ceil(rank); + if (lowerIndex == upperIndex) { + return sorted.get(lowerIndex); + } + + double weight = rank - lowerIndex; + return sorted.get(lowerIndex) * (1.0 - weight) + sorted.get(upperIndex) * weight; + } +} diff --git a/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatisticsTest.java b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatisticsTest.java new file mode 100644 index 0000000..30e5b04 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/e2e/WorkerIndexingThroughputStatisticsTest.java @@ -0,0 +1,41 @@ +package com.opensource.docgrid.e2e; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.util.List; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +/** + * Worker 처리량 Benchmark가 사용하는 선형 보간 Percentile과 입력 검증 계약을 확인한다. + */ +@DisplayName("Worker 인덱싱 처리량 통계 단위 테스트") +class WorkerIndexingThroughputStatisticsTest { + + @Test + @DisplayName("정렬되지 않은 짝수 표본의 중앙값과 p95를 선형 보간한다") + void percentile_interpolatesUnsortedEvenSamples() { + List samples = List.of(40.0, 10.0, 30.0, 20.0); + + assertThat(WorkerIndexingThroughputStatistics.median(samples)).isEqualTo(25.0); + assertThat(WorkerIndexingThroughputStatistics.percentile(samples, 95.0)).isEqualTo(38.5); + } + + @Test + @DisplayName("단일 표본은 모든 Percentile에서 같은 값을 반환한다") + void percentile_returnsOnlySample() { + assertThat(WorkerIndexingThroughputStatistics.percentile(List.of(17.0), 99.0)) + .isEqualTo(17.0); + } + + @Test + @DisplayName("빈 표본이나 유효 범위 밖 Percentile은 거부한다") + void percentile_rejectsInvalidInput() { + assertThatThrownBy(() -> WorkerIndexingThroughputStatistics.percentile(List.of(), 50.0)) + .isInstanceOf(IllegalArgumentException.class); + assertThatThrownBy(() -> WorkerIndexingThroughputStatistics.percentile(List.of(1.0), 101.0)) + .isInstanceOf(IllegalArgumentException.class); + } +}