diff --git a/.github/workflows/cicd-test.yml b/.github/workflows/cicd-test.yml new file mode 100644 index 0000000..4c6e97f --- /dev/null +++ b/.github/workflows/cicd-test.yml @@ -0,0 +1,51 @@ +name: CI/CD Test + +on: + # Manual trigger: needed to (re)deploy manifest-only changes, since k8s/** is + # path-ignored below and therefore never auto-triggers build+deploy. + workflow_dispatch: + push: + branches: [feat/ci] + paths-ignore: + - '**.md' + pull_request: + branches: [feat/ci] + paths-ignore: + - '**.md' + +jobs: + build: + uses: 404-Factory/.github/.github/workflows/reusable-build.yml@dev + permissions: + contents: read + packages: write + with: + use-ecr: true + aws-region: ap-northeast-2 + ecr-registry: ${{ vars.ECR_REGISTRY }} + platforms: ${{ vars.PLATFORM }} # 레포/조직 Variables에서 관리 (예: linux/arm64) + java-version: "17" + pre-build-command: "./gradlew build -x test --no-daemon" + push-on-ref: refs/heads/feat/ci + secrets: + aws-access-key-id: ${{ secrets.AWS_ACCESS_KEY_ID }} + aws-secret-access-key: ${{ secrets.AWS_SECRET_ACCESS_KEY }} + gpr-token: ${{ secrets.GPR_TOKEN }} + + deploy: + needs: build + if: (github.event_name == 'push' || github.event_name == 'workflow_dispatch') && github.ref == 'refs/heads/feat/ci' + uses: 404-Factory/.github/.github/workflows/reusable-deploy.yml@dev + permissions: + contents: read + with: + eks-cluster-name: sigma-eks-cluster + deploy-environment: production + aws-region: ap-northeast-2 + image: ${{ needs.build.outputs.image }} + image-tag: ${{ needs.build.outputs.tag }} + manifest-dir: k8s/app + namespace: analysis + secrets: + aws-access-key-id: ${{ secrets.AWS_ACCESS_KEY_ID }} + aws-secret-access-key: ${{ secrets.AWS_SECRET_ACCESS_KEY }} diff --git a/.github/workflows/cicd.yml b/.github/workflows/cicd.yml new file mode 100644 index 0000000..727096f --- /dev/null +++ b/.github/workflows/cicd.yml @@ -0,0 +1,52 @@ +name: CI/CD + +on: + # Manual trigger: needed to (re)deploy manifest-only changes, since k8s/** is + # path-ignored below and therefore never auto-triggers build+deploy. + workflow_dispatch: + push: + branches: [main, dev] + paths-ignore: + - 'k8s/**' + - '**.md' + pull_request: + branches: [main, dev] + paths-ignore: + - 'k8s/**' + - '**.md' + +jobs: + build: + uses: 404-Factory/.github/.github/workflows/reusable-build.yml@dev + permissions: + contents: read + packages: write + with: + use-ecr: true + aws-region: ap-northeast-2 + ecr-registry: ${{ vars.ECR_REGISTRY }} + platforms: ${{ vars.PLATFORM }} # 레포/조직 Variables에서 관리 (예: linux/arm64) + java-version: "17" + pre-build-command: "./gradlew build -x test --no-daemon" + secrets: + aws-access-key-id: ${{ secrets.AWS_ACCESS_KEY_ID }} + aws-secret-access-key: ${{ secrets.AWS_SECRET_ACCESS_KEY }} + gpr-token: ${{ secrets.GPR_TOKEN }} + + deploy: + needs: build + if: (github.event_name == 'push' || github.event_name == 'workflow_dispatch') && github.ref == 'refs/heads/main' + uses: 404-Factory/.github/.github/workflows/reusable-deploy.yml@dev + permissions: + contents: read + with: + eks-cluster-name: sigma-eks-cluster + deploy-environment: production + aws-region: ap-northeast-2 + namespace: analysis + image: ${{ needs.build.outputs.image }} + image-tag: ${{ needs.build.outputs.tag }} + manifest-dir: k8s/app + secrets: + aws-access-key-id: ${{ secrets.AWS_ACCESS_KEY_ID }} + aws-secret-access-key: ${{ secrets.AWS_SECRET_ACCESS_KEY }} diff --git a/Dockerfile b/Dockerfile index f1b81e6..bfdc87a 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,9 +1,9 @@ -FROM eclipse-temurin:17-jdk-alpine AS extractor +FROM eclipse-temurin:17-jdk-jammy AS extractor WORKDIR /builder COPY build/libs/*.jar application.jar RUN java -Djarmode=tools -jar application.jar extract --layers --destination extracted -FROM eclipse-temurin:17-jre-alpine +FROM eclipse-temurin:17-jre-jammy WORKDIR /app COPY --from=extractor /builder/extracted/dependencies/ ./ COPY --from=extractor /builder/extracted/spring-boot-loader/ ./ diff --git a/README.md b/README.md new file mode 100644 index 0000000..b2fee55 --- /dev/null +++ b/README.md @@ -0,0 +1,212 @@ +# SIGMA 스마트 팩토리 월간 요약 분석 서비스 (Analysis-Service) + +이 프로젝트는 **Spring Boot** 기반의 마이크로서비스로, 스마트 팩토리 설비별 **월간 운영 요약 리포트**를 제공합니다. S3에 일자별로 적재된 센서 요약(Parquet)을 집계하고, `management-service`/`anomaly-service`와 통신해 최근 30일간 **불량 건수·이상 건수·센서 평균값**을 하나의 응답으로 조립합니다. + +데이터 흐름은 전부 **읽기 전용(inbound)** 이며, S3·외부 서비스 어디에도 쓰기/발행을 하지 않습니다. 대량의 S3 왕복을 견디기 위해 **병렬 조회 / 인메모리 Parquet 파싱 / 외부 호출 동시화 / 일자별 영구 캐시** 4가지 최적화를 적용했습니다. + +--- + +## 아키텍처 및 시스템 흐름도 + +월간 요약 1건을 조립하기까지의 흐름입니다. 캐시 미적중 일자만 S3에서 병렬로 읽고, 외부 count 2건은 S3 집계와 동시에 진행됩니다. + +```mermaid +sequenceDiagram + autonumber + participant Client as Client + participant A as Analysis Service + participant DB as MariaDB (일자별 캐시) + participant S3 as AWS S3 (Parquet) + participant M as Management Service + participant AN as Anomaly Service + + Client->>A: GET /api/summary/{equipmentName}/monthly + par 외부 count 동시 호출 (Mono.zip) + A->>M: 불량 건수 조회 (defects/count) + A->>AN: 이상 건수 조회 (anomalies/count) + and 센서 집계 (캐시 + 병렬 S3) + A->>DB: 최근 30일 캐시 일괄 조회 (1 query) + alt 캐시 미적중 일자 존재 + A->>S3: 미적중 일자만 병렬 다운로드 (ThreadPool) + S3-->>A: Parquet bytes (인메모리 파싱) + A->>DB: 신규 일자 요약 캐시 저장 + end + end + A-->>Client: { totalDefects, totalAnomalies, sensors[] } +``` + +--- + +## 프로젝트 구조 + +```text +src/main/java/com/factory/analysis_service/ + ├── config/ # 아키텍처 설정 + │ ├── S3Config.java # S3Client 빈 (StaticCredentials / IAM 분기) + │ ├── AwsS3Properties.java # cloud.aws.* 바인딩 + │ ├── WebClientConfig.java # 외부 서비스 호출용 WebClient + │ └── AsyncConfig.java # S3 병렬 조회 전용 ExecutorService + ├── controller/ # REST API (월간 요약 / S3 디버그) + ├── dto/ # MonthlySummaryResponseDTO, SensorSummaryDTO + ├── entity/ # DailySensorSummary (일자별 요약 캐시, 복합키) + ├── repository/ # DailySensorSummaryRepository (JPA) + ├── parquet/ # InMemoryInputFile (디스크 없는 Parquet 리더) + └── service/ # S3SummaryService (집계 핵심 로직) + +src/main/resources/ + ├── application.yml # 기본 설정 (UTC 타임존 고정 포함) + ├── application-dev.yml # dev 프로파일 (MariaDB) + ├── application-local.yml # local 프로파일 + └── application-prod.yml # prod 프로파일 +``` + +데이터 적재 규약 — S3는 `설비/일자` 폴더당 Parquet 1개를 둡니다. + +```text +summary-data/date=2026-05-29/equipmentId=EQP-DEPOSITION-001/part-0.parquet + 스키마: { sensorType: string, unit: string, avg_value: double } +``` + +--- + +## 핵심 기술 및 구현 상세 + +### 1. 30일 S3 조회 병렬화 (`AsyncConfig`, `S3SummaryService`) + +과거 일자의 요약은 서로 독립적이고 불변이므로 완전 병렬화가 가능합니다. `parallelStream`의 공용 ForkJoinPool 대신 **전용 고정 스레드풀(8)** 을 두어, blocking S3 I/O가 애플리케이션의 다른 병렬 작업을 굶기지 않게 합니다. 미적중 일자들을 `CompletableFuture`로 동시 다운로드하고 마지막에 병합합니다. + +> 직렬 30회 왕복 → "가장 느린 한 번" 수준으로 단축. (호출 스레드와 S3 풀이 분리되어 join 데드락도 방지) + +### 2. 디스크 없는 인메모리 Parquet 파싱 (`InMemoryInputFile`) + +기존 구현은 `getObjectAsBytes()`로 이미 메모리에 올라온 바이트를 다시 임시파일에 쓰고(`Files.write`) 디스크에서 재차 읽은 뒤 삭제했습니다. Parquet은 footer seek 때문에 random-access가 필요할 뿐 디스크가 필요한 것은 아니므로, **byte[] 위에서 seek/read를 직접 제공**하는 `InputFile`을 구현해 디스크 I/O를 0으로 만들었습니다. (read-only 컨테이너 환경에도 안전) + +### 3. 외부 count 동시 호출 (`fetchCountsAsync`) + +`management`/`anomaly` 두 count 호출은 서로 독립적이므로 `Mono.zip`으로 **동시에 구독**해 병렬 실행합니다. 나아가 이 두 호출을 S3 집계 작업과도 겹쳐 전체 레이턴시를 "가장 오래 걸리는 한 워크스트림"으로 수렴시킵니다. 호출 실패 시 해당 count는 `0`으로 폴백되어 전체 응답은 죽지 않습니다. + +### 4. 일자별 영구 캐시 (`DailySensorSummary`) + +월간 요약은 항상 "1~30일 전"을 보므로 매 호출마다 29/30이 겹칩니다. 과거 일자의 S3 집계값은 불변이라 `(equipmentId, date, sensorType)` 단위로 영구 캐시하면 반복 호출 시 캐시 적중률이 매우 높습니다. 30일치 캐시는 **단 한 번의 범위 쿼리**로 일괄 로드하며, 미적중 일자만 S3로 갑니다. (가변값인 defect/anomaly count는 캐시하지 않고 매번 실시간 조회) + +> 복합키를 직접 할당하므로 `Persistable`을 구현해 `save()`가 `merge`(UPDATE)가 아닌 `persist`(INSERT)로 동작하도록 했습니다. + +--- + +## REST API 명세 + +### 1. 설비별 월간 요약 조회 + +#### Request + +```http +GET /api/summary/EQP-DEPOSITION-001/monthly +``` + +#### Response + +```json +{ + "success": true, + "status": 200, + "message": "success", + "data": { + "totalDefects": 11, + "totalAnomalies": 5, + "sensors": [ + { "sensorType": "TEMP", "unit": "C", "avgValue": 20.5 }, + { "sensorType": "PRESSURE", "unit": "kPa", "avgValue": 100.0 } + ] + }, + "timestamp": "2026-06-12T04:24:15.918" +} +``` + +### 2. S3 연결/스키마 진단 (디버그용) + +#### Request + +```http +GET /api/summary/debug/EQP-DEPOSITION-001 +``` + +#### Response + +```json +{ + "bucket": "sigma-factory-bucket", + "region": "ap-northeast-2", + "prefix": "summary-data/date=2026-05-28/equipmentId=EQP-DEPOSITION-001/", + "foundFiles": "summary-data/date=2026-05-28/.../part-0.parquet", + "fileSize": "1234 bytes", + "schema": "message summary { ... }", + "rowCount": "2", + "firstRow": "sensorType: TEMP, unit: C, avg_value: 20.0" +} +``` + +--- + +## 검증 및 테스트 + +총 **25개 테스트**가 단위·통합·성능을 커버하며, JaCoCo 기준 **Instruction 93% / Branch(조건) 89%** 를 달성합니다. + +| 검증 항목 | 테스트 | +| :--- | :--- | +| S3 Parquet 데이터 조회·30일 평균 집계 | `S3SummaryServiceTest` | +| 외부 서비스 count 통신 및 실패 폴백 | `S3SummaryServiceTest` (stub WebClient) | +| 병렬 수행 안정성 (16스레드 동시 호출) | `S3SummaryServiceTest#concurrentCallsAreConsistent` | +| MariaDB(H2) 캐시 저장·재사용 | `S3SummaryServiceIntegrationTest` | +| 인메모리 Parquet 읽기 (EOF/seek/ByteBuffer) | `InMemoryInputFileTest` | +| 병렬 vs 직렬 성능 (지연 주입) | `S3SummaryServicePerformanceTest` | +| S3 자격증명 분기 (Static / IAM) | `S3ConfigTest` | + +### 테스트 명령어 + +```bash +./gradlew test # 전체 테스트 + JaCoCo 리포트 +# 리포트: build/reports/jacoco/test/html/index.html +``` + +> 테스트 JVM은 `user.timezone=UTC`로 기동합니다. 앱이 런타임에 JVM 기본 타임존을 UTC로 바꾸므로, 테스트도 UTC로 맞춰 `LocalDate`↔JDBC 날짜 변환의 타임존 오프바이원을 방지합니다. + +--- + +## 환경 변수 및 설정 + +보안을 위해 실제 자격 증명은 README에 기술하지 않으며, 루트의 `.env.example`을 참고해 `.env`를 생성합니다. + +| 환경 변수명 | 설명 | 예시 / 기본값 | +| :--- | :--- | :--- | +| `S3_BUCKET_NAME` | 요약 Parquet이 저장된 S3 버킷명 | `sigma-factory-bucket` | +| `AWS_ACCESS_KEY_ID` | AWS 액세스 키 (비우면 IAM 역할 사용) | *(빈 값)* | +| `AWS_SECRET_ACCESS_KEY` | AWS 시크릿 키 (비우면 IAM 역할 사용) | *(빈 값)* | +| `MANAGEMENT_SERVICE_URL` | 불량 건수 조회 대상 서비스 | `http://localhost:8086` | +| `ANOMALY_SERVICE_URL` | 이상 건수 조회 대상 서비스 | `http://localhost:8085` | +| `DB_HOST` / `DB_PORT` | MariaDB 호스트 / 포트 | `localhost` / `3306` | +| `DB_NAME` | 캐시 데이터베이스명 | `analysis_db` | +| `DB_USERNAME` / `DB_PASSWORD` | DB 접속 계정 / 패스워드 | `root` / `your_password` | +| `DB_MAX_POOL_SIZE` / `DB_MIN_IDLE` | HikariCP 풀 크기 / 최소 유휴 | `5` / `2` | +| `KAFKA_BOOTSTRAP_SERVERS` | Kafka 브로커 주소 | `localhost:9092` | + +* AWS 리전은 `ap-northeast-2`로 고정되어 있습니다. +* 운영 컨테이너 타임존과 무관하게 캐시 날짜가 어긋나지 않도록 `hibernate.jdbc.time_zone=UTC`가 기본 설정되어 있습니다. + +--- + +## 로컬 실행 방법 + +### 1. 의존성 및 빌드 검증 + +```bash +./gradlew clean build -x test +``` + +### 2. 로컬 실행 + +```bash +./gradlew bootRun +``` + +* 기본 포트: **8087** +* 활성 프로파일: `dev` (변경: `--spring.profiles.active=local`) diff --git a/build.gradle b/build.gradle index bffe0bc..cca51e9 100644 --- a/build.gradle +++ b/build.gradle @@ -1,6 +1,7 @@ plugins { id "com.factory.spring-application-conventions" id "com.factory.maven-consumer-conventions" + id "jacoco" } version = '1.0.0' @@ -11,9 +12,10 @@ repositories { dependencies { // common modules - implementation "com.factory.common:contract:1.0.1" - implementation "com.factory.common:outbox-jpa:1.0.1" - implementation "com.factory.common:inbox-jpa:1.0.1" + implementation "com.factory.common:contract:1.1.0" + implementation "com.factory.common:kafka:1.1.0" + implementation "com.factory.common:outbox-jpa:1.1.0" + implementation "com.factory.common:inbox-jpa:1.1.0" // Spring Web implementation "org.springframework.boot:spring-boot-starter-web" @@ -34,6 +36,8 @@ dependencies { // AWS SDK v2 - S3 implementation platform('software.amazon.awssdk:bom:2.25.60') implementation 'software.amazon.awssdk:s3' + // IRSA(WebIdentityToken)로 자격증명 받으려면 STS 모듈이 클래스패스에 필요 + implementation 'software.amazon.awssdk:sts' // Parquet reading (LocalInputFile - Hadoop FS 미사용) implementation 'org.apache.parquet:parquet-hadoop:1.14.1' @@ -69,10 +73,33 @@ dependencies { // Test testImplementation 'org.springframework.boot:spring-boot-starter-test' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' + testImplementation 'com.h2database:h2' } tasks.named('test', Test) { useJUnitPlatform() + // 앱이 런타임에 JVM 기본 TZ를 UTC로 바꾸므로, 테스트 JVM도 처음부터 UTC로 띄워 + // LocalDate <-> JDBC(H2) 날짜 변환의 타임존 오프바이원을 방지한다. + systemProperty 'user.timezone', 'UTC' + finalizedBy jacocoTestReport +} + +jacocoTestReport { + dependsOn test + reports { + xml.required = true + html.required = true + } +} + +jacocoTestCoverageVerification { + violationRules { + rule { + limit { + minimum = 0.80 + } + } + } } tasks.named('bootJar') { diff --git a/k8s/app/analysis-service-deployment.yml b/k8s/app/analysis-service-deployment.yml new file mode 100644 index 0000000..7e51813 --- /dev/null +++ b/k8s/app/analysis-service-deployment.yml @@ -0,0 +1,53 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: analysis-service + namespace: analysis + labels: + app: analysis-service +spec: + replicas: 1 + selector: + matchLabels: + app: analysis-service + template: + metadata: + labels: + app: analysis-service + spec: + serviceAccountName: analysis-sa + containers: + - name: analysis-service + image: $IMAGE:$IMAGE_TAG + ports: + - containerPort: 8080 + envFrom: + - secretRef: + name: s3-secret + - secretRef: + name: rdb-secret + env: + - name: DB_PORT + value: "3306" + - name: DB_NAME + value: "analysis_db" + - name: KAFKA_BOOTSTRAP_SERVERS + value: "sigma-kafka-kafka-bootstrap.default.svc.cluster.local:9092" + - name: DB_MAX_POOL_SIZE + value: "5" + - name: DB_MIN_IDLE + value: "2" + - name: SERVER_PORT + value: "8080" + - name: MANAGEMENT_SERVICE_URL + value: "http://management-service.management.svc.cluster.local:8086" + - name: ANOMALY_SERVICE_URL + value: "http://anomaly-service.anomaly.svc.cluster.local:8083" + resources: + limits: + memory: "512Mi" + cpu: "500m" + requests: + memory: "256Mi" + cpu: "250m" + restartPolicy: Always diff --git a/k8s/app/analysis-service-ingress.yml b/k8s/app/analysis-service-ingress.yml new file mode 100644 index 0000000..5726e1d --- /dev/null +++ b/k8s/app/analysis-service-ingress.yml @@ -0,0 +1,35 @@ +apiVersion: v1 +kind: Service +metadata: + name: analysis-service + namespace: analysis +spec: + selector: + app: analysis-service + ports: + - port: 8087 + targetPort: 8080 + type: ClusterIP +--- +apiVersion: networking.k8s.io/v1 +kind: Ingress +metadata: + name: analysis-ingress + namespace: analysis + annotations: + kubernetes.io/ingress.class: alb + alb.ingress.kubernetes.io/scheme: internet-facing + alb.ingress.kubernetes.io/target-type: ip + alb.ingress.kubernetes.io/group.name: sigma + alb.ingress.kubernetes.io/group.order: "6" +spec: + rules: + - http: + paths: + - path: /api/summary + pathType: Prefix + backend: + service: + name: analysis-service + port: + number: 8087 diff --git a/src/main/java/com/factory/analysis_service/config/AsyncConfig.java b/src/main/java/com/factory/analysis_service/config/AsyncConfig.java new file mode 100644 index 0000000..f1a483e --- /dev/null +++ b/src/main/java/com/factory/analysis_service/config/AsyncConfig.java @@ -0,0 +1,35 @@ +package com.factory.analysis_service.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * 30일치 S3 조회를 병렬화하기 위한 전용 스레드풀. + * + *

{@code parallelStream}의 공용 ForkJoinPool을 쓰지 않는 이유: S3 호출은 blocking I/O라 + * 공용 풀을 점유하면 애플리케이션의 다른 병렬 작업을 굶길 수 있다. I/O 바운드라 CPU 코어 수보다 + * 크게 잡아 동시 다운로드 처리량을 높인다. + */ +@Configuration +public class AsyncConfig { + + @Bean(name = "s3FetchExecutor", destroyMethod = "shutdown") + public ExecutorService s3FetchExecutor() { + ThreadFactory factory = new ThreadFactory() { + private final AtomicInteger seq = new AtomicInteger(); + + @Override + public Thread newThread(Runnable r) { + Thread t = new Thread(r, "s3-fetch-" + seq.incrementAndGet()); + t.setDaemon(true); + return t; + } + }; + return Executors.newFixedThreadPool(8, factory); + } +} diff --git a/src/main/java/com/factory/analysis_service/config/JpaConfig.java b/src/main/java/com/factory/analysis_service/config/JpaConfig.java new file mode 100644 index 0000000..3c86803 --- /dev/null +++ b/src/main/java/com/factory/analysis_service/config/JpaConfig.java @@ -0,0 +1,11 @@ +package com.factory.analysis_service.config; + +import org.springframework.boot.autoconfigure.domain.EntityScan; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.jpa.repository.config.EnableJpaRepositories; + +@Configuration +@EnableJpaRepositories(basePackages = "com.factory.analysis_service") +@EntityScan(basePackages = "com.factory.analysis_service") +public class JpaConfig { +} diff --git a/src/main/java/com/factory/analysis_service/entity/DailySensorSummary.java b/src/main/java/com/factory/analysis_service/entity/DailySensorSummary.java new file mode 100644 index 0000000..b910a54 --- /dev/null +++ b/src/main/java/com/factory/analysis_service/entity/DailySensorSummary.java @@ -0,0 +1,91 @@ +package com.factory.analysis_service.entity; + +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.Id; +import jakarta.persistence.IdClass; +import jakarta.persistence.PostLoad; +import jakarta.persistence.PostPersist; +import jakarta.persistence.Table; +import jakarta.persistence.Transient; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; +import org.springframework.data.domain.Persistable; + +import java.io.Serializable; +import java.time.LocalDate; + +/** + * 설비/일자/센서타입별 일일 평균 요약값 캐시. + * + *

과거 일자의 S3 Parquet 집계 결과는 불변이므로 한 번 읽으면 영구 캐시할 수 있다. + * 월간 요약은 항상 "1~30일 전"을 보기 때문에, 매 호출마다 29/30이 겹쳐 캐시 적중률이 높다. + * defect/anomaly count는 가변이라 이 캐시에 넣지 않고 매번 실시간 조회한다. + * + *

복합키({@link IdClass})를 직접 할당하므로 Spring Data의 기본 isNew 판정이 항상 false가 되어 + * {@code save()}가 {@code merge}(UPDATE)로 동작한다. 캐시는 항상 신규 INSERT이므로 + * {@link Persistable}을 구현해 {@code persist}가 되도록 한다. + */ +@Entity +@Table(name = "daily_sensor_summary") +@IdClass(DailySensorSummary.Pk.class) +@Getter +@Setter +@NoArgsConstructor +@AllArgsConstructor +@Builder +public class DailySensorSummary implements Persistable { + + @Id + @Column(name = "equipment_id", nullable = false, length = 100) + private String equipmentId; + + @Id + @Column(name = "summary_date", nullable = false) + private LocalDate summaryDate; + + @Id + @Column(name = "sensor_type", nullable = false, length = 100) + private String sensorType; + + @Column(name = "unit", length = 50) + private String unit; + + @Column(name = "avg_value", nullable = false) + private double avgValue; + + // 자연 기본값 false → 빌더/생성자로 만든 신규 엔티티는 isNew()=true (persist). + // Hibernate가 조회로 적재하면 @PostLoad로 true가 되어 isNew()=false. + @Transient + private boolean loaded; + + @Override + public Pk getId() { + return new Pk(equipmentId, summaryDate, sensorType); + } + + @Override + public boolean isNew() { + return !loaded; + } + + @PostLoad + void markLoaded() { + this.loaded = true; + } + + @Getter + @Setter + @NoArgsConstructor + @AllArgsConstructor + @EqualsAndHashCode + public static class Pk implements Serializable { + private String equipmentId; + private LocalDate summaryDate; + private String sensorType; + } +} diff --git a/src/main/java/com/factory/analysis_service/parquet/InMemoryInputFile.java b/src/main/java/com/factory/analysis_service/parquet/InMemoryInputFile.java new file mode 100644 index 0000000..f6f8e1f --- /dev/null +++ b/src/main/java/com/factory/analysis_service/parquet/InMemoryInputFile.java @@ -0,0 +1,114 @@ +package com.factory.analysis_service.parquet; + +import org.apache.parquet.io.InputFile; +import org.apache.parquet.io.SeekableInputStream; + +import java.io.EOFException; +import java.io.IOException; +import java.nio.ByteBuffer; + +/** + * S3에서 받은 byte[]를 디스크 임시파일 없이 그대로 Parquet 리더에 넘기기 위한 {@link InputFile} 구현. + * + *

기존 구현은 {@code getObjectAsBytes()}로 이미 메모리에 올라온 바이트를 다시 임시파일에 쓰고 + * ({@code Files.write}) {@code LocalInputFile}로 디스크에서 재차 읽은 뒤 삭제했다. + * Parquet은 footer seek 때문에 random-access가 필요할 뿐, 디스크가 필요한 것은 아니므로 + * 메모리 바이트 배열 위에서 seek/read를 직접 제공하여 디스크 I/O를 0으로 만든다. + */ +public class InMemoryInputFile implements InputFile { + + private final byte[] data; + + public InMemoryInputFile(byte[] data) { + this.data = data; + } + + @Override + public long getLength() { + return data.length; + } + + @Override + public SeekableInputStream newStream() { + return new InMemorySeekableInputStream(data); + } + + private static final class InMemorySeekableInputStream extends SeekableInputStream { + + private final byte[] data; + private int pos = 0; + + private InMemorySeekableInputStream(byte[] data) { + this.data = data; + } + + @Override + public long getPos() { + return pos; + } + + @Override + public void seek(long newPos) throws IOException { + if (newPos < 0 || newPos > data.length) { + throw new IOException("Invalid seek position: " + newPos + " (length=" + data.length + ")"); + } + pos = (int) newPos; + } + + @Override + public int read() { + return pos < data.length ? (data[pos++] & 0xFF) : -1; + } + + @Override + public int read(byte[] b, int off, int len) { + if (pos >= data.length) { + return -1; + } + int n = Math.min(len, data.length - pos); + System.arraycopy(data, pos, b, off, n); + pos += n; + return n; + } + + @Override + public void readFully(byte[] bytes) throws IOException { + readFully(bytes, 0, bytes.length); + } + + @Override + public void readFully(byte[] bytes, int start, int len) throws IOException { + if (pos + len > data.length) { + throw new EOFException("Reached end of stream with " + len + " bytes left to read"); + } + System.arraycopy(data, pos, bytes, start, len); + pos += len; + } + + @Override + public int read(ByteBuffer buf) { + if (pos >= data.length) { + return -1; + } + int n = Math.min(buf.remaining(), data.length - pos); + buf.put(data, pos, n); + pos += n; + return n; + } + + @Override + public void readFully(ByteBuffer buf) throws IOException { + int n = buf.remaining(); + if (pos + n > data.length) { + throw new EOFException("Reached end of stream with " + n + " bytes left to read"); + } + buf.put(data, pos, n); + pos += n; + } + + @Override + public int available() { + return data.length - pos; + } + } +} diff --git a/src/main/java/com/factory/analysis_service/repository/DailySensorSummaryRepository.java b/src/main/java/com/factory/analysis_service/repository/DailySensorSummaryRepository.java new file mode 100644 index 0000000..40f7495 --- /dev/null +++ b/src/main/java/com/factory/analysis_service/repository/DailySensorSummaryRepository.java @@ -0,0 +1,18 @@ +package com.factory.analysis_service.repository; + +import com.factory.analysis_service.entity.DailySensorSummary; +import org.springframework.data.jpa.repository.JpaRepository; + +import java.time.LocalDate; +import java.util.List; + +public interface DailySensorSummaryRepository + extends JpaRepository { + + /** + * 한 번의 쿼리로 설비의 기간 내 캐시된 모든 일자/센서 요약을 가져온다. + * (30일치를 일자별로 30번 조회하지 않기 위함) + */ + List findByEquipmentIdAndSummaryDateBetween( + String equipmentId, LocalDate from, LocalDate to); +} diff --git a/src/main/java/com/factory/analysis_service/service/S3SummaryService.java b/src/main/java/com/factory/analysis_service/service/S3SummaryService.java index 5511ef4..b9619b4 100644 --- a/src/main/java/com/factory/analysis_service/service/S3SummaryService.java +++ b/src/main/java/com/factory/analysis_service/service/S3SummaryService.java @@ -1,9 +1,12 @@ package com.factory.analysis_service.service; import com.factory.analysis_service.config.AwsS3Properties; -import com.factory.common.core.dto.ApiResponse; +import com.factory.analysis_service.entity.DailySensorSummary; +import com.factory.analysis_service.parquet.InMemoryInputFile; +import com.factory.analysis_service.repository.DailySensorSummaryRepository; import com.factory.analysis_service.dto.MonthlySummaryResponseDTO; import com.factory.analysis_service.dto.SensorSummaryDTO; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.parquet.column.page.PageReadStore; @@ -11,40 +14,49 @@ import org.apache.parquet.example.data.simple.convert.GroupRecordConverter; import org.apache.parquet.hadoop.ParquetFileReader; import org.apache.parquet.io.ColumnIOFactory; -import org.apache.parquet.io.LocalInputFile; +import org.apache.parquet.io.InputFile; import org.apache.parquet.io.MessageColumnIO; import org.apache.parquet.io.RecordReader; import org.apache.parquet.schema.MessageType; import org.springframework.beans.factory.annotation.Value; -import org.springframework.core.ParameterizedTypeReference; +import org.springframework.dao.DataAccessException; import org.springframework.stereotype.Service; import org.springframework.web.reactive.function.client.WebClient; +import reactor.core.publisher.Mono; import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.GetObjectRequest; import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; import software.amazon.awssdk.services.s3.model.S3Object; import java.io.IOException; -import java.nio.file.Files; -import java.nio.file.Path; import java.time.LocalDate; -import java.time.LocalDateTime; import java.time.format.DateTimeFormatter; -import java.util.*; +import java.util.ArrayList; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; import java.util.stream.Collectors; @Slf4j @Service @RequiredArgsConstructor public class S3SummaryService { - + static { System.setProperty("hadoop.home.dir", "/"); } + private static final int LOOKBACK_DAYS = 30; + private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("yyyy-MM-dd"); + private final S3Client s3Client; private final AwsS3Properties awsS3Properties; private final WebClient webClient; + private final DailySensorSummaryRepository cacheRepository; + private final ExecutorService s3FetchExecutor; @Value("${services.management-service.url}") private String managementServiceUrl; @@ -52,122 +64,171 @@ public class S3SummaryService { @Value("${services.anomaly-service.url}") private String anomalyServiceUrl; - private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("yyyy-MM-dd"); - public MonthlySummaryResponseDTO getMonthlySummary(String equipmentName) { - Map> avgByType = new LinkedHashMap<>(); - Map unitByType = new LinkedHashMap<>(); - LocalDate today = LocalDate.now(); - for (int i = 1; i <= 30; i++) { - String dateStr = today.minusDays(i).format(DATE_FMT); - String prefix = String.format("summary-data/date=%s/equipmentId=%s/", dateStr, equipmentName); + String sinceStr = today.minusDays(LOOKBACK_DAYS).atStartOfDay().toString(); - for (ParquetRow row : fetchRowsByPrefix(prefix)) { - avgByType.computeIfAbsent(row.sensorType(), k -> new ArrayList<>()).add(row.avgValue()); - unitByType.put(row.sensorType(), row.unit()); - } + // (3) defect/anomaly count 두 호출을 S3 집계와 동시에 시작한다. (논블로킹으로 먼저 띄움) + CompletableFuture countsFuture = fetchCountsAsync(equipmentName, sinceStr); + + // (1)(2)(4) 캐시 조회 → 미적중 일자만 병렬 S3 조회 → 인메모리 Parquet 파싱 → 캐시 저장 + List sensors = aggregateSensors(equipmentName, today); + + long[] counts = countsFuture.join(); + + return MonthlySummaryResponseDTO.builder() + .totalDefects(counts[0]) + .totalAnomalies(counts[1]) + .sensors(sensors) + .build(); + } + + // ---------------------------------------------------------------------------------- + // 센서 집계: 캐시 + 병렬 S3 + // ---------------------------------------------------------------------------------- + + private List aggregateSensors(String equipmentName, LocalDate today) { + List dates = new ArrayList<>(LOOKBACK_DAYS); + for (int i = 1; i <= LOOKBACK_DAYS; i++) { + dates.add(today.minusDays(i)); + } + LocalDate from = today.minusDays(LOOKBACK_DAYS); + LocalDate to = today.minusDays(1); + + // (4) 캐시를 단 한 번의 쿼리로 일괄 로딩 + Map> rowsByDate = new LinkedHashMap<>(); + for (DailySensorSummary cached : cacheRepository + .findByEquipmentIdAndSummaryDateBetween(equipmentName, from, to)) { + rowsByDate.computeIfAbsent(cached.getSummaryDate(), k -> new ArrayList<>()) + .add(new ParquetRow(cached.getSensorType(), cached.getUnit(), cached.getAvgValue())); } - LocalDateTime from = today.minusDays(30).atStartOfDay(); - String sinceStr = from.toString(); + // 캐시에 없는 일자만 추려서 + List missing = dates.stream() + .filter(d -> !rowsByDate.containsKey(d)) + .collect(Collectors.toList()); - long totalDefects = 0; - try { - ApiResponse defectRes = webClient.get() - .uri(managementServiceUrl + "/api/management/defects/count?equipmentName=" + equipmentName + "&since=" + sinceStr) - .retrieve() - .bodyToMono(new ParameterizedTypeReference>() {}) - .block(); - if (defectRes != null && defectRes.getData() != null) { - totalDefects = defectRes.getData(); - } - } catch (Exception e) { - log.error("### management-service defect count API 호출 실패: {}", e.getMessage()); + // (1) 미적중 일자들을 병렬로 S3 조회 + if (!missing.isEmpty()) { + Map> fetched = fetchDatesInParallel(equipmentName, missing); + rowsByDate.putAll(fetched); + persistCache(equipmentName, fetched); } - long totalAnomalies = 0; - try { - ApiResponse anomalyRes = webClient.get() - .uri(anomalyServiceUrl + "/api/anomalies/count?equipmentName=" + equipmentName + "&since=" + sinceStr) - .retrieve() - .bodyToMono(new ParameterizedTypeReference>() {}) - .block(); - if (anomalyRes != null && anomalyRes.getData() != null) { - totalAnomalies = anomalyRes.getData(); + // 일자 순서대로 센서타입별 평균값 누적 + Map> avgByType = new LinkedHashMap<>(); + Map unitByType = new LinkedHashMap<>(); + for (LocalDate date : dates) { + List rows = rowsByDate.get(date); + if (rows == null) { + continue; + } + for (ParquetRow row : rows) { + avgByType.computeIfAbsent(row.sensorType(), k -> new ArrayList<>()).add(row.avgValue()); + unitByType.put(row.sensorType(), row.unit()); } - } catch (Exception e) { - log.error("### anomaly-service anomaly count API 호출 실패: {}", e.getMessage()); } - List sensors = avgByType.entrySet().stream() + return avgByType.entrySet().stream() .map(e -> SensorSummaryDTO.builder() .sensorType(e.getKey()) .unit(unitByType.get(e.getKey())) - .avgValue(roundTwo(e.getValue().stream().mapToDouble(Double::doubleValue).average().orElse(0.0))) + .avgValue(roundTwo(e.getValue().stream() + .mapToDouble(Double::doubleValue).average().orElse(0.0))) .build()) .collect(Collectors.toList()); - - return MonthlySummaryResponseDTO.builder() - .totalDefects(totalDefects) - .totalAnomalies(totalAnomalies) - .sensors(sensors) - .build(); } - public Map debugS3(String equipmentName) { - Map result = new LinkedHashMap<>(); - String bucket = awsS3Properties.getS3().getBucket(); - String prefix = String.format("summary-data/date=2026-05-28/equipmentId=%s/", equipmentName); - - result.put("bucket", bucket); - result.put("region", awsS3Properties.getRegion()); - result.put("prefix", prefix); + private Map> fetchDatesInParallel(String equipmentName, List dates) { + List>>> futures = dates.stream() + .map(date -> CompletableFuture.supplyAsync( + () -> Map.entry(date, fetchRowsByPrefix(buildPrefix(date, equipmentName))), + s3FetchExecutor)) + .collect(Collectors.toList()); - try { - List keys = listParquetKeys(prefix); - if (keys.isEmpty()) { - result.put("foundFiles", "없음"); - return result; + Map> result = new LinkedHashMap<>(); + for (CompletableFuture>> future : futures) { + Map.Entry> entry = future.join(); + // 데이터가 없는 일자는 캐싱/집계에서 제외 + if (!entry.getValue().isEmpty()) { + result.put(entry.getKey(), entry.getValue()); } - result.put("foundFiles", keys.get(0)); - - // 파일 직접 읽어서 스키마와 첫 번째 행 확인 - byte[] bytes = s3Client.getObjectAsBytes( - GetObjectRequest.builder().bucket(bucket).key(keys.get(0)).build() - ).asByteArray(); - result.put("fileSize", bytes.length + " bytes"); + } + return result; + } - Path tempFile = Files.createTempFile("parquet-debug-", ".parquet"); - try { - Files.write(tempFile, bytes); - try (ParquetFileReader fileReader = ParquetFileReader.open(new LocalInputFile(tempFile))) { - MessageType schema = fileReader.getFooter().getFileMetaData().getSchema(); - result.put("schema", schema.toString()); - - PageReadStore pages = fileReader.readNextRowGroup(); - if (pages != null) { - result.put("rowCount", String.valueOf(pages.getRowCount())); - MessageColumnIO columnIO = new ColumnIOFactory().getColumnIO(schema); - RecordReader rr = columnIO.getRecordReader(pages, new GroupRecordConverter(schema)); - Group g = rr.read(); - if (g != null) result.put("firstRow", g.toString()); - } - } - } finally { - Files.deleteIfExists(tempFile); + private void persistCache(String equipmentName, Map> fetched) { + List toSave = new ArrayList<>(); + fetched.forEach((date, rows) -> { + for (ParquetRow row : rows) { + toSave.add(DailySensorSummary.builder() + .equipmentId(equipmentName) + .summaryDate(date) + .sensorType(row.sensorType()) + .unit(row.unit()) + .avgValue(row.avgValue()) + .build()); } - } catch (Throwable e) { - result.put("error", e.getClass().getSimpleName() + ": " + e.getMessage()); + }); + if (toSave.isEmpty()) { + return; } - return result; + try { + cacheRepository.saveAll(toSave); + } catch (DataAccessException e) { + // 동시 요청이 같은 일자를 동시에 적재하려다 충돌한 경우 — 캐시일 뿐이므로 무시 + log.warn("### 일일 요약 캐시 저장 충돌(동시 요청 추정), 무시: {}", e.getMessage()); + } + } + + private String buildPrefix(LocalDate date, String equipmentName) { + return String.format("summary-data/date=%s/equipmentId=%s/", date.format(DATE_FMT), equipmentName); + } + + // ---------------------------------------------------------------------------------- + // 외부 서비스 count 동시 호출 + // ---------------------------------------------------------------------------------- + + private CompletableFuture fetchCountsAsync(String equipmentName, String sinceStr) { + Mono defects = countMono( + managementServiceUrl + "/api/management/defects/count?equipmentName=" + equipmentName + "&since=" + sinceStr, + "management-service"); + Mono anomalies = countMono( + anomalyServiceUrl + "/api/anomalies/count?equipmentName=" + equipmentName + "&since=" + sinceStr, + "anomaly-service"); + + // zip 은 두 Mono 를 동시에 구독 → 두 호출이 병렬 실행된다. + return Mono.zip(defects, anomalies, (d, a) -> new long[]{d, a}).toFuture(); } + private Mono countMono(String uri, String serviceName) { + return webClient.get() + .uri(uri) + .retrieve() + // ApiResponse 전체를 역직렬화하지 않고 필요한 data 필드만 추출 (역직렬화 견고성↑) + .bodyToMono(CountResponse.class) + .map(res -> res.data() != null ? res.data() : 0L) + .onErrorResume(e -> { + log.error("### {} count API 호출 실패: {}", serviceName, e.getMessage()); + return Mono.just(0L); + }); + } + + /** management/anomaly count 응답에서 data 값만 받기 위한 최소 DTO. (success/status 등은 무시) */ + @JsonIgnoreProperties(ignoreUnknown = true) + private record CountResponse(Long data) {} + + // ---------------------------------------------------------------------------------- + // S3 + Parquet (인메모리) + // ---------------------------------------------------------------------------------- + private List fetchRowsByPrefix(String prefix) { try { List keys = listParquetKeys(prefix); - if (keys.isEmpty()) return Collections.emptyList(); - + if (keys.isEmpty()) { + return Collections.emptyList(); + } + // 폴더(date+equipment)당 parquet 파일은 1개라는 데이터 적재 규약에 따라 첫 파일을 읽는다. return fetchRows(keys.get(0)); } catch (Throwable e) { log.error("### prefix 조회 실패: {}", prefix, e); @@ -187,7 +248,6 @@ private List listParquetKeys(String prefix) { } private List fetchRows(String s3Key) { - Path tempFile = null; try { byte[] bytes = s3Client.getObjectAsBytes( GetObjectRequest.builder() @@ -196,23 +256,17 @@ private List fetchRows(String s3Key) { .build() ).asByteArray(); - tempFile = Files.createTempFile("parquet-", ".parquet"); - Files.write(tempFile, bytes); - return readParquet(tempFile); - + // (2) 임시파일 없이 메모리 바이트에서 바로 Parquet 파싱 + return readParquet(new InMemoryInputFile(bytes)); } catch (Throwable e) { log.error("### S3 Parquet 읽기 실패: {}", s3Key, e); return Collections.emptyList(); - } finally { - if (tempFile != null) { - try { Files.deleteIfExists(tempFile); } catch (IOException ignored) {} - } } } - private List readParquet(Path path) throws IOException { + private List readParquet(InputFile inputFile) throws IOException { List rows = new ArrayList<>(); - try (ParquetFileReader fileReader = ParquetFileReader.open(new LocalInputFile(path))) { + try (ParquetFileReader fileReader = ParquetFileReader.open(inputFile)) { MessageType schema = fileReader.getFooter().getFileMetaData().getSchema(); MessageColumnIO columnIO = new ColumnIOFactory().getColumnIO(schema); PageReadStore pages; @@ -239,5 +293,52 @@ private double roundTwo(double value) { return Math.round(value * 100.0) / 100.0; } + // ---------------------------------------------------------------------------------- + // 디버그 엔드포인트 (인메모리 Parquet 읽기로 통일) + // ---------------------------------------------------------------------------------- + + public Map debugS3(String equipmentName) { + Map result = new LinkedHashMap<>(); + String bucket = awsS3Properties.getS3().getBucket(); + String prefix = String.format("summary-data/date=2026-05-28/equipmentId=%s/", equipmentName); + + result.put("bucket", bucket); + result.put("region", awsS3Properties.getRegion()); + result.put("prefix", prefix); + + try { + List keys = listParquetKeys(prefix); + if (keys.isEmpty()) { + result.put("foundFiles", "없음"); + return result; + } + result.put("foundFiles", keys.get(0)); + + byte[] bytes = s3Client.getObjectAsBytes( + GetObjectRequest.builder().bucket(bucket).key(keys.get(0)).build() + ).asByteArray(); + result.put("fileSize", bytes.length + " bytes"); + + try (ParquetFileReader fileReader = ParquetFileReader.open(new InMemoryInputFile(bytes))) { + MessageType schema = fileReader.getFooter().getFileMetaData().getSchema(); + result.put("schema", schema.toString()); + + PageReadStore pages = fileReader.readNextRowGroup(); + if (pages != null) { + result.put("rowCount", String.valueOf(pages.getRowCount())); + MessageColumnIO columnIO = new ColumnIOFactory().getColumnIO(schema); + RecordReader rr = columnIO.getRecordReader(pages, new GroupRecordConverter(schema)); + Group g = rr.read(); + if (g != null) { + result.put("firstRow", g.toString()); + } + } + } + } catch (Throwable e) { + result.put("error", e.getClass().getSimpleName() + ": " + e.getMessage()); + } + return result; + } + private record ParquetRow(String sensorType, String unit, double avgValue) {} } diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index cd29ee0..ca28b28 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -17,6 +17,13 @@ spring: profiles: active: dev + jpa: + properties: + hibernate: + # 캐시(LocalDate) 저장/조회 시 JDBC 날짜 타임존 오프바이원 방지 (컨테이너 TZ 무관하게 UTC 고정) + jdbc: + time_zone: UTC + kafka: bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092} producer: @@ -40,4 +47,4 @@ services: management-service: url: ${MANAGEMENT_SERVICE_URL:http://localhost:8086} anomaly-service: - url: ${ANOMALY_SERVICE_URL:http://localhost:8085} \ No newline at end of file + url: ${ANOMALY_SERVICE_URL:http://localhost:8083} \ No newline at end of file diff --git a/src/test/java/com/factory/analysis_service/config/S3ConfigTest.java b/src/test/java/com/factory/analysis_service/config/S3ConfigTest.java new file mode 100644 index 0000000..fcbba4e --- /dev/null +++ b/src/test/java/com/factory/analysis_service/config/S3ConfigTest.java @@ -0,0 +1,36 @@ +package com.factory.analysis_service.config; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import software.amazon.awssdk.services.s3.S3Client; + +import static org.assertj.core.api.Assertions.assertThat; + +@DisplayName("S3Config - 자격증명 분기") +class S3ConfigTest { + + private AwsS3Properties props(String accessKey, String secretKey) { + AwsS3Properties p = new AwsS3Properties(); + p.setRegion("ap-northeast-2"); + p.getS3().setBucket("test-bucket"); + p.getCredentials().setAccessKey(accessKey); + p.getCredentials().setSecretKey(secretKey); + return p; + } + + @Test + @DisplayName("access/secret 키가 있으면 StaticCredentialsProvider로 S3Client를 만든다") + void staticCredentials() { + S3Client client = new S3Config(props("AKIA-TEST", "secret")).s3Client(); + assertThat(client).isNotNull(); + client.close(); + } + + @Test + @DisplayName("키가 비어있으면 DefaultCredentialsProvider(IAM)로 S3Client를 만든다") + void defaultCredentials() { + S3Client client = new S3Config(props("", "")).s3Client(); + assertThat(client).isNotNull(); + client.close(); + } +} diff --git a/src/test/java/com/factory/analysis_service/controller/SummaryControllerTest.java b/src/test/java/com/factory/analysis_service/controller/SummaryControllerTest.java new file mode 100644 index 0000000..89ec2c4 --- /dev/null +++ b/src/test/java/com/factory/analysis_service/controller/SummaryControllerTest.java @@ -0,0 +1,65 @@ +package com.factory.analysis_service.controller; + +import com.factory.analysis_service.dto.MonthlySummaryResponseDTO; +import com.factory.analysis_service.dto.SensorSummaryDTO; +import com.factory.analysis_service.service.S3SummaryService; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.autoconfigure.web.servlet.WebMvcTest; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.bean.override.mockito.MockitoBean; +import org.springframework.test.web.servlet.MockMvc; + +import java.util.List; +import java.util.Map; + +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.when; +import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; + +@WebMvcTest(SummaryController.class) +@ActiveProfiles("test") +@DisplayName("SummaryController") +class SummaryControllerTest { + + @Autowired + private MockMvc mockMvc; + + @MockitoBean + private S3SummaryService s3SummaryService; + + @Test + @DisplayName("월간 요약을 ApiResponse로 감싸 반환한다") + void getMonthlySummary() throws Exception { + MonthlySummaryResponseDTO dto = MonthlySummaryResponseDTO.builder() + .totalDefects(7L) + .totalAnomalies(3L) + .sensors(List.of(SensorSummaryDTO.builder() + .sensorType("TEMP").unit("C").avgValue(20.5).build())) + .build(); + when(s3SummaryService.getMonthlySummary(eq("EQP-01"))).thenReturn(dto); + + mockMvc.perform(get("/api/summary/{equipmentName}/monthly", "EQP-01")) + .andExpect(status().isOk()) + .andExpect(jsonPath("$.success").value(true)) + .andExpect(jsonPath("$.data.totalDefects").value(7)) + .andExpect(jsonPath("$.data.totalAnomalies").value(3)) + .andExpect(jsonPath("$.data.sensors[0].sensorType").value("TEMP")) + .andExpect(jsonPath("$.data.sensors[0].avgValue").value(20.5)); + } + + @Test + @DisplayName("debug 엔드포인트는 S3 진단 맵을 반환한다") + void debug() throws Exception { + when(s3SummaryService.debugS3(eq("EQP-01"))) + .thenReturn(Map.of("bucket", "test-bucket", "foundFiles", "x.parquet")); + + mockMvc.perform(get("/api/summary/debug/{equipmentName}", "EQP-01")) + .andExpect(status().isOk()) + .andExpect(jsonPath("$.bucket").value("test-bucket")) + .andExpect(jsonPath("$.foundFiles").value("x.parquet")); + } +} diff --git a/src/test/java/com/factory/analysis_service/parquet/InMemoryInputFileTest.java b/src/test/java/com/factory/analysis_service/parquet/InMemoryInputFileTest.java new file mode 100644 index 0000000..76700a6 --- /dev/null +++ b/src/test/java/com/factory/analysis_service/parquet/InMemoryInputFileTest.java @@ -0,0 +1,157 @@ +package com.factory.analysis_service.parquet; + +import com.factory.analysis_service.support.ParquetTestSupport; +import org.apache.parquet.column.page.PageReadStore; +import org.apache.parquet.example.data.Group; +import org.apache.parquet.example.data.simple.convert.GroupRecordConverter; +import org.apache.parquet.hadoop.ParquetFileReader; +import org.apache.parquet.io.ColumnIOFactory; +import org.apache.parquet.io.MessageColumnIO; +import org.apache.parquet.io.RecordReader; +import org.apache.parquet.io.SeekableInputStream; +import org.apache.parquet.schema.MessageType; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +@DisplayName("InMemoryInputFile - 인메모리 Parquet 읽기") +class InMemoryInputFileTest { + + @Test + @DisplayName("#5 디스크 임시파일 없이 메모리 바이트에서 Parquet 전체 행을 읽는다") + void readsAllRowsFromMemory() throws Exception { + byte[] bytes = ParquetTestSupport.parquetBytes(List.of( + new ParquetTestSupport.Row("TEMP", "C", 21.5), + new ParquetTestSupport.Row("PRESSURE", "kPa", 101.3), + new ParquetTestSupport.Row("VIBRATION", "mm/s", 0.42) + )); + + List rows = readRows(bytes); + + assertThat(rows).hasSize(3); + assertThat(rows.get(0)).containsExactly("TEMP", "C", "21.5"); + assertThat(rows.get(1)).containsExactly("PRESSURE", "kPa", "101.3"); + assertThat(rows.get(2)).containsExactly("VIBRATION", "mm/s", "0.42"); + } + + @Test + @DisplayName("getLength는 바이트 길이를 반환한다") + void getLengthReturnsByteCount() { + assertThat(new InMemoryInputFile(new byte[]{1, 2, 3, 4}).getLength()).isEqualTo(4L); + } + + @Nested + @DisplayName("SeekableInputStream 동작") + class StreamBehavior { + + private final byte[] data = {10, 20, 30, 40, 50}; + + @Test + @DisplayName("seek 후 getPos가 일치하고 이후 read가 해당 위치부터 읽는다") + void seekAndRead() throws IOException { + try (SeekableInputStream s = new InMemoryInputFile(data).newStream()) { + s.seek(2); + assertThat(s.getPos()).isEqualTo(2L); + assertThat(s.read()).isEqualTo(30); + assertThat(s.getPos()).isEqualTo(3L); + } + } + + @Test + @DisplayName("EOF에서 read()는 -1, read(byte[])는 -1을 반환한다") + void readAtEof() throws IOException { + try (SeekableInputStream s = new InMemoryInputFile(data).newStream()) { + s.seek(5); + assertThat(s.read()).isEqualTo(-1); + assertThat(s.read(new byte[4], 0, 4)).isEqualTo(-1); + assertThat(s.read(ByteBuffer.allocate(4))).isEqualTo(-1); + } + } + + @Test + @DisplayName("readFully(byte[])는 정확히 채우고, 범위를 넘으면 EOFException") + void readFully() throws IOException { + try (SeekableInputStream s = new InMemoryInputFile(data).newStream()) { + byte[] buf = new byte[3]; + s.readFully(buf); + assertThat(buf).containsExactly(10, 20, 30); + assertThat(s.getPos()).isEqualTo(3L); + + assertThatThrownBy(() -> s.readFully(new byte[10])) + .isInstanceOf(IOException.class); + } + } + + @Test + @DisplayName("read(ByteBuffer)와 readFully(ByteBuffer)가 위치를 전진시킨다") + void byteBufferReads() throws IOException { + try (SeekableInputStream s = new InMemoryInputFile(data).newStream()) { + ByteBuffer b1 = ByteBuffer.allocate(2); + assertThat(s.read(b1)).isEqualTo(2); + assertThat(b1.array()).containsExactly(10, 20); + + ByteBuffer b2 = ByteBuffer.allocate(2); + s.readFully(b2); + assertThat(b2.array()).containsExactly(30, 40); + assertThat(s.getPos()).isEqualTo(4L); + } + } + + @Test + @DisplayName("available은 남은 바이트 수를 반환한다") + void available() throws IOException { + try (SeekableInputStream s = new InMemoryInputFile(data).newStream()) { + assertThat(s.available()).isEqualTo(5); + s.seek(3); + assertThat(s.available()).isEqualTo(2); + } + } + + @Test + @DisplayName("음수/범위초과 seek는 IOException을 던진다") + void invalidSeek() throws IOException { + try (SeekableInputStream s = new InMemoryInputFile(data).newStream()) { + assertThatThrownBy(() -> s.seek(-1)).isInstanceOf(IOException.class); + assertThatThrownBy(() -> s.seek(99)).isInstanceOf(IOException.class); + } + } + } + + // ParquetFileReader로 InMemoryInputFile을 읽어 (sensorType, unit, avg_value)를 추출 + private List readRows(byte[] bytes) throws IOException { + List rows = new ArrayList<>(); + try (ParquetFileReader reader = ParquetFileReader.open(new InMemoryInputFile(bytes))) { + MessageType schema = reader.getFooter().getFileMetaData().getSchema(); + MessageColumnIO columnIO = new ColumnIOFactory().getColumnIO(schema); + PageReadStore pages; + while ((pages = reader.readNextRowGroup()) != null) { + RecordReader rr = columnIO.getRecordReader(pages, new GroupRecordConverter(schema)); + long count = pages.getRowCount(); + for (long i = 0; i < count; i++) { + Group g = rr.read(); + rows.add(new String[]{ + g.getString("sensorType", 0), + g.getString("unit", 0), + trimDouble(g.getDouble("avg_value", 0)) + }); + } + } + } + return rows; + } + + private String trimDouble(double v) { + if (v == Math.rint(v)) { + return String.valueOf((long) v); + } + return String.valueOf(v); + } +} diff --git a/src/test/java/com/factory/analysis_service/service/S3SummaryServiceIntegrationTest.java b/src/test/java/com/factory/analysis_service/service/S3SummaryServiceIntegrationTest.java new file mode 100644 index 0000000..953feeb --- /dev/null +++ b/src/test/java/com/factory/analysis_service/service/S3SummaryServiceIntegrationTest.java @@ -0,0 +1,116 @@ +package com.factory.analysis_service.service; + +import com.factory.analysis_service.dto.MonthlySummaryResponseDTO; +import com.factory.analysis_service.entity.DailySensorSummary; +import com.factory.analysis_service.repository.DailySensorSummaryRepository; +import com.factory.analysis_service.support.ParquetTestSupport; +import com.factory.analysis_service.support.StubWebClients; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Primary; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.bean.override.mockito.MockitoBean; +import org.springframework.web.reactive.function.client.WebClient; +import software.amazon.awssdk.core.ResponseBytes; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.model.GetObjectRequest; +import software.amazon.awssdk.services.s3.model.GetObjectResponse; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Response; +import software.amazon.awssdk.services.s3.model.S3Object; + +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.clearInvocations; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@SpringBootTest +@ActiveProfiles("test") +@DisplayName("S3SummaryService 통합 테스트 (H2 캐시)") +class S3SummaryServiceIntegrationTest { + + @TestConfiguration + static class StubConfig { + @Bean + @Primary + WebClient stubWebClient() { + return StubWebClients.counting(11, 5); + } + } + + @MockitoBean + private S3Client s3Client; + + @Autowired + private S3SummaryService service; + + @Autowired + private DailySensorSummaryRepository repository; + + @BeforeEach + void setUp() { + repository.deleteAll(); + + ListObjectsV2Response listResp = ListObjectsV2Response.builder() + .contents(S3Object.builder().key("summary-data/part-0.parquet").build()) + .build(); + when(s3Client.listObjectsV2(any(ListObjectsV2Request.class))).thenReturn(listResp); + + byte[] bytes = ParquetTestSupport.parquetBytes(List.of( + new ParquetTestSupport.Row("TEMP", "C", 20.0), + new ParquetTestSupport.Row("PRESSURE", "kPa", 100.0))); + ResponseBytes resp = + ResponseBytes.fromByteArray(GetObjectResponse.builder().build(), bytes); + when(s3Client.getObjectAsBytes(any(GetObjectRequest.class))).thenReturn(resp); + } + + @AfterEach + void tearDown() { + repository.deleteAllInBatch(); + } + + @Test + @DisplayName("#1#2#4 S3 데이터 조회 + 외부 count 통신 결과를 반환하고, 일자별 요약을 H2에 저장한다") + void fetchesAggregatesAndPersistsCache() { + MonthlySummaryResponseDTO result = service.getMonthlySummary("EQP-01"); + + // #1 센서 집계 + assertThat(result.getSensors()).hasSize(2); + // #2 외부 서비스 count + assertThat(result.getTotalDefects()).isEqualTo(11L); + assertThat(result.getTotalAnomalies()).isEqualTo(5L); + + // #4 MariaDB(H2)에 30일 × 2센서 = 60행 저장됨 + List persisted = repository.findAll(); + assertThat(persisted).hasSize(60); + assertThat(persisted).allSatisfy(row -> + assertThat(row.getEquipmentId()).isEqualTo("EQP-01")); + } + + @Test + @DisplayName("#4 두 번째 호출은 H2 캐시를 사용하여 S3를 다시 조회하지 않는다") + void secondCallUsesCacheFromDb() { + // 1차 호출 — S3에서 읽어 캐시 적재 + service.getMonthlySummary("EQP-01"); + assertThat(repository.count()).isEqualTo(60); + + clearInvocations(s3Client); + + // 2차 호출 — 같은 날짜 범위라 전부 캐시 적중 + MonthlySummaryResponseDTO result = service.getMonthlySummary("EQP-01"); + + assertThat(result.getSensors()).hasSize(2); + verify(s3Client, never()).listObjectsV2(any(ListObjectsV2Request.class)); + verify(s3Client, never()).getObjectAsBytes(any(GetObjectRequest.class)); + } +} diff --git a/src/test/java/com/factory/analysis_service/service/S3SummaryServicePerformanceTest.java b/src/test/java/com/factory/analysis_service/service/S3SummaryServicePerformanceTest.java new file mode 100644 index 0000000..7e8ebcf --- /dev/null +++ b/src/test/java/com/factory/analysis_service/service/S3SummaryServicePerformanceTest.java @@ -0,0 +1,120 @@ +package com.factory.analysis_service.service; + +import com.factory.analysis_service.config.AwsS3Properties; +import com.factory.analysis_service.dto.MonthlySummaryResponseDTO; +import com.factory.analysis_service.repository.DailySensorSummaryRepository; +import com.factory.analysis_service.support.ParquetTestSupport; +import com.factory.analysis_service.support.StubWebClients; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.test.util.ReflectionTestUtils; +import software.amazon.awssdk.core.ResponseBytes; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.model.GetObjectRequest; +import software.amazon.awssdk.services.s3.model.GetObjectResponse; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Response; +import software.amazon.awssdk.services.s3.model.S3Object; + +import java.util.Collections; +import java.util.List; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * 성능 테스트: S3 호출에 인위적 지연(각 40ms)을 주입하여 30일 조회가 '직렬'이 아니라 '병렬'로 + * 수행되는지 확인한다. + * + *

+ * 임계값 1500ms로, 직렬이라면 절대 통과할 수 없고 병렬이면 충분히 통과하도록 둔다. + */ +@DisplayName("S3SummaryService 성능 - 병렬 조회 검증") +class S3SummaryServicePerformanceTest { + + private static final long IO_LATENCY_MS = 40; + private static final long SERIAL_LOWER_BOUND_MS = 30 * (IO_LATENCY_MS * 2); // 2400ms + private static final long PARALLEL_THRESHOLD_MS = 1500; + + private S3Client s3Client; + private ExecutorService executor; + private S3SummaryService service; + + @BeforeEach + void setUp() { + s3Client = mock(S3Client.class); + DailySensorSummaryRepository cacheRepository = mock(DailySensorSummaryRepository.class); + when(cacheRepository.findByEquipmentIdAndSummaryDateBetween(any(), any(), any())) + .thenReturn(Collections.emptyList()); + when(cacheRepository.saveAll(any())).thenAnswer(inv -> inv.getArgument(0)); + + executor = Executors.newFixedThreadPool(8); + + AwsS3Properties props = new AwsS3Properties(); + props.getS3().setBucket("test-bucket"); + props.setRegion("ap-northeast-2"); + + byte[] parquet = ParquetTestSupport.parquetBytes( + List.of(new ParquetTestSupport.Row("TEMP", "C", 20.0))); + + ListObjectsV2Response listResp = ListObjectsV2Response.builder() + .contents(S3Object.builder().key("summary-data/part-0.parquet").build()) + .build(); + when(s3Client.listObjectsV2(any(ListObjectsV2Request.class))).thenAnswer(inv -> { + sleep(IO_LATENCY_MS); + return listResp; + }); + ResponseBytes objResp = + ResponseBytes.fromByteArray(GetObjectResponse.builder().build(), parquet); + when(s3Client.getObjectAsBytes(any(GetObjectRequest.class))).thenAnswer(inv -> { + sleep(IO_LATENCY_MS); + return objResp; + }); + + service = new S3SummaryService( + s3Client, props, StubWebClients.counting(0, 0), cacheRepository, executor); + ReflectionTestUtils.setField(service, "managementServiceUrl", "http://mgmt"); + ReflectionTestUtils.setField(service, "anomalyServiceUrl", "http://anomaly"); + } + + @AfterEach + void tearDown() { + executor.shutdownNow(); + } + + @Test + @DisplayName("30일 S3 조회가 병렬로 수행되어 직렬 하한보다 현저히 빠르다") + void monthlyFetchRunsInParallel() { + // 워밍업 1회 (JIT/클래스 로딩 영향 제거) + service.getMonthlySummary("EQP-01"); + + long start = System.nanoTime(); + MonthlySummaryResponseDTO result = service.getMonthlySummary("EQP-01"); + long elapsedMs = (System.nanoTime() - start) / 1_000_000; + + System.out.printf("[perf] 30일 월간요약 소요=%dms (직렬 하한=%dms, 임계=%dms)%n", + elapsedMs, SERIAL_LOWER_BOUND_MS, PARALLEL_THRESHOLD_MS); + + assertThat(result.getSensors()).hasSize(1); + assertThat(elapsedMs) + .as("병렬이면 %dms 미만, 직렬이면 ~%dms", PARALLEL_THRESHOLD_MS, SERIAL_LOWER_BOUND_MS) + .isLessThan(PARALLEL_THRESHOLD_MS); + } + + private static void sleep(long ms) { + try { + Thread.sleep(ms); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/src/test/java/com/factory/analysis_service/service/S3SummaryServiceTest.java b/src/test/java/com/factory/analysis_service/service/S3SummaryServiceTest.java new file mode 100644 index 0000000..c0b8cdd --- /dev/null +++ b/src/test/java/com/factory/analysis_service/service/S3SummaryServiceTest.java @@ -0,0 +1,292 @@ +package com.factory.analysis_service.service; + +import com.factory.analysis_service.config.AwsS3Properties; +import com.factory.analysis_service.dto.MonthlySummaryResponseDTO; +import com.factory.analysis_service.dto.SensorSummaryDTO; +import com.factory.analysis_service.entity.DailySensorSummary; +import com.factory.analysis_service.repository.DailySensorSummaryRepository; +import com.factory.analysis_service.support.ParquetTestSupport; +import com.factory.analysis_service.support.StubWebClients; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.test.util.ReflectionTestUtils; +import software.amazon.awssdk.core.ResponseBytes; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.model.GetObjectRequest; +import software.amazon.awssdk.services.s3.model.GetObjectResponse; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Response; +import software.amazon.awssdk.services.s3.model.S3Object; + +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@DisplayName("S3SummaryService - 월간 요약 집계") +class S3SummaryServiceTest { + + private static final String EQUIPMENT = "EQP-DEPOSITION-001"; + + private S3Client s3Client; + private DailySensorSummaryRepository cacheRepository; + private ExecutorService executor; + private AwsS3Properties props; + private S3SummaryService service; + + @BeforeEach + void setUp() { + s3Client = mock(S3Client.class); + cacheRepository = mock(DailySensorSummaryRepository.class); + executor = Executors.newFixedThreadPool(4); + + props = new AwsS3Properties(); + props.getS3().setBucket("test-bucket"); + props.setRegion("ap-northeast-2"); + + // 기본: 캐시 미적중 + when(cacheRepository.findByEquipmentIdAndSummaryDateBetween(any(), any(), any())) + .thenReturn(Collections.emptyList()); + when(cacheRepository.saveAll(any())).thenAnswer(inv -> new ArrayList<>(inv.getArgument(0))); + + service = new S3SummaryService( + s3Client, props, StubWebClients.counting(7, 3), cacheRepository, executor); + ReflectionTestUtils.setField(service, "managementServiceUrl", "http://mgmt"); + ReflectionTestUtils.setField(service, "anomalyServiceUrl", "http://anomaly"); + } + + @AfterEach + void tearDown() { + executor.shutdownNow(); + } + + private void stubS3Returns(ParquetTestSupport.Row... rows) { + ListObjectsV2Response listResp = ListObjectsV2Response.builder() + .contents(S3Object.builder().key("summary-data/part-0.parquet").build()) + .build(); + when(s3Client.listObjectsV2(any(ListObjectsV2Request.class))).thenReturn(listResp); + + byte[] bytes = ParquetTestSupport.parquetBytes(List.of(rows)); + ResponseBytes response = + ResponseBytes.fromByteArray(GetObjectResponse.builder().build(), bytes); + when(s3Client.getObjectAsBytes(any(GetObjectRequest.class))).thenReturn(response); + } + + @Test + @DisplayName("#1 S3 Parquet에서 센서 데이터를 정상적으로 가져와 30일 평균을 집계한다") + void aggregatesSensorAveragesFromS3() { + // 모든 일자가 동일하게 TEMP=20.0, PRESSURE=100.0 을 반환 + stubS3Returns( + new ParquetTestSupport.Row("TEMP", "C", 20.0), + new ParquetTestSupport.Row("PRESSURE", "kPa", 100.0)); + + MonthlySummaryResponseDTO result = service.getMonthlySummary(EQUIPMENT); + + Map byType = result.getSensors().stream() + .collect(Collectors.toMap(SensorSummaryDTO::getSensorType, s -> s)); + assertThat(byType).containsOnlyKeys("TEMP", "PRESSURE"); + assertThat(byType.get("TEMP").getAvgValue()).isEqualTo(20.0); + assertThat(byType.get("TEMP").getUnit()).isEqualTo("C"); + assertThat(byType.get("PRESSURE").getAvgValue()).isEqualTo(100.0); + + // 30일치를 모두 조회했는지 (병렬 fetch) + verify(s3Client, times(30)).getObjectAsBytes(any(GetObjectRequest.class)); + } + + @Test + @DisplayName("#1 여러 일자의 서로 다른 값을 올바른 평균으로 계산한다 (소수 2자리 반올림)") + void averagesAcrossDifferentDailyValues() { + // 병렬 실행에서도 결정적이도록, 값은 호출 순서가 아니라 '키에 박힌 날짜'로 결정한다. + // listObjectsV2는 prefix(=날짜 포함)를 그대로 키에 실어 돌려준다. + when(s3Client.listObjectsV2(any(ListObjectsV2Request.class))).thenAnswer(inv -> { + ListObjectsV2Request req = inv.getArgument(0); + return ListObjectsV2Response.builder() + .contents(S3Object.builder().key(req.prefix() + "part.parquet").build()) + .build(); + }); + // getObjectAsBytes는 키의 date=YYYY-MM-dd 에서 일(day-of-month)을 값으로 사용 + when(s3Client.getObjectAsBytes(any(GetObjectRequest.class))).thenAnswer(inv -> { + GetObjectRequest req = inv.getArgument(0); + double v = dayFromKey(req.key()); + byte[] bytes = ParquetTestSupport.parquetBytes( + List.of(new ParquetTestSupport.Row("TEMP", "C", v))); + return ResponseBytes.fromByteArray(GetObjectResponse.builder().build(), bytes); + }); + + // 기대값: 최근 30일(어제~30일 전)의 day-of-month 평균을 소수 2자리로 반올림 + LocalDate today = LocalDate.now(); + double expected = IntStream.rangeClosed(1, 30) + .mapToDouble(i -> today.minusDays(i).getDayOfMonth()) + .average().orElse(0.0); + expected = Math.round(expected * 100.0) / 100.0; + + MonthlySummaryResponseDTO result = service.getMonthlySummary(EQUIPMENT); + + assertThat(result.getSensors()).hasSize(1); + assertThat(result.getSensors().get(0).getAvgValue()).isEqualTo(expected); + } + + private static int dayFromKey(String key) { + int idx = key.indexOf("date="); + String date = key.substring(idx + 5, idx + 15); // YYYY-MM-dd + return Integer.parseInt(date.substring(8, 10)); + } + + @Test + @DisplayName("#2 management/anomaly 서비스와 통신해 defect/anomaly count를 가져온다") + void fetchesCountsFromOtherServices() { + stubS3Returns(new ParquetTestSupport.Row("TEMP", "C", 20.0)); + + MonthlySummaryResponseDTO result = service.getMonthlySummary(EQUIPMENT); + + assertThat(result.getTotalDefects()).isEqualTo(7L); + assertThat(result.getTotalAnomalies()).isEqualTo(3L); + } + + @Test + @DisplayName("#2 외부 서비스 실패/응답누락 시 count는 0으로 폴백한다") + void fallsBackToZeroOnExternalFailure() { + stubS3Returns(new ParquetTestSupport.Row("TEMP", "C", 20.0)); + service = new S3SummaryService(s3Client, props, + StubWebClients.failingManagementAndNullAnomaly(), cacheRepository, executor); + ReflectionTestUtils.setField(service, "managementServiceUrl", "http://mgmt"); + ReflectionTestUtils.setField(service, "anomalyServiceUrl", "http://anomaly"); + + MonthlySummaryResponseDTO result = service.getMonthlySummary(EQUIPMENT); + + assertThat(result.getTotalDefects()).isZero(); + assertThat(result.getTotalAnomalies()).isZero(); + // 외부 호출이 실패해도 센서 집계는 정상 + assertThat(result.getSensors()).hasSize(1); + } + + @Test + @DisplayName("#4 S3에서 읽은 일자별 데이터를 캐시에 저장한다") + void persistsFetchedDataToCache() { + stubS3Returns( + new ParquetTestSupport.Row("TEMP", "C", 20.0), + new ParquetTestSupport.Row("PRESSURE", "kPa", 100.0)); + + service.getMonthlySummary(EQUIPMENT); + + // 30일 × 2센서 = 60행 저장 + @SuppressWarnings("unchecked") + java.util.List saved = captureSaved(); + assertThat(saved).hasSize(60); + assertThat(saved).allSatisfy(row -> { + assertThat(row.getEquipmentId()).isEqualTo(EQUIPMENT); + assertThat(row.getSensorType()).isIn("TEMP", "PRESSURE"); + }); + } + + @SuppressWarnings("unchecked") + private List captureSaved() { + var captor = org.mockito.ArgumentCaptor.forClass(List.class); + verify(cacheRepository).saveAll(captor.capture()); + return captor.getValue(); + } + + @Test + @DisplayName("#4 캐시가 적중하면 S3를 조회하지 않고 캐시 값으로 집계한다") + void usesCacheWhenPresentAndSkipsS3() { + LocalDate today = LocalDate.now(); + List cached = new ArrayList<>(); + for (int i = 1; i <= 30; i++) { + cached.add(DailySensorSummary.builder() + .equipmentId(EQUIPMENT).summaryDate(today.minusDays(i)) + .sensorType("TEMP").unit("C").avgValue(25.0).build()); + } + when(cacheRepository.findByEquipmentIdAndSummaryDateBetween(any(), any(), any())) + .thenReturn(cached); + + MonthlySummaryResponseDTO result = service.getMonthlySummary(EQUIPMENT); + + assertThat(result.getSensors()).hasSize(1); + assertThat(result.getSensors().get(0).getAvgValue()).isEqualTo(25.0); + // 전부 캐시 적중 → S3 호출 없음, 저장도 없음 + verify(s3Client, never()).listObjectsV2(any(ListObjectsV2Request.class)); + verify(s3Client, never()).getObjectAsBytes(any(GetObjectRequest.class)); + verify(cacheRepository, never()).saveAll(any()); + } + + @Test + @DisplayName("S3 prefix에 parquet이 없으면 해당 일자는 건너뛴다") + void skipsDatesWithoutParquet() { + when(s3Client.listObjectsV2(any(ListObjectsV2Request.class))) + .thenReturn(ListObjectsV2Response.builder().contents(Collections.emptyList()).build()); + + MonthlySummaryResponseDTO result = service.getMonthlySummary(EQUIPMENT); + + assertThat(result.getSensors()).isEmpty(); + verify(s3Client, never()).getObjectAsBytes(any(GetObjectRequest.class)); + // 저장할 것도 없음 + verify(cacheRepository, never()).saveAll(any()); + } + + @Test + @DisplayName("debugS3는 파일을 찾으면 스키마/행수/첫행을 진단 맵에 담는다") + void debugS3ReturnsDiagnostics() { + stubS3Returns(new ParquetTestSupport.Row("TEMP", "C", 20.0)); + + Map result = service.debugS3(EQUIPMENT); + + assertThat(result.get("bucket")).isEqualTo("test-bucket"); + assertThat(result.get("region")).isEqualTo("ap-northeast-2"); + assertThat(result.get("foundFiles")).isNotBlank(); + assertThat(result.get("rowCount")).isEqualTo("1"); + assertThat(result).containsKey("schema"); + assertThat(result).containsKey("firstRow"); + } + + @Test + @DisplayName("debugS3는 parquet이 없으면 foundFiles=없음을 반환한다") + void debugS3NoFiles() { + when(s3Client.listObjectsV2(any(ListObjectsV2Request.class))) + .thenReturn(ListObjectsV2Response.builder().contents(Collections.emptyList()).build()); + + Map result = service.debugS3(EQUIPMENT); + + assertThat(result.get("foundFiles")).isEqualTo("없음"); + assertThat(result).doesNotContainKey("schema"); + } + + @Test + @DisplayName("#3 동시 호출에도 일관된 결과를 반환한다 (병렬 수행 안정성)") + void concurrentCallsAreConsistent() throws Exception { + stubS3Returns(new ParquetTestSupport.Row("TEMP", "C", 20.0)); + + ExecutorService callers = Executors.newFixedThreadPool(8); + try { + List> futures = IntStream.range(0, 16) + .mapToObj(i -> CompletableFuture.supplyAsync( + () -> service.getMonthlySummary(EQUIPMENT), callers)) + .collect(Collectors.toList()); + + for (CompletableFuture f : futures) { + MonthlySummaryResponseDTO r = f.get(); + assertThat(r.getSensors()).hasSize(1); + assertThat(r.getSensors().get(0).getAvgValue()).isEqualTo(20.0); + assertThat(r.getTotalDefects()).isEqualTo(7L); + assertThat(r.getTotalAnomalies()).isEqualTo(3L); + } + } finally { + callers.shutdownNow(); + } + } +} diff --git a/src/test/java/com/factory/analysis_service/support/InMemoryOutputFile.java b/src/test/java/com/factory/analysis_service/support/InMemoryOutputFile.java new file mode 100644 index 0000000..577dc63 --- /dev/null +++ b/src/test/java/com/factory/analysis_service/support/InMemoryOutputFile.java @@ -0,0 +1,62 @@ +package com.factory.analysis_service.support; + +import org.apache.parquet.io.OutputFile; +import org.apache.parquet.io.PositionOutputStream; + +import java.io.ByteArrayOutputStream; + +/** + * 테스트에서 Parquet 바이트를 Hadoop FileSystem(winutils 의존) 없이 메모리에 직접 쓰기 위한 OutputFile. + * 프로덕션 읽기 경로({@code InMemoryInputFile})를 디스크 없이 검증하는 데 사용한다. + */ +public class InMemoryOutputFile implements OutputFile { + + private final ByteArrayOutputStream baos = new ByteArrayOutputStream(); + + @Override + public PositionOutputStream create(long blockSizeHint) { + return stream(); + } + + @Override + public PositionOutputStream createOrOverwrite(long blockSizeHint) { + return stream(); + } + + @Override + public boolean supportsBlockSize() { + return false; + } + + @Override + public long defaultBlockSize() { + return 0; + } + + public byte[] toByteArray() { + return baos.toByteArray(); + } + + private PositionOutputStream stream() { + return new PositionOutputStream() { + private long pos = 0; + + @Override + public long getPos() { + return pos; + } + + @Override + public void write(int b) { + baos.write(b); + pos++; + } + + @Override + public void write(byte[] b, int off, int len) { + baos.write(b, off, len); + pos += len; + } + }; + } +} diff --git a/src/test/java/com/factory/analysis_service/support/ParquetTestSupport.java b/src/test/java/com/factory/analysis_service/support/ParquetTestSupport.java new file mode 100644 index 0000000..bcbf65a --- /dev/null +++ b/src/test/java/com/factory/analysis_service/support/ParquetTestSupport.java @@ -0,0 +1,59 @@ +package com.factory.analysis_service.support; + +import org.apache.parquet.example.data.Group; +import org.apache.parquet.example.data.simple.SimpleGroupFactory; +import org.apache.parquet.hadoop.ParquetWriter; +import org.apache.parquet.hadoop.example.ExampleParquetWriter; +import org.apache.parquet.hadoop.metadata.CompressionCodecName; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.Types; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.util.List; + +import static org.apache.parquet.schema.LogicalTypeAnnotation.stringType; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.BINARY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.DOUBLE; + +/** + * 프로덕션이 읽는 것과 동일한 스키마(sensorType/unit/avg_value)의 Parquet 바이트를 생성하는 테스트 헬퍼. + */ +public final class ParquetTestSupport { + + public static final MessageType SCHEMA = Types.buildMessage() + .required(BINARY).as(stringType()).named("sensorType") + .required(BINARY).as(stringType()).named("unit") + .required(DOUBLE).named("avg_value") + .named("summary"); + + private ParquetTestSupport() { + } + + public record Row(String sensorType, String unit, double avgValue) { + } + + /** + * 주어진 행들로 압축 없는 Parquet 파일 바이트를 만든다. (UNCOMPRESSED — 네이티브 코덱 의존 제거) + */ + public static byte[] parquetBytes(List rows) { + InMemoryOutputFile outputFile = new InMemoryOutputFile(); + SimpleGroupFactory factory = new SimpleGroupFactory(SCHEMA); + + try (ParquetWriter writer = ExampleParquetWriter.builder(outputFile) + .withType(SCHEMA) + .withCompressionCodec(CompressionCodecName.UNCOMPRESSED) + .build()) { + for (Row row : rows) { + Group group = factory.newGroup() + .append("sensorType", row.sensorType()) + .append("unit", row.unit()) + .append("avg_value", row.avgValue()); + writer.write(group); + } + } catch (IOException e) { + throw new UncheckedIOException(e); + } + return outputFile.toByteArray(); + } +} diff --git a/src/test/java/com/factory/analysis_service/support/StubWebClients.java b/src/test/java/com/factory/analysis_service/support/StubWebClients.java new file mode 100644 index 0000000..f57302b --- /dev/null +++ b/src/test/java/com/factory/analysis_service/support/StubWebClients.java @@ -0,0 +1,58 @@ +package com.factory.analysis_service.support; + +import org.springframework.http.HttpStatus; +import org.springframework.http.HttpHeaders; +import org.springframework.http.MediaType; +import org.springframework.web.reactive.function.client.ClientResponse; +import org.springframework.web.reactive.function.client.ExchangeFunction; +import org.springframework.web.reactive.function.client.WebClient; +import reactor.core.publisher.Mono; + +/** + * 실제 네트워크 없이 management/anomaly count 응답을 흉내내는 WebClient 스텁. + * URL 경로로 라우팅하여 ApiResponse<Long> JSON을 돌려준다. + */ +public final class StubWebClients { + + private StubWebClients() { + } + + /** 두 서비스가 각각 defects/anomalies count를 정상 반환하는 WebClient. */ + public static WebClient counting(long defects, long anomalies) { + return fromExchange(request -> { + String url = request.url().toString(); + if (url.contains("/defects/count")) { + return Mono.just(json("{\"success\":true,\"status\":200,\"data\":" + defects + "}")); + } + if (url.contains("/anomalies/count")) { + return Mono.just(json("{\"success\":true,\"status\":200,\"data\":" + anomalies + "}")); + } + return Mono.just(ClientResponse.create(HttpStatus.NOT_FOUND).build()); + }); + } + + /** management는 5xx로 실패, anomaly는 data 필드가 없는(null) 응답 — 둘 다 0으로 폴백되어야 함. */ + public static WebClient failingManagementAndNullAnomaly() { + return fromExchange(request -> { + String url = request.url().toString(); + if (url.contains("/defects/count")) { + return Mono.just(ClientResponse.create(HttpStatus.INTERNAL_SERVER_ERROR).build()); + } + if (url.contains("/anomalies/count")) { + return Mono.just(json("{\"success\":true,\"status\":200}")); // data 없음 → null + } + return Mono.just(ClientResponse.create(HttpStatus.NOT_FOUND).build()); + }); + } + + private static WebClient fromExchange(ExchangeFunction exchange) { + return WebClient.builder().exchangeFunction(exchange).build(); + } + + private static ClientResponse json(String body) { + return ClientResponse.create(HttpStatus.OK) + .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) + .body(body) + .build(); + } +} diff --git a/src/test/resources/application-test.yml b/src/test/resources/application-test.yml new file mode 100644 index 0000000..35ae8be --- /dev/null +++ b/src/test/resources/application-test.yml @@ -0,0 +1,43 @@ +server: + port: 0 + +cloud: + aws: + s3: + bucket: test-bucket + region: ap-northeast-2 + credentials: + access-key: "" + secret-key: "" + +spring: + application: + name: analysis-service + + datasource: + driver-class-name: org.h2.Driver + url: jdbc:h2:mem:analysisdb;MODE=MariaDB;DB_CLOSE_DELAY=-1;NON_KEYWORDS=VALUE + username: sa + password: "" + + jpa: + hibernate: + ddl-auto: create-drop + show-sql: false + properties: + hibernate: + dialect: org.hibernate.dialect.H2Dialect + + autoconfigure: + exclude: + - org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration + - org.springframework.kafka.annotation.EnableKafkaBootstrapConfiguration + # 요약 기능과 무관하며, kafka 공통모듈(EventPublisher) 미의존으로 컨텍스트 로딩을 깨는 오토컨피그 제외 + - com.factory.common.outbox.jpa.autoconfigure.OutboxJpaAutoConfiguration + - com.factory.common.inbox.jpa.autoconfigure.InboxJpaAutoConfiguration + +services: + management-service: + url: http://localhost:18086 + anomaly-service: + url: http://localhost:18085