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"] diff --git a/build.gradle b/build.gradle index 48471f1..b9d9eb8 100644 --- a/build.gradle +++ b/build.gradle @@ -30,6 +30,7 @@ dependencies { implementation 'org.springframework.kafka:spring-kafka' 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' @@ -39,6 +40,7 @@ dependencies { 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' } 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 74d0eb2..107c493 100644 --- a/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java +++ b/src/main/java/com/holliverse/logserver/consumer/SpeedLayerConsumer.java @@ -21,9 +21,9 @@ public class SpeedLayerConsumer { 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,6 +31,7 @@ public class SpeedLayerConsumer { ) public void consume(ConsumerRecord record) { try { + // 원본 로그 역직렬화 LogEvent event = objectMapper.readValue(record.value(), LogEvent.class); postgresLogService.process(event); } catch (Exception e) { @@ -41,6 +42,7 @@ public void consume(ConsumerRecord record) { 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/entity/ProductViewHistoryId.java b/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java index e7b57e8..5411d2e 100644 --- a/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java +++ b/src/main/java/com/holliverse/logserver/entity/ProductViewHistoryId.java @@ -13,7 +13,6 @@ @NoArgsConstructor @AllArgsConstructor @EqualsAndHashCode -// (member_id, product_id) 복합 PK @Embeddable public class ProductViewHistoryId implements Serializable { @Column(name = "member_id", nullable = false) diff --git a/src/main/java/com/holliverse/logserver/service/PostgresLogService.java b/src/main/java/com/holliverse/logserver/service/PostgresLogService.java index b48b23d..f144c0f 100644 --- a/src/main/java/com/holliverse/logserver/service/PostgresLogService.java +++ b/src/main/java/com/holliverse/logserver/service/PostgresLogService.java @@ -51,7 +51,6 @@ public void process(LogEvent event) throws JsonProcessingException { repository.upsert(memberId, productId, productName, productType, tags, viewedAt, lastEventId); repository.trimOldRecords(memberId, MAX_RECENT_VIEWS); - log.debug("[PostgresLog] UPSERT 완료 memberId={} productId={}", memberId, productId); } diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index 65f5f50..94ae370 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -25,7 +25,12 @@ app: listener: max-poll-records: ${KAFKA_MAX_POLL_RECORDS:1} ack-mode: RECORD + auto-startup: ${KAFKA_LISTENER_AUTO_STARTUP:true} producer: 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