Skip to content
Merged
17 changes: 16 additions & 1 deletion build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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
}
}
23 changes: 1 addition & 22 deletions docker-compose.local.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -17,7 +17,7 @@
public class SpeedLayerConsumer {

private final ObjectMapper objectMapper;
private final RedisLogService redisLogService;
private final PostgresLogService postgresLogService;
private final KafkaTemplate<String, String> dlqKafkaTemplate;
private final KafkaAppProperties kafkaAppProperties;

Expand All @@ -32,9 +32,9 @@ public class SpeedLayerConsumer {
public void consume(ConsumerRecord<String, String> 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(),
Expand Down
12 changes: 3 additions & 9 deletions src/main/java/com/holliverse/logserver/dto/LogEvent.java
Original file line number Diff line number Diff line change
Expand Up @@ -8,23 +8,17 @@
public class LogEvent {

@JsonProperty("event_id")
private String eventId;
private Long eventId;

private String timestamp;
private String event;

@JsonProperty("event_name")
private String eventName;

@JsonProperty("member_properties")
private MemberProperties memberProperties;
@JsonProperty("member_id")
private Long memberId;

@JsonProperty("event_properties")
private Map<String, Object> eventProperties;

@Data
public static class MemberProperties {
@JsonProperty("member_id")
private String memberId;
}
}
Original file line number Diff line number Diff line change
@@ -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<String>을 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;
}
Original file line number Diff line number Diff line change
@@ -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
Comment thread
rettooo marked this conversation as resolved.
public class ProductViewHistoryId implements Serializable {

@Column(name = "member_id", nullable = false)
private Long memberId;

@Column(name = "product_id", nullable = false)
private Long productId;
}
Original file line number Diff line number Diff line change
@@ -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<ProductViewHistory, ProductViewHistoryId> {

/**
* 복합 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
);
}
Original file line number Diff line number Diff line change
@@ -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<String> 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<String> tags;
}
}
Loading