diff --git a/build.gradle b/build.gradle index 29994c5..48471f1 100644 --- a/build.gradle +++ b/build.gradle @@ -2,6 +2,7 @@ plugins { id 'java' id 'org.springframework.boot' version '3.5.2' id 'io.spring.dependency-management' version '1.1.7' + id 'jacoco' } group = 'com.holliverse' @@ -27,17 +28,31 @@ repositories { dependencies { implementation 'org.springframework.boot:spring-boot-starter-web' implementation 'org.springframework.kafka:spring-kafka' - implementation 'org.springframework.boot:spring-boot-starter-data-redis' + implementation 'org.springframework.boot:spring-boot-starter-data-jpa' implementation 'com.fasterxml.jackson.core:jackson-databind' + runtimeOnly 'org.postgresql:postgresql' compileOnly 'org.projectlombok:lombok' annotationProcessor 'org.projectlombok:lombok' testImplementation 'org.springframework.boot:spring-boot-starter-test' testImplementation 'org.springframework.kafka:spring-kafka-test' + testImplementation 'org.testcontainers:postgresql' + testImplementation 'org.testcontainers:junit-jupiter' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' } tasks.named('test') { useJUnitPlatform() } + +jacoco { + toolVersion = '0.8.11' +} + +tasks.named('jacocoTestReport') { + reports { + xml.required = true + html.required = true + } +} diff --git a/docker-compose.local.yml b/docker-compose.local.yml index ab2d1ea..856ab35 100644 --- a/docker-compose.local.yml +++ b/docker-compose.local.yml @@ -40,28 +40,7 @@ services: depends_on: - kafka - redis: # ECS 0.5 vCPU 환경 - image: redis:7.2-alpine - container_name: redis - ports: - - "6379:6379" # 영속성 설정, 256mb이상 쓰지 못하게 제한 - command: > - redis-server - --appendonly yes - --appendfsync everysec - --maxmemory 256mb - --maxmemory-policy allkeys-lru - --requirepass ${REDIS_PASSWORD:-local_redis_pass} - environment: - REDIS_PASSWORD: ${REDIS_PASSWORD:-local_redis_pass} - healthcheck: - test: ["CMD-SHELL", "redis-cli -a ${REDIS_PASSWORD:-local_redis_pass} ping | grep PONG"] - interval: 10s - timeout: 3s - retries: 5 - volumes: - - redis_data:/data + volumes: kafka_data: - redis_data: \ No newline at end of file diff --git a/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java b/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java index 222bab0..74d0eb2 100644 --- a/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java +++ b/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java @@ -3,7 +3,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.holliverse.logserver.config.properties.KafkaAppProperties; import com.holliverse.logserver.dto.LogEvent; -import com.holliverse.logserver.service.RedisLogService; +import com.holliverse.logserver.service.PostgresLogService; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; @@ -17,7 +17,7 @@ public class SpeedLayerConsumer { private final ObjectMapper objectMapper; - private final RedisLogService redisLogService; + private final PostgresLogService postgresLogService; private final KafkaTemplate dlqKafkaTemplate; private final KafkaAppProperties kafkaAppProperties; @@ -32,9 +32,9 @@ public class SpeedLayerConsumer { public void consume(ConsumerRecord record) { try { LogEvent event = objectMapper.readValue(record.value(), LogEvent.class); - redisLogService.process(event); + postgresLogService.process(event); } catch (Exception e) { - // 역직렬화 or Redis 처리 실패 → DLQ 토픽으로 원본 페이로드 전송 + // 역직렬화 or PostgreSQL 처리 실패 → DLQ 토픽으로 원본 페이로드 전송 // 예외를 re-throw하지 않아 다음 메시지 처리가 중단되지 않음 (무한 루프 방지) dlqKafkaTemplate.send( kafkaAppProperties.getTopics().getError(), diff --git a/src/main/java/com/holliverse/logserver/dto/LogEvent.java b/src/main/java/com/holliverse/logserver/dto/LogEvent.java index 43e1493..1ee9526 100644 --- a/src/main/java/com/holliverse/logserver/dto/LogEvent.java +++ b/src/main/java/com/holliverse/logserver/dto/LogEvent.java @@ -8,7 +8,7 @@ public class LogEvent { @JsonProperty("event_id") - private String eventId; + private Long eventId; private String timestamp; private String event; @@ -16,15 +16,9 @@ public class LogEvent { @JsonProperty("event_name") private String eventName; - @JsonProperty("member_properties") - private MemberProperties memberProperties; + @JsonProperty("member_id") + private Long memberId; @JsonProperty("event_properties") private Map eventProperties; - - @Data - public static class MemberProperties { - @JsonProperty("member_id") - private String memberId; - } } diff --git a/src/main/java/com/holliverse/logserver/entity/ProductViewHistory.java b/src/main/java/com/holliverse/logserver/entity/ProductViewHistory.java new file mode 100644 index 0000000..02db801 --- /dev/null +++ b/src/main/java/com/holliverse/logserver/entity/ProductViewHistory.java @@ -0,0 +1,39 @@ +package com.holliverse.logserver.entity; + +import jakarta.persistence.Column; +import jakarta.persistence.EmbeddedId; +import jakarta.persistence.Entity; +import jakarta.persistence.Table; +import java.time.OffsetDateTime; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Getter; +import lombok.NoArgsConstructor; + +@Entity +@Table(name = "product_view_history") +@Getter +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class ProductViewHistory { + + @EmbeddedId + private ProductViewHistoryId id; + + @Column(name = "product_name", nullable = false, length = 100) + private String productName; + + @Column(name = "product_type", nullable = false, length = 50) + private String productType; + + // List을 JSON 문자열로 직렬화해 저장. DB DDL에서 jsonb 타입으로 선언됨. + @Column(columnDefinition = "jsonb") + private String tags; + + @Column(name = "viewed_at", nullable = false, columnDefinition = "TIMESTAMPTZ") + private OffsetDateTime viewedAt; + + @Column(name = "last_event_id", nullable = false) + private Long lastEventId; +} diff --git a/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java b/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java new file mode 100644 index 0000000..e7b57e8 --- /dev/null +++ b/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java @@ -0,0 +1,24 @@ +package com.holliverse.logserver.entity; + +import jakarta.persistence.Column; +import jakarta.persistence.Embeddable; +import java.io.Serializable; +import lombok.AllArgsConstructor; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; + +@Embeddable +@Getter +@NoArgsConstructor +@AllArgsConstructor +@EqualsAndHashCode +// (member_id, product_id) 복합 PK @Embeddable +public class ProductViewHistoryId implements Serializable { + + @Column(name = "member_id", nullable = false) + private Long memberId; + + @Column(name = "product_id", nullable = false) + private Long productId; +} diff --git a/src/main/java/com/holliverse/logserver/repository/ProductViewHistoryRepository.java b/src/main/java/com/holliverse/logserver/repository/ProductViewHistoryRepository.java new file mode 100644 index 0000000..fd6f018 --- /dev/null +++ b/src/main/java/com/holliverse/logserver/repository/ProductViewHistoryRepository.java @@ -0,0 +1,64 @@ +package com.holliverse.logserver.repository; + +import com.holliverse.logserver.entity.ProductViewHistory; +import com.holliverse.logserver.entity.ProductViewHistoryId; +import java.time.OffsetDateTime; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +public interface ProductViewHistoryRepository + extends JpaRepository { + + /** + * 복합 PK (member_id, product_id) 기준 UPSERT. + * 동일 사용자가 동일 상품을 재조회하면 나머지 컬럼 전체를 최신값으로 덮어씀. + * tags는 Java String → CAST AS jsonb 로 PostgreSQL 레벨에서 변환됨. + */ + @Modifying + @Query(value = """ + INSERT INTO product_view_history + (member_id, product_id, product_name, product_type, tags, viewed_at, last_event_id) + VALUES + (:memberId, :productId, :productName, :productType, + CAST(:tags AS jsonb), :viewedAt, :lastEventId) + ON CONFLICT (member_id, product_id) + DO UPDATE SET + product_name = EXCLUDED.product_name, + product_type = EXCLUDED.product_type, + tags = EXCLUDED.tags, + viewed_at = EXCLUDED.viewed_at, + last_event_id = EXCLUDED.last_event_id + """, nativeQuery = true) + void upsert( + @Param("memberId") Long memberId, + @Param("productId") Long productId, + @Param("productName") String productName, + @Param("productType") String productType, + @Param("tags") String tags, + @Param("viewedAt") OffsetDateTime viewedAt, + @Param("lastEventId") Long lastEventId + ); + + /** + * 유저당 최신 N개(viewed_at DESC)를 초과하는 오래된 레코드 삭제. + * 복합 PK 구조라 id 컬럼이 없으므로 product_id NOT IN 서브쿼리로 대상을 특정함. + */ + @Modifying + @Query(value = """ + DELETE FROM product_view_history + WHERE member_id = :memberId + AND product_id NOT IN ( + SELECT product_id + FROM product_view_history + WHERE member_id = :memberId + ORDER BY viewed_at DESC + LIMIT :maxCount + ) + """, nativeQuery = true) + void trimOldRecords( + @Param("memberId") Long memberId, + @Param("maxCount") int maxCount + ); +} diff --git a/src/main/java/com/holliverse/logserver/service/PostgresLogService.java b/src/main/java/com/holliverse/logserver/service/PostgresLogService.java new file mode 100644 index 0000000..b48b23d --- /dev/null +++ b/src/main/java/com/holliverse/logserver/service/PostgresLogService.java @@ -0,0 +1,78 @@ +package com.holliverse.logserver.service; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.holliverse.logserver.dto.LogEvent; +import com.holliverse.logserver.repository.ProductViewHistoryRepository; +import java.time.OffsetDateTime; +import java.util.Collections; +import java.util.List; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.Setter; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Slf4j +@Service +@RequiredArgsConstructor +public class PostgresLogService { + + private static final int MAX_RECENT_VIEWS = 30; + + private final ProductViewHistoryRepository repository; + private final ObjectMapper objectMapper; + + /** + * click_product_detail 이벤트 1건을 product_view_history 테이블에 UPSERT 후 + * 해당 유저의 레코드가 MAX_RECENT_VIEWS 를 초과하면 오래된 항목 자동 삭제. + */ + @Transactional + public void process(LogEvent event) throws JsonProcessingException { + Long memberId = event.getMemberId(); + ClickProductDetailProperties props = objectMapper.convertValue( + event.getEventProperties(), + ClickProductDetailProperties.class + ); + + Long productId = props.getProductId(); + if (productId == null) { + throw new IllegalArgumentException("product_id 변환 실패: null"); + } + String productName = props.getProductName(); + String productType = props.getProductType(); + List tagList = props.getTags() != null ? props.getTags() : Collections.emptyList(); + String tags = objectMapper.writeValueAsString(tagList); + + OffsetDateTime viewedAt = OffsetDateTime.parse(event.getTimestamp()); + Long lastEventId = event.getEventId(); + + repository.upsert(memberId, productId, productName, productType, tags, viewedAt, lastEventId); + repository.trimOldRecords(memberId, MAX_RECENT_VIEWS); + + log.debug("[PostgresLog] UPSERT 완료 memberId={} productId={}", memberId, productId); + } + + /** + * click_product_detail 이벤트의 event_properties 전용 DTO. + * ObjectMapper.convertValue()로 변환하여 Map 직접 캐스팅 및 ClassCastException을 방지. + * JSON 키(snake_case)는 @JsonProperty로 Java CamelCase 필드에 매핑. + */ + @Getter + @Setter + private static class ClickProductDetailProperties { + + @JsonProperty("product_id") + private Long productId; + + @JsonProperty("product_name") + private String productName; + + @JsonProperty("product_type") + private String productType; + + private List tags; + } +} diff --git a/src/main/java/com/holliverse/logserver/service/RedisLogService.java b/src/main/java/com/holliverse/logserver/service/RedisLogService.java deleted file mode 100644 index 7131234..0000000 --- a/src/main/java/com/holliverse/logserver/service/RedisLogService.java +++ /dev/null @@ -1,81 +0,0 @@ -package com.holliverse.logserver.service; - -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.holliverse.logserver.dto.LogEvent; -import java.time.Instant; -import java.util.Collections; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; -import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.springframework.data.redis.core.StringRedisTemplate; -import org.springframework.stereotype.Service; - -@Slf4j -@Service -@RequiredArgsConstructor -public class RedisLogService { - - private final StringRedisTemplate redisTemplate; - private final ObjectMapper objectMapper; - - // recent_views 최대 보관 개수 - private static final int MAX_RECENT_VIEWS = 30; - - /** - * click_product_detail 이벤트 1건에 대해 2개 ZSET 동시 업데이트 - * 1) user:{member_id}:recent_views — 최근 본 상품 (최신순 정렬, 중복 제거) - * 2) user:{member_id}:tag_scores — 태그 선호도 누적 점수 - */ - public void process(LogEvent event) throws JsonProcessingException { - String memberId = event.getMemberProperties().getMemberId(); - Map props = event.getEventProperties(); - - updateRecentViews(memberId, props, event.getTimestamp()); - updateTagScores(memberId, props); - } - - /** - * ZSET: user:{member_id}:recent_views - * Score = 이벤트 타임스탬프 (epoch ms) → 높을수록 최신 - * Value = 프론트엔드 렌더링용 최소 JSON string - * 최대 30개 유지: ZREMRANGEBYRANK로 오래된 항목 자동 삭제 - */ - private void updateRecentViews(String memberId, Map props, String timestamp) - throws JsonProcessingException { - - String key = "user:" + memberId + ":recent_views"; - - // ISO 8601 → epoch ms (ZSET score는 double) - double score = (double) Instant.parse(timestamp).toEpochMilli(); - - // 프론트가 화면을 그릴 최소 데이터만 포함 - Map viewData = new LinkedHashMap<>(); - viewData.put("product_id", props.get("product_id")); - viewData.put("target_url", props.getOrDefault("page_url", "")); - String valueJson = objectMapper.writeValueAsString(viewData); - - // ZADD: 이미 같은 value가 있으면 score(timestamp)만 갱신 → 멱등성 보장 - redisTemplate.opsForZSet().add(key, valueJson, score); - - // 최대 30개 초과분 삭제: rank 0(가장 오래된) ~ -(MAX+1) 범위 제거 - // ex) 31개가 되는 순간 rank 0 항목 1개가 삭제되어 항상 30개 유지 - redisTemplate.opsForZSet().removeRange(key, 0, -(MAX_RECENT_VIEWS + 1)); - } - - /** - * ZSET: user:{member_id}:tag_scores - * ZINCRBY: 태그마다 +1점씩 누적 - * 같은 태그가 중복 도착해도 score 합산으로 자연스럽게 선호도 반영 - */ - @SuppressWarnings("unchecked") - private void updateTagScores(String memberId, Map props) { - String key = "user:" + memberId + ":tag_scores"; - List tags = (List) props.getOrDefault("tags", Collections.emptyList()); - for (String tag : tags) { - redisTemplate.opsForZSet().incrementScore(key, tag, 1.0); - } - } -} diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index dbdf3fd..65f5f50 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -2,11 +2,17 @@ spring: application: name: log-server - data: - redis: - host: ${REDIS_HOST:localhost} - port: ${REDIS_PORT:6379} - password: ${REDIS_PASSWORD:local_redis_pass} + datasource: + url: ${DB_URL:jdbc:postgresql://localhost:5435/holliverse} + username: ${DB_USERNAME:postgres} + password: ${DB_PASSWORD:postgres} + driver-class-name: org.postgresql.Driver + jpa: + hibernate: + ddl-auto: ${JPA_DDL_AUTO:validate} + properties: + hibernate.dialect: org.hibernate.dialect.PostgreSQLDialect + show-sql: false app: kafka: @@ -22,3 +28,4 @@ app: producer: dlq-acks: ${KAFKA_DLQ_ACKS:all} dlq-retries: ${KAFKA_DLQ_RETRIES:3} + diff --git a/src/test/java/com/holliverse/logserver/consumer/SpeedLayerConsumerIntegrationTest.java b/src/test/java/com/holliverse/logserver/consumer/SpeedLayerConsumerIntegrationTest.java new file mode 100644 index 0000000..2361573 --- /dev/null +++ b/src/test/java/com/holliverse/logserver/consumer/SpeedLayerConsumerIntegrationTest.java @@ -0,0 +1,335 @@ +package com.holliverse.logserver.consumer; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +import com.holliverse.logserver.entity.ProductViewHistory; +import com.holliverse.logserver.entity.ProductViewHistoryId; +import com.holliverse.logserver.repository.ProductViewHistoryRepository; +import java.time.Duration; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.TimeUnit; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.apache.kafka.common.serialization.StringSerializer; +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.kafka.core.DefaultKafkaConsumerFactory; +import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.test.EmbeddedKafkaBroker; +import org.springframework.kafka.test.context.EmbeddedKafka; +import org.springframework.kafka.test.utils.KafkaTestUtils; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; +import org.testcontainers.containers.PostgreSQLContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; + +/** + * Kafka → SpeedLayerConsumer → PostgreSQL 전체 파이프라인 통합 테스트. + * + * 인프라: + * - EmbeddedKafka : 인메모리 Kafka 브로커 (포트 랜덤) + * - PostgreSQLContainer : TestContainers PostgreSQL 16 + * + * 부하 근거 (Case 6): + * - 일 사용자 3만 명 × 10회/일 = 30만 건/일 + * - 30만 / 86,400초 ≈ 3.47 TPS (평균) + * - 피크타임(전체의 30% 집중, 4시간) ≈ 6.25 TPS + * → MAX_POLL_RECORDS=1, ACK_MODE=RECORD 설정으로도 충분히 처리 가능함을 검증 + */ +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE) +@EmbeddedKafka( + partitions = 1, + topics = {"client-event-logs", "error-logs"}, + bootstrapServersProperty = "app.kafka.bootstrap-servers" +) +@Testcontainers +@DirtiesContext +class SpeedLayerConsumerIntegrationTest { + + // ---------------------------------------------------------------- + // TestContainers PostgreSQL (클래스 전체 공유, 최초 1회 기동) + // ---------------------------------------------------------------- + @Container + static PostgreSQLContainer postgres = + new PostgreSQLContainer<>("postgres:16-alpine") + .withDatabaseName("logdb") + .withUsername("loguser") + .withPassword("logpass"); + + // Hibernate ddl-auto=create → 테이블 자동 생성 (테스트 전용) + @DynamicPropertySource + static void overrideProps(DynamicPropertyRegistry registry) { + registry.add("spring.datasource.url", postgres::getJdbcUrl); + registry.add("spring.datasource.username", postgres::getUsername); + registry.add("spring.datasource.password", postgres::getPassword); + registry.add("spring.jpa.hibernate.ddl-auto", () -> "create"); + } + + // @EmbeddedKafka 가 테스트 컨텍스트에서만 빈을 등록하므로 IDE가 미인식 → 런타임에는 정상 주입됨 + @Autowired + @SuppressWarnings("SpringJavaInjectionPointsAutowiringInspection") + private EmbeddedKafkaBroker embeddedKafkaBroker; + + @Autowired + private ProductViewHistoryRepository repository; + + // 테스트에서 client-event-logs 토픽으로 메시지를 발행할 프로듀서 + private KafkaTemplate testProducer; + + @BeforeEach + void setUp() { + Map producerProps = KafkaTestUtils.producerProps(embeddedKafkaBroker); + producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + testProducer = new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(producerProps)); + + repository.deleteAll(); + } + + // ================================================================ + // Case 1: 신규 삽입 — DB에 1건 적재 확인 + // ================================================================ + @Test + @DisplayName("Case 1 | 신규 삽입 — product_view_history에 1건 적재") + void case1_normalInsert_savedToDatabase() { + String payload = buildPayload(1001L, 45L, 10L, + "2026-03-02T16:30:00.000Z", "5G 요금제", "mobile", "[\"영상OTT\",\"인기\"]"); + + testProducer.send("client-event-logs", payload); + + await().atMost(15, TimeUnit.SECONDS) + .untilAsserted(() -> { + ProductViewHistory row = repository + .findById(new ProductViewHistoryId(45L, 10L)) + .orElseThrow(); + + assertThat(row.getProductName()).isEqualTo("5G 요금제"); + assertThat(row.getProductType()).isEqualTo("mobile"); + assertThat(row.getTags()).contains("영상OTT"); + assertThat(row.getLastEventId()).isEqualTo(1001L); + }); + } + + // ================================================================ + // Case 2: UPSERT — 동일 (member_id, product_id) 재조회 시 viewed_at 갱신 + // ================================================================ + @Test + @DisplayName("Case 2 | UPSERT — 동일 상품 재조회 시 viewed_at·last_event_id 갱신") + void case2_upsert_updatesViewedAtOnDuplicate() { + String first = buildPayload(2001L, 45L, 20L, + "2026-03-02T09:00:00.000Z", "LTE 요금제", "mobile", "[]"); + String second = buildPayload(2002L, 45L, 20L, + "2026-03-02T18:00:00.000Z", "LTE 요금제", "mobile", "[]"); + + testProducer.send("client-event-logs", first); + await().atMost(15, TimeUnit.SECONDS) + .until(() -> repository.findById(new ProductViewHistoryId(45L, 20L)).isPresent()); + + testProducer.send("client-event-logs", second); + + await().atMost(15, TimeUnit.SECONDS) + .untilAsserted(() -> { + ProductViewHistory row = repository + .findById(new ProductViewHistoryId(45L, 20L)) + .orElseThrow(); + + // 최신 이벤트로 덮어써진 상태 + assertThat(row.getLastEventId()).isEqualTo(2002L); + assertThat(row.getViewedAt()).isEqualTo(java.time.OffsetDateTime.parse("2026-03-02T18:00:00Z")); + }); + + // 레코드는 여전히 1건 + assertThat(repository.count()).isEqualTo(1L); + } + + // ================================================================ + // Case 3: TRIM — 유저당 31번째 상품 추가 시 가장 오래된 1개 삭제 + // ================================================================ + @Test + @DisplayName("Case 3 | TRIM — 31개 발행 후 유저당 최신 30개만 유지") + void case3_trim_keepsMax30RecordsPerMember() throws InterruptedException { + // 30개 먼저 삽입 (product_id 1~30, 시간 순차 증가) + for (int i = 1; i <= 30; i++) { + String ts = "2026-03-02T%02d:00:00.000Z".formatted(i % 24); + testProducer.send("client-event-logs", + buildPayload(3000L + i, 99L, (long) i, ts, + "상품" + i, "mobile", "[]")); + } + + // 30개 적재 완료 대기 + await().atMost(60, TimeUnit.SECONDS) + .until(() -> repository.count() == 30L); + + // 31번째 상품 추가 → TRIM 발동 + testProducer.send("client-event-logs", + buildPayload(3031L, 99L, 31L, + "2026-03-02T23:59:00.000Z", "신상품31", "mobile", "[]")); + + // TRIM 후 30개만 남아야 함 + await().atMost(15, TimeUnit.SECONDS) + .until(() -> repository.count() == 30L); + + assertThat(repository.count()).isEqualTo(30L); + } + + // ================================================================ + // Case 4: event_name 필터링 — click_product_detail 외 이벤트 무시 + // ================================================================ + @Test + @DisplayName("Case 4 | event_name 필터링 — page_view 이벤트는 DB 미적재") + void case4_eventNameFilter_nonTargetEventIgnored() { + // 1) 무시될 메시지 (page_view) — product_id 40 + String wrongEvent = """ + { + "event_id": 4001, + "timestamp": "2026-03-02T16:30:00.000Z", + "event": "page_view", + "event_name": "page_view", + "member_id": 45, + "event_properties": { + "product_id": 40, + "product_name": "무시상품", + "product_type": "mobile", + "tags": [] + } + } + """; + testProducer.send("client-event-logs", wrongEvent); + + // 2) 정상 처리될 메시지 (click_product_detail) — product_id 41 + String normalEvent = buildPayload(4002L, 45L, 41L, + "2026-03-02T16:31:00.000Z", "정상상품", "mobile", "[]"); + testProducer.send("client-event-logs", normalEvent); + + // 3) 정상 메시지가 처리될 때까지 대기 (Thread.sleep 대신 Awaitility) + await().atMost(15, TimeUnit.SECONDS) + .untilAsserted(() -> + assertThat(repository.findById(new ProductViewHistoryId(45L, 41L))).isPresent()); + + // 4) 무시된 메시지(product_id 40)는 DB에 없고, 정상 메시지(41)만 1건 존재 + assertThat(repository.findById(new ProductViewHistoryId(45L, 40L))).isEmpty(); + assertThat(repository.count()).isEqualTo(1L); + } + + // ================================================================ + // Case 5: event_id 타입 오류 (String UUID) → DLQ(error-logs) 전송, DB 미적재 + // ================================================================ + @Test + @DisplayName("Case 5 | event_id 타입 오류 — 역직렬화 실패 시 error-logs DLQ 전송") + void case5_invalidEventIdType_sentToDlq() { + String brokenPayload = """ + { + "event_id": "uuid-1234-5678", + "timestamp": "2026-03-02T17:10:00.000Z", + "event": "click", + "event_name": "click_product_detail", + "member_id": 45, + "event_properties": { + "product_id": 47, + "product_name": "타입오류", + "product_type": "mobile", + "tags": [] + } + } + """; + + // DLQ(error-logs) 전용 테스트 컨슈머 구성 + Map consumerProps = KafkaTestUtils.consumerProps( + "dlq-test-group", "true", embeddedKafkaBroker); + consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + + Consumer dlqConsumer = + new DefaultKafkaConsumerFactory(consumerProps).createConsumer(); + dlqConsumer.subscribe(Collections.singletonList("error-logs")); + + testProducer.send("client-event-logs", brokenPayload); + + // error-logs 토픽에 원본 메시지가 들어왔는지 확인 + await().atMost(15, TimeUnit.SECONDS) + .untilAsserted(() -> { + ConsumerRecords records = + dlqConsumer.poll(Duration.ofMillis(500)); + assertThat(records.count()).isGreaterThan(0); + assertThat(records.iterator().next().value()).contains("uuid-1234-5678"); + }); + + dlqConsumer.close(); + + // DB에는 적재되지 않음 + assertThat(repository.count()).isEqualTo(0L); + } + + // ================================================================ + // Case 6: 부하 시나리오 — 50건 처리 시간으로 TPS 검증 + // + // 목표: 일 3만명 × 10회 = 30만 건/일 ≈ 3.47 TPS (평균), 6.25 TPS (피크) + // 검증: 50건을 처리하는 실제 시간으로 가용 TPS 측정 + // ================================================================ + @Test + @DisplayName("Case 6 | 부하 시나리오 — 50건 처리 TPS ≥ 10 (운영 요구 6.25 TPS의 1.6배 여유)") + void case6_throughput_50MessagesProcessedWithSufficientTps() { + int messageCount = 50; + + long startMs = System.currentTimeMillis(); + + for (int i = 1; i <= messageCount; i++) { + testProducer.send("client-event-logs", + buildPayload(6000L + i, 77000L + i, (long) i, + "2026-03-02T16:30:00.000Z", "상품" + i, "mobile", "[]")); + } + + await().atMost(60, TimeUnit.SECONDS) + .until(() -> repository.count() == (long) messageCount); + + long elapsedMs = System.currentTimeMillis() - startMs; + double tps = messageCount / (elapsedMs / 1000.0); + + System.out.printf( + "[부하 테스트] %d건 처리 완료 | 소요시간: %dms | 처리량: %.2f TPS%n", + messageCount, elapsedMs, tps); + System.out.printf( + "[운영 요구사항] 평균 3.47 TPS / 피크 6.25 TPS → 현재 %.2f TPS (%.1f배 여유)%n", + tps, tps / 6.25); + + // MAX_POLL_RECORDS=1 + ACK_MODE=RECORD 설정으로도 피크 TPS의 1.6배 이상 처리 가능해야 함 + assertThat(tps).isGreaterThan(10.0); + } + + // ---------------------------------------------------------------- + // 헬퍼 — Kafka 발행용 JSON 페이로드 생성 + // ---------------------------------------------------------------- + private String buildPayload(long eventId, long memberId, long productId, + String timestamp, String productName, + String productType, String tagsJson) { + return """ + { + "event_id": %d, + "timestamp": "%s", + "event": "click", + "event_name": "click_product_detail", + "member_id": %d, + "event_properties": { + "product_id": %d, + "product_name": "%s", + "product_type": "%s", + "tags": %s + } + } + """.formatted(eventId, timestamp, memberId, productId, + productName, productType, tagsJson); + } +} diff --git a/src/test/java/com/holliverse/logserver/service/PostgresLogServiceTest.java b/src/test/java/com/holliverse/logserver/service/PostgresLogServiceTest.java new file mode 100644 index 0000000..ecb5971 --- /dev/null +++ b/src/test/java/com/holliverse/logserver/service/PostgresLogServiceTest.java @@ -0,0 +1,192 @@ +package com.holliverse.logserver.service; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.holliverse.logserver.dto.LogEvent; +import com.holliverse.logserver.repository.ProductViewHistoryRepository; +import java.time.OffsetDateTime; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +/** + * PostgresLogService 단위 테스트. + * Repository는 Mock, ObjectMapper는 실제 인스턴스를 사용하여 + * JSON 직렬화 로직까지 검증한다. + */ +@ExtendWith(MockitoExtension.class) +class PostgresLogServiceTest { + + @Mock + private ProductViewHistoryRepository repository; + + private PostgresLogService service; + + @BeforeEach + void setUp() { + // ObjectMapper는 실제 인스턴스 — tags 직렬화 결과까지 검증 + service = new PostgresLogService(repository, new ObjectMapper()); + } + + // ---------------------------------------------------------------- + // Case 1: 정상 신규 삽입 — upsert() 호출 인자 전체 검증 + // ---------------------------------------------------------------- + @Test + @DisplayName("Case 1 | 정상 신규 삽입 — upsert() 인자 정확히 전달") + void case1_normalInsert_callsUpsertWithCorrectArgs() throws Exception { + LogEvent event = buildEvent( + 1001L, 45L, 10L, + "2026-03-02T16:30:00.000Z", + List.of("영상OTT", "구독결제", "인기") + ); + + service.process(event); + + // upsert() 인자 캡처 후 검증 + ArgumentCaptor tagsCaptor = ArgumentCaptor.forClass(String.class); + ArgumentCaptor viewedAtCaptor = ArgumentCaptor.forClass(OffsetDateTime.class); + + verify(repository).upsert( + eq(45L), // memberId + eq(10L), // productId + eq("5G 요금제"), // productName + eq("mobile"), // productType + tagsCaptor.capture(), // tags JSON + viewedAtCaptor.capture(), // viewedAt + eq(1001L) // lastEventId + ); + + // tags가 JSON 배열 문자열로 직렬화되었는지 확인 + assertThat(tagsCaptor.getValue()).isEqualTo("[\"영상OTT\",\"구독결제\",\"인기\"]"); + + // timestamp가 OffsetDateTime으로 올바르게 파싱되었는지 확인 + assertThat(viewedAtCaptor.getValue()) + .isEqualTo(OffsetDateTime.parse("2026-03-02T16:30:00.000Z")); + } + + + // ---------------------------------------------------------------- + // Case 3: tags 빈 배열 → "[]" 로 직렬화 + // ---------------------------------------------------------------- + @Test + @DisplayName("Case 3 | tags 빈 배열 — JSON 빈 배열 \"[]\" 로 저장") + void case3_emptyTags_serializedAsEmptyJsonArray() throws Exception { + LogEvent event = buildEvent(3001L, 45L, 30L, + "2026-03-02T16:32:00.000Z", List.of()); + + service.process(event); + + ArgumentCaptor tagsCaptor = ArgumentCaptor.forClass(String.class); + verify(repository).upsert(anyLong(), anyLong(), anyString(), anyString(), + tagsCaptor.capture(), any(), anyLong()); + + assertThat(tagsCaptor.getValue()).isEqualTo("[]"); + } + + // ---------------------------------------------------------------- + // Case 4: event_properties에 tags 키 자체가 없는 경우 → "[]" + // ---------------------------------------------------------------- + @Test + @DisplayName("Case 4 | tags 키 누락 — null 대신 빈 배열 \"[]\" 로 처리") + void case4_tagsMissing_defaultsToEmptyJsonArray() throws Exception { + LogEvent event = buildEvent(4001L, 45L, 40L, + "2026-03-02T16:33:00.000Z", null); // null → 키 미포함으로 설정 + + service.process(event); + + ArgumentCaptor tagsCaptor = ArgumentCaptor.forClass(String.class); + verify(repository).upsert(anyLong(), anyLong(), anyString(), anyString(), + tagsCaptor.capture(), any(), anyLong()); + + assertThat(tagsCaptor.getValue()).isEqualTo("[]"); + } + + // ---------------------------------------------------------------- + // Case 5: product_id가 Integer(Jackson 기본 숫자 타입)로 역직렬화된 경우 → Long 변환 + // ---------------------------------------------------------------- + @Test + @DisplayName("Case 5 | product_id Integer → Long 변환 — ClassCastException 없음") + void case5_productIdAsInteger_convertedToLong() throws Exception { + LogEvent event = new LogEvent(); + event.setEventId(5001L); + event.setTimestamp("2026-03-02T16:34:00.000Z"); + event.setMemberId(45L); + + Map props = new HashMap<>(); + props.put("product_id", 50); // Jackson이 역직렬화할 때 Integer로 오는 케이스 + props.put("product_name", "인터넷"); + props.put("product_type", "internet"); + props.put("tags", List.of()); + event.setEventProperties(props); + + service.process(event); + + verify(repository).upsert(eq(45L), eq(50L), anyString(), anyString(), + anyString(), any(), anyLong()); + } + + // ---------------------------------------------------------------- + // Case 6: product_id null → IllegalArgumentException 발생, upsert 미호출 + // ---------------------------------------------------------------- + @Test + @DisplayName("Case 6 | product_id null — IllegalArgumentException 발생, upsert 미호출") + void case6_productIdNull_throwsIllegalArgumentException() { + LogEvent event = new LogEvent(); + event.setEventId(6001L); + event.setTimestamp("2026-03-02T16:35:00.000Z"); + event.setMemberId(45L); + + Map props = new HashMap<>(); + props.put("product_id", null); + props.put("product_name", "오류상품"); + props.put("product_type", "mobile"); + props.put("tags", List.of()); + event.setEventProperties(props); + + assertThatThrownBy(() -> service.process(event)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("product_id 변환 실패"); + + verify(repository, never()).upsert(anyLong(), anyLong(), anyString(), anyString(), + anyString(), any(), anyLong()); + } + + // ---------------------------------------------------------------- + // 헬퍼 메서드 + // ---------------------------------------------------------------- + + private LogEvent buildEvent(long eventId, long memberId, long productId, + String timestamp, List tags) { + LogEvent event = new LogEvent(); + event.setEventId(eventId); + event.setTimestamp(timestamp); + event.setEvent("click"); + event.setEventName("click_product_detail"); + event.setMemberId(memberId); + + Map props = new HashMap<>(); + props.put("product_id", productId); + props.put("product_name", "5G 요금제"); + props.put("product_type", "mobile"); + if (tags != null) { + props.put("tags", tags); + } + event.setEventProperties(props); + return event; + } +}