From d052c060d2667fe916d6eee3f1036d5ff1d93226 Mon Sep 17 00:00:00 2001 From: tkv00 Date: Fri, 13 Mar 2026 00:14:26 +0900 Subject: [PATCH 1/2] =?UTF-8?q?[HSC-290]=20feat:=20=EB=A1=9C=EA=B7=B8=20?= =?UTF-8?q?=EC=84=9C=EB=B2=84=20=EC=A0=81=EC=9E=AC=20=EA=B5=AC=EC=A1=B0=20?= =?UTF-8?q?=EC=A0=84=ED=99=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- build.gradle | 19 +++- .../config/KafkaDlqProducerConfig.java | 27 +++++- .../config/KafkaPropertiesConfig.java | 5 +- .../config/KafkaSpeedConsumerConfig.java | 54 +++++++++--- .../config/properties/KafkaAppProperties.java | 11 +++ .../consumer/SpeedLayerConsumer.java | 18 ++-- .../holliverse/logserver/dto/LogEvent.java | 12 +-- .../logserver/entity/ProductViewHistory.java | 38 ++++++++ .../entity/ProductViewHistoryId.java | 23 +++++ .../ProductViewHistoryRepository.java | 61 +++++++++++++ .../logserver/service/PostgresLogService.java | 87 +++++++++++++++++++ .../logserver/service/RedisLogService.java | 81 ----------------- src/main/resources/application.yaml | 34 +++++--- .../logserver/LogServerApplicationTests.java | 9 +- 14 files changed, 351 insertions(+), 128 deletions(-) create mode 100644 src/main/java/com/holliverse/logserver/entity/ProductViewHistory.java create mode 100644 src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java create mode 100644 src/main/java/com/holliverse/logserver/repository/ProductViewHistoryRepository.java create mode 100644 src/main/java/com/holliverse/logserver/service/PostgresLogService.java delete mode 100644 src/main/java/com/holliverse/logserver/service/RedisLogService.java diff --git a/build.gradle b/build.gradle index 29994c5..b9d9eb8 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,33 @@ 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' + implementation 'software.amazon.msk:aws-msk-iam-auth:2.3.5' + 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 'com.h2database:h2' 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/src/main/java/com/holliverse/logserver/config/KafkaDlqProducerConfig.java b/src/main/java/com/holliverse/logserver/config/KafkaDlqProducerConfig.java index e8e6f8d..3a1ba9e 100644 --- a/src/main/java/com/holliverse/logserver/config/KafkaDlqProducerConfig.java +++ b/src/main/java/com/holliverse/logserver/config/KafkaDlqProducerConfig.java @@ -3,13 +3,16 @@ import com.holliverse.logserver.config.properties.KafkaAppProperties; import java.util.HashMap; import java.util.Map; +import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.config.SaslConfigs; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; +import org.springframework.util.StringUtils; @Configuration public class KafkaDlqProducerConfig { @@ -23,12 +26,34 @@ public KafkaDlqProducerConfig(KafkaAppProperties kafkaAppProperties) { @Bean public ProducerFactory dlqProducerFactory() { Map props = new HashMap<>(); + // 브로커 주소 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaAppProperties.getBootstrapServers()); + // 키 직렬화 props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + // 값 직렬화 props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); + // 보안 프로토콜 + props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, kafkaAppProperties.getSecurity().getProtocol()); - // 에러 로그는 전송 성공이 중요하므로 기본값 acks=all, retries=3 + if (StringUtils.hasText(kafkaAppProperties.getSecurity().getSaslMechanism())) { + // SASL 메커니즘 + props.put(SaslConfigs.SASL_MECHANISM, kafkaAppProperties.getSecurity().getSaslMechanism()); + } + if (StringUtils.hasText(kafkaAppProperties.getSecurity().getSaslJaasConfig())) { + // JAAS 설정 + props.put(SaslConfigs.SASL_JAAS_CONFIG, kafkaAppProperties.getSecurity().getSaslJaasConfig()); + } + if (StringUtils.hasText(kafkaAppProperties.getSecurity().getSaslCallbackHandlerClass())) { + // 콜백 핸들러 + props.put( + SaslConfigs.SASL_CLIENT_CALLBACK_HANDLER_CLASS, + kafkaAppProperties.getSecurity().getSaslCallbackHandlerClass() + ); + } + + // DLQ ack 설정 props.put(ProducerConfig.ACKS_CONFIG, kafkaAppProperties.getProducer().getDlqAcks()); + // DLQ retry 설정 props.put(ProducerConfig.RETRIES_CONFIG, kafkaAppProperties.getProducer().getDlqRetries()); return new DefaultKafkaProducerFactory<>(props); diff --git a/src/main/java/com/holliverse/logserver/config/KafkaPropertiesConfig.java b/src/main/java/com/holliverse/logserver/config/KafkaPropertiesConfig.java index a620e21..6a73563 100644 --- a/src/main/java/com/holliverse/logserver/config/KafkaPropertiesConfig.java +++ b/src/main/java/com/holliverse/logserver/config/KafkaPropertiesConfig.java @@ -9,9 +9,7 @@ public class KafkaPropertiesConfig { /** - * app.kafka.* 설정을 KafkaAppProperties에 바인딩하는 전용 Bean. - * Bean 이름을 명시적으로 'kafkaAppProperties'로 고정해서 - * SpEL(@kafkaAppProperties)에서 안정적으로 참조할 수 있게 한다. + * app.kafka 설정 바인딩 빈. */ @Bean("kafkaAppProperties") @ConfigurationProperties(prefix = "app.kafka") @@ -19,4 +17,3 @@ public KafkaAppProperties kafkaAppProperties() { return new KafkaAppProperties(); } } - diff --git a/src/main/java/com/holliverse/logserver/config/KafkaSpeedConsumerConfig.java b/src/main/java/com/holliverse/logserver/config/KafkaSpeedConsumerConfig.java index 22a0f27..647d7b6 100644 --- a/src/main/java/com/holliverse/logserver/config/KafkaSpeedConsumerConfig.java +++ b/src/main/java/com/holliverse/logserver/config/KafkaSpeedConsumerConfig.java @@ -3,12 +3,12 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.holliverse.logserver.config.properties.KafkaAppProperties; - -import lombok.extern.slf4j.Slf4j; - import java.util.HashMap; import java.util.Map; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.common.config.SaslConfigs; import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -16,6 +16,7 @@ import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; +import org.springframework.util.StringUtils; @Slf4j @Configuration @@ -32,41 +33,70 @@ public KafkaSpeedConsumerConfig(KafkaAppProperties kafkaAppProperties, ObjectMap this.objectMapper = objectMapper; } - /* consumer factory bean */ - // MAX_POLL_RECORDS_CONFIG, 1 (중요): 폴링할때 한건만 가져옴 - // ENABLE_AUTO_COMMIT_CONFIG, false: 자동 커밋 비활성화 → AckMode.RECORD로 처리 성공 후 커밋 - // AUTO_OFFSET_RESET_CONFIG, earliest: 컨슈머 그룹 최초 생성 시 맨 처음 데이터부터 읽어 유실 방지 + /** + * speed consumer factory 빈. + */ @Bean public ConsumerFactory speedConsumerFactory() { Map props = new HashMap<>(); + // 브로커 주소 props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaAppProperties.getBootstrapServers()); + // 그룹 아이디 props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaAppProperties.getGroups().getSpeed()); + // 키 역직렬화 props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + // 값 역직렬화 props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + // 자동 커밋 비활성 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); + // poll 크기 props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, kafkaAppProperties.getListener().getMaxPollRecords()); + // 초기 오프셋 정책 props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + // 보안 프로토콜 + props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, kafkaAppProperties.getSecurity().getProtocol()); + + if (StringUtils.hasText(kafkaAppProperties.getSecurity().getSaslMechanism())) { + // SASL 메커니즘 + props.put(SaslConfigs.SASL_MECHANISM, kafkaAppProperties.getSecurity().getSaslMechanism()); + } + if (StringUtils.hasText(kafkaAppProperties.getSecurity().getSaslJaasConfig())) { + // JAAS 설정 + props.put(SaslConfigs.SASL_JAAS_CONFIG, kafkaAppProperties.getSecurity().getSaslJaasConfig()); + } + if (StringUtils.hasText(kafkaAppProperties.getSecurity().getSaslCallbackHandlerClass())) { + // 콜백 핸들러 + props.put( + SaslConfigs.SASL_CLIENT_CALLBACK_HANDLER_CLASS, + kafkaAppProperties.getSecurity().getSaslCallbackHandlerClass() + ); + } + return new DefaultKafkaConsumerFactory<>(props); } - /* container factory bean — 실제 Kafka Listener 엔진 */ - // RecordFilterStrategy: event_name != "click_product_detail" 이면 메시지를 @KafkaListener 전에 폐기 - // → true 반환 = 폐기(discard), false 반환 = 리스너로 전달 - // AckMode.RECORD: 1건 처리 완료 즉시 커밋 → 재처리 범위를 최소화 + /** + * speed listener container 빈. + */ @Bean public ConcurrentKafkaListenerContainerFactory speedLayerContainerFactory() { ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); + // consumer factory 연결 factory.setConsumerFactory(speedConsumerFactory()); + // ack 모드 factory.getContainerProperties().setAckMode( kafkaAppProperties.getListener().getAckMode()); + // 자동 시작 여부 + factory.setAutoStartup(kafkaAppProperties.getListener().isAutoStartup()); + // 이벤트 필터 factory.setRecordFilterStrategy(record -> { try { JsonNode node = objectMapper.readTree(record.value()); return !FILTER_EVENT_NAME.equals(node.path("event_name").asText("")); } catch (Exception e) { - // 파싱 불가 = 깨진 JSON → 폐기 (Consumer의 DLQ와 역할 분리) + // 파싱 실패 폐기 log.warn("레코드 필터링 중 value 매칭 실패로 폐기합니다. value={}", record.value(), e); return true; } diff --git a/src/main/java/com/holliverse/logserver/config/properties/KafkaAppProperties.java b/src/main/java/com/holliverse/logserver/config/properties/KafkaAppProperties.java index d35b8ec..74cff05 100644 --- a/src/main/java/com/holliverse/logserver/config/properties/KafkaAppProperties.java +++ b/src/main/java/com/holliverse/logserver/config/properties/KafkaAppProperties.java @@ -13,6 +13,7 @@ public class KafkaAppProperties { private Groups groups = new Groups(); private Listener listener = new Listener(); private Producer producer = new Producer(); + private Security security = new Security(); @Getter @Setter @@ -32,6 +33,7 @@ public static class Groups { public static class Listener { private int maxPollRecords = 1; private ContainerProperties.AckMode ackMode = ContainerProperties.AckMode.RECORD; + private boolean autoStartup = true; } @Getter @@ -40,4 +42,13 @@ public static class Producer { private String dlqAcks = "all"; private int dlqRetries = 3; } + + @Getter + @Setter + public static class Security { + private String protocol = "PLAINTEXT"; + private String saslMechanism; + private String saslJaasConfig; + private String saslCallbackHandlerClass; + } } diff --git a/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java b/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java index 222bab0..a42a67b 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,13 +17,13 @@ public class SpeedLayerConsumer { private final ObjectMapper objectMapper; - private final RedisLogService redisLogService; + private final PostgresLogService postgresLogService; private final KafkaTemplate dlqKafkaTemplate; private final KafkaAppProperties kafkaAppProperties; - // topics, groupId를 SpEL로 application.yaml의 app.kafka 값에서 주입 - // RecordFilterStrategy가 설정 레벨에서 click_product_detail만 통과시키므로 - // 이 메서드에 도달하는 메시지는 이미 필터링된 상태임 + /** + * 클릭 로그 소비 메서드. + */ @KafkaListener( topics = "#{@kafkaAppProperties.topics.clientEvents}", groupId = "#{@kafkaAppProperties.groups.speed}", @@ -31,16 +31,18 @@ public class SpeedLayerConsumer { ) public void consume(ConsumerRecord record) { try { + // 원본 로그 역직렬화 LogEvent event = objectMapper.readValue(record.value(), LogEvent.class); - redisLogService.process(event); + // DB 반영 + postgresLogService.process(event); } catch (Exception e) { - // 역직렬화 or Redis 처리 실패 → DLQ 토픽으로 원본 페이로드 전송 - // 예외를 re-throw하지 않아 다음 메시지 처리가 중단되지 않음 (무한 루프 방지) + // DLQ 전송 dlqKafkaTemplate.send( kafkaAppProperties.getTopics().getError(), record.key(), record.value() ); + // 에러 로그 log.error("[SpeedLayer DLQ] topic={}, partition={}, offset={}, err={}", record.topic(), record.partition(), record.offset(), e.getMessage()); } 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..d19aa3f --- /dev/null +++ b/src/main/java/com/holliverse/logserver/entity/ProductViewHistory.java @@ -0,0 +1,38 @@ +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; + + @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..5411d2e --- /dev/null +++ b/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java @@ -0,0 +1,23 @@ +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 +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..3471fc0 --- /dev/null +++ b/src/main/java/com/holliverse/logserver/repository/ProductViewHistoryRepository.java @@ -0,0 +1,61 @@ +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 { + + /** + * 최근 본 상품 upsert 쿼리. + */ + @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 + ); + + /** + * 유저별 오래된 기록 정리 쿼리. + */ + @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..6531422 --- /dev/null +++ b/src/main/java/com/holliverse/logserver/service/PostgresLogService.java @@ -0,0 +1,87 @@ +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; + + /** + * 클릭 로그 DB 반영 메서드. + */ + @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(); + + // 최근 본 상품 upsert + repository.upsert(memberId, productId, productName, productType, tags, viewedAt, lastEventId); + // 오래된 기록 정리 + repository.trimOldRecords(memberId, MAX_RECENT_VIEWS); + + // 처리 로그 + log.debug("[PostgresLog] UPSERT 완료 memberId={} productId={}", memberId, productId); + } + + /** + * 클릭 상세 속성 DTO. + */ + @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 da9e1c1..36984ee 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -1,23 +1,35 @@ 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: bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092} topics: - client-events: client-event-logs - error: error-logs + client-events: ${KAFKA_TOPIC_CLIENT_EVENTS:client-event-logs} + error: ${KAFKA_TOPIC_ERROR:error-logs} groups: - speed: speed-layer-group + speed: ${KAFKA_GROUP_SPEED:speed-layer-group} listener: - max-poll-records: 1 + max-poll-records: ${KAFKA_MAX_POLL_RECORDS:1} ack-mode: RECORD + auto-startup: ${KAFKA_LISTENER_AUTO_STARTUP:true} producer: - dlq-acks: all - dlq-retries: 3 + dlq-acks: ${KAFKA_DLQ_ACKS:all} + dlq-retries: ${KAFKA_DLQ_RETRIES:3} + security: + protocol: ${KAFKA_SECURITY_PROTOCOL:PLAINTEXT} + sasl-mechanism: ${KAFKA_SASL_MECHANISM:} + sasl-jaas-config: ${KAFKA_SASL_JAAS_CONFIG:} + sasl-callback-handler-class: ${KAFKA_SASL_CALLBACK_HANDLER_CLASS:} diff --git a/src/test/java/com/holliverse/logserver/LogServerApplicationTests.java b/src/test/java/com/holliverse/logserver/LogServerApplicationTests.java index c0bcf6c..7dfd2ba 100644 --- a/src/test/java/com/holliverse/logserver/LogServerApplicationTests.java +++ b/src/test/java/com/holliverse/logserver/LogServerApplicationTests.java @@ -3,7 +3,14 @@ import org.junit.jupiter.api.Test; import org.springframework.boot.test.context.SpringBootTest; -@SpringBootTest +@SpringBootTest(properties = { + "spring.datasource.url=jdbc:h2:mem:logserver;MODE=PostgreSQL;DB_CLOSE_DELAY=-1", + "spring.datasource.username=sa", + "spring.datasource.password=", + "spring.datasource.driver-class-name=org.h2.Driver", + "spring.jpa.hibernate.ddl-auto=none", + "app.kafka.listener.auto-startup=false" +}) class LogServerApplicationTests { @Test From 2bcbcdc5315f7a84f4193260f74b5112aa6b569e Mon Sep 17 00:00:00 2001 From: tkv00 Date: Fri, 13 Mar 2026 00:14:42 +0900 Subject: [PATCH 2/2] =?UTF-8?q?[HSC-290]=20feat:=20=EB=A1=9C=EA=B7=B8=20?= =?UTF-8?q?=EC=84=9C=EB=B2=84=20=EC=A4=91=EC=95=99=20=EB=B0=B0=ED=8F=AC=20?= =?UTF-8?q?=EC=97=B0=EA=B2=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .dockerignore | 6 +++++ .github/workflows/main-deploy-dispatch.yml | 31 ++++++++++++++++++++++ Dockerfile | 21 +++++++++++++++ 3 files changed, 58 insertions(+) create mode 100644 .dockerignore create mode 100644 .github/workflows/main-deploy-dispatch.yml create mode 100644 Dockerfile diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..2ec9a8d --- /dev/null +++ b/.dockerignore @@ -0,0 +1,6 @@ +.git +.github +.gradle +build +node_modules +coverage diff --git a/.github/workflows/main-deploy-dispatch.yml b/.github/workflows/main-deploy-dispatch.yml new file mode 100644 index 0000000..e0c0954 --- /dev/null +++ b/.github/workflows/main-deploy-dispatch.yml @@ -0,0 +1,31 @@ +name: main-deploy-dispatch + +on: + push: + branches: [main] + +permissions: + contents: read + +concurrency: + group: log-server-main-deploy-dispatch + cancel-in-progress: false + +jobs: + dispatch: + runs-on: ubuntu-latest + steps: + - name: Request central deployment + uses: actions/github-script@v7 + with: + github-token: ${{ secrets.CENTRAL_REPO_TOKEN }} + script: | + await github.rest.repos.createDispatchEvent({ + owner: "one-year-gap", + repo: "infra", + event_type: "log-server-main-deploy-request", + client_payload: { + source_repo: `${context.repo.owner}/${context.repo.repo}`, + source_sha: context.sha + } + }); diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..3e98591 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,21 @@ +FROM gradle:8.14.3-jdk17 AS builder + +WORKDIR /app + +COPY gradlew gradlew +COPY gradle gradle +COPY build.gradle settings.gradle ./ +COPY src src + +RUN chmod +x gradlew +RUN ./gradlew bootJar --no-daemon + +FROM eclipse-temurin:17-jre + +WORKDIR /app + +COPY --from=builder /app/build/libs/*.jar app.jar + +EXPOSE 8080 + +ENTRYPOINT ["java", "-jar", "/app/app.jar"]