Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion .claude/settings.local.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,10 @@
"Bash(gh pr list:*)",
"Bash(ls:*)",
"Bash(wc:*)",
"Bash(brew list:*)"
"Bash(brew list:*)",
"WebFetch(domain:hey-mi.oopy.io)",
"Bash(mdfind:*)",
"Bash(done)"
]
}
}
30 changes: 30 additions & 0 deletions .editorconfig
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
root = true

[*.{kt,kts,java,js,jsx,ts,tsx}]
charset = utf-8
end_of_line = lf
indent_style = space
indent_size = 4
insert_final_newline = true
max_line_length = 120
trim_trailing_whitespace = true

[*.{js,jsx,ts,tsx,json,md,yml,yaml}]
indent_size = 2

[*.md]
max_line_length = off

[*.{kt,kts}]
ij_kotlin_name_count_to_use_star_import = 999
ij_kotlin_name_count_to_use_star_import_for_members = 999
ktlint_standard_trailing-comma-on-declaration-site = disabled
ktlint_standard_trailing-comma-on-call-site = disabled
ktlint_standard_function-expression-body = disabled
ktlint_standard_multiline-expression-wrapping = disabled
ktlint_standard_chain-method-continuation = disabled
ktlint_standard_function-signature = disabled
ktlint_standard_argument-list-wrapping = disabled
ktlint_standard_function-literal = disabled
ktlint_standard_if-else-wrapping = disabled
ktlint_standard_max-line-length = disabled
6 changes: 3 additions & 3 deletions behavior-consumer/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@ dependencies {
// Spring Boot
implementation("org.springframework.boot:spring-boot-starter")
implementation("org.springframework.boot:spring-boot-starter-actuator")
implementation("org.springframework.boot:spring-boot-starter-validation") // Bean Validation
implementation("org.springframework.boot:spring-boot-starter-webflux") // WebClient for Embedding Service
implementation("org.springframework.boot:spring-boot-starter-validation") // Bean Validation
implementation("org.springframework.boot:spring-boot-starter-webflux") // WebClient for Embedding Service
implementation("org.springframework.kafka:spring-kafka")

// Kotlin
Expand All @@ -44,7 +44,7 @@ dependencies {

// Logging
implementation("io.github.microutils:kotlin-logging-jvm:3.0.5")
implementation("net.logstash.logback:logstash-logback-encoder:8.0") // Phase 5: JSON 로깅
implementation("net.logstash.logback:logstash-logback-encoder:8.0") // Phase 5: JSON 로깅

// Micrometer for metrics
implementation("io.micrometer:micrometer-registry-prometheus")
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
package com.rep.consumer.client

import com.fasterxml.jackson.annotation.JsonProperty
import kotlinx.coroutines.reactor.awaitSingleOrNull
import mu.KotlinLogging
import org.springframework.stereotype.Component
Expand All @@ -19,11 +18,11 @@ private val log = KotlinLogging.logger {}
*/
@Component
class EmbeddingClient(
private val embeddingWebClient: WebClient
private val embeddingWebClient: WebClient,
) {
companion object {
const val QUERY_PREFIX = "query: " // 검색 쿼리용 (유저 취향)
const val PASSAGE_PREFIX = "passage: " // 문서용 (상품 정보)
const val QUERY_PREFIX = "query: " // 검색 쿼리용 (유저 취향)
const val PASSAGE_PREFIX = "passage: " // 문서용 (상품 정보)
}

/**
Expand All @@ -33,16 +32,21 @@ class EmbeddingClient(
* @param prefix e5 모델용 prefix (query: 또는 passage:)
* @return 벡터 목록
*/
suspend fun embed(texts: List<String>, prefix: String = QUERY_PREFIX): List<FloatArray>? {
suspend fun embed(
texts: List<String>,
prefix: String = QUERY_PREFIX,
): List<FloatArray>? {
if (texts.isEmpty()) return emptyList()

return try {
val response = embeddingWebClient.post()
.uri("/embed")
.bodyValue(EmbedRequest(texts = texts, prefix = prefix))
.retrieve()
.bodyToMono<EmbedResponse>()
.awaitSingleOrNull()
val response =
embeddingWebClient
.post()
.uri("/embed")
.bodyValue(EmbedRequest(texts = texts, prefix = prefix))
.retrieve()
.bodyToMono<EmbedResponse>()
.awaitSingleOrNull()

response?.embeddings?.map { it.toFloatArray() }
} catch (e: Exception) {
Expand All @@ -54,41 +58,43 @@ class EmbeddingClient(
/**
* 단일 텍스트를 벡터로 변환합니다.
*/
suspend fun embedSingle(text: String, prefix: String = QUERY_PREFIX): FloatArray? {
return embed(listOf(text), prefix)?.firstOrNull()
}
suspend fun embedSingle(
text: String,
prefix: String = QUERY_PREFIX,
): FloatArray? = embed(listOf(text), prefix)?.firstOrNull()

/**
* 헬스체크
*/
suspend fun healthCheck(): Boolean {
return try {
val response = embeddingWebClient.get()
.uri("/health")
.retrieve()
.bodyToMono<HealthResponse>()
.awaitSingleOrNull()
suspend fun healthCheck(): Boolean =
try {
val response =
embeddingWebClient
.get()
.uri("/health")
.retrieve()
.bodyToMono<HealthResponse>()
.awaitSingleOrNull()

response?.status == "ok"
} catch (e: Exception) {
log.warn(e) { "Embedding service health check failed" }
false
}
}
}

data class EmbedRequest(
val texts: List<String>,
val prefix: String = "query: "
val prefix: String = "query: ",
)

data class EmbedResponse(
val embeddings: List<List<Float>>,
val dims: Int = 768
val dims: Int = 768,
)

data class HealthResponse(
val status: String,
val model: String,
val dims: Int
val dims: Int,
)
Original file line number Diff line number Diff line change
Expand Up @@ -10,30 +10,22 @@ import org.springframework.validation.annotation.Validated
data class ConsumerProperties(
@field:NotBlank(message = "topic must not be blank")
val topic: String = "user.action.v1",

@field:NotBlank(message = "dlqTopic must not be blank")
val dlqTopic: String = "user.action.v1.dlq",

@field:Positive(message = "bulkSize must be positive")
val bulkSize: Int = 500,

@field:Positive(message = "concurrency must be positive")
val concurrency: Int = 3,

@field:Positive(message = "maxRetries must be positive")
val maxRetries: Int = 3,

@field:Positive(message = "retryDelayMs must be positive")
val retryDelayMs: Long = 1000,

// DLQ 파일 백업 설정
@field:Positive(message = "dlqFileMaxSizeBytes must be positive")
val dlqFileMaxSizeBytes: Long = 10 * 1024 * 1024L, // 10MB

val dlqFileMaxSizeBytes: Long = 10 * 1024 * 1024L, // 10MB
@field:NotBlank(message = "dlqLogsDir must not be blank")
val dlqLogsDir: String = "logs",

// 벡터 설정 (multilingual-e5-base)
@field:Positive(message = "vectorDimensions must be positive")
val vectorDimensions: Int = 768
val vectorDimensions: Int = 768,
)
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import kotlinx.coroutines.CloseableCoroutineDispatcher
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.asCoroutineDispatcher
import mu.KotlinLogging
import org.springframework.beans.factory.annotation.Qualifier
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
import java.util.concurrent.Executors
Expand All @@ -23,7 +22,6 @@ private val log = KotlinLogging.logger {}
@Configuration
@OptIn(ExperimentalCoroutinesApi::class)
class DispatcherConfig {

private var dispatcher: CloseableCoroutineDispatcher? = null

/**
Expand All @@ -32,11 +30,11 @@ class DispatcherConfig {
* Java 25 Virtual Threads를 활용하여 blocking I/O 호출 시에도
* 시스템 처리량을 유지합니다.
*/
@Bean
@Qualifier("virtualThreadDispatcher")
@Bean("virtualThreadDispatcher")
fun virtualThreadDispatcher(): CloseableCoroutineDispatcher {
log.info { "Creating Virtual Thread Coroutine Dispatcher" }
return Executors.newVirtualThreadPerTaskExecutor()
return Executors
.newVirtualThreadPerTaskExecutor()
.asCoroutineDispatcher()
.also { dispatcher = it }
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ private val log = KotlinLogging.logger {}

@Configuration
class ElasticsearchConfig {

@Value("\${elasticsearch.host}")
private lateinit var host: String

Expand All @@ -29,11 +28,12 @@ class ElasticsearchConfig {
private var transport: RestClientTransport? = null

@Bean
fun restClient(): RestClient {
return RestClient.builder(
HttpHost(host, port, scheme)
).build().also { restClient = it }
}
fun restClient(): RestClient =
RestClient
.builder(
HttpHost(host, port, scheme),
).build()
.also { restClient = it }

@Bean
fun elasticsearchClient(restClient: RestClient): ElasticsearchClient {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ import org.springframework.validation.annotation.Validated
data class EmbeddingProperties(
@field:NotBlank(message = "url must not be blank")
val url: String = "http://localhost:8000",

@field:Positive(message = "timeoutMs must be positive")
val timeoutMs: Long = 5000
val timeoutMs: Long = 5000,
)
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,8 @@ import org.springframework.kafka.listener.ContainerProperties

@Configuration
class KafkaConsumerConfig(
private val consumerProperties: ConsumerProperties
private val consumerProperties: ConsumerProperties,
) {

@Value("\${spring.kafka.bootstrap-servers}")
private lateinit var bootstrapServers: String

Expand All @@ -26,31 +25,29 @@ class KafkaConsumerConfig(

@Bean
fun consumerFactory(): ConsumerFactory<String, UserActionEvent> {
val props = mapOf(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG to bootstrapServers,
// Group ID는 application.yml 또는 @KafkaListener에서 설정
ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG to StringDeserializer::class.java,
ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG to KafkaAvroDeserializer::class.java,

// 성능 튜닝 - Phase 2 문서 기준
ConsumerConfig.MAX_POLL_RECORDS_CONFIG to consumerProperties.bulkSize, // Bulk 크기와 일치
ConsumerConfig.FETCH_MIN_BYTES_CONFIG to 1024, // 최소 fetch 크기
ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG to 500, // 최대 대기 시간

// 안정성 - 수동 커밋으로 메시지 유실 방지
ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG to false,
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG to "earliest",

// Avro Schema Registry
KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG to schemaRegistryUrl,
KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG to true
)
val props =
mapOf(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG to bootstrapServers,
// Group ID는 application.yml 또는 @KafkaListener에서 설정
ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG to StringDeserializer::class.java,
ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG to KafkaAvroDeserializer::class.java,
// 성능 튜닝 - Phase 2 문서 기준
ConsumerConfig.MAX_POLL_RECORDS_CONFIG to consumerProperties.bulkSize, // Bulk 크기와 일치
ConsumerConfig.FETCH_MIN_BYTES_CONFIG to 1024, // 최소 fetch 크기
ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG to 500, // 최대 대기 시간
// 안정성 - 수동 커밋으로 메시지 유실 방지
ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG to false,
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG to "earliest",
// Avro Schema Registry
KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG to schemaRegistryUrl,
KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG to true,
)
return DefaultKafkaConsumerFactory(props)
}

@Bean
fun kafkaListenerContainerFactory(): ConcurrentKafkaListenerContainerFactory<String, UserActionEvent> {
return ConcurrentKafkaListenerContainerFactory<String, UserActionEvent>().apply {
fun kafkaListenerContainerFactory(): ConcurrentKafkaListenerContainerFactory<String, UserActionEvent> =
ConcurrentKafkaListenerContainerFactory<String, UserActionEvent>().apply {
consumerFactory = consumerFactory()
// 수동 커밋 - 배치 처리 완료 즉시 커밋
containerProperties.ackMode = ContainerProperties.AckMode.MANUAL_IMMEDIATE
Expand All @@ -62,5 +59,4 @@ class KafkaConsumerConfig(
// 분산 트레이싱 Observation 활성화
containerProperties.isObservationEnabled = true
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@ import org.springframework.kafka.core.ProducerFactory
*/
@Configuration
class KafkaProducerConfig {

@Value("\${spring.kafka.bootstrap-servers}")
private lateinit var bootstrapServers: String

Expand All @@ -26,26 +25,24 @@ class KafkaProducerConfig {

@Bean
fun producerFactory(): ProducerFactory<String, UserActionEvent> {
val configProps = mapOf(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG to bootstrapServers,
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG to StringSerializer::class.java,
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG to KafkaAvroSerializer::class.java,

// Reliability settings
ProducerConfig.ACKS_CONFIG to "all",
ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG to true,
ProducerConfig.RETRIES_CONFIG to 3,

// Schema Registry
KafkaAvroSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG to schemaRegistryUrl
)
val configProps =
mapOf(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG to bootstrapServers,
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG to StringSerializer::class.java,
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG to KafkaAvroSerializer::class.java,
// Reliability settings
ProducerConfig.ACKS_CONFIG to "all",
ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG to true,
ProducerConfig.RETRIES_CONFIG to 3,
// Schema Registry
KafkaAvroSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG to schemaRegistryUrl,
)
return DefaultKafkaProducerFactory(configProps)
}

@Bean
fun kafkaTemplate(): KafkaTemplate<String, UserActionEvent> {
return KafkaTemplate(producerFactory()).apply {
fun kafkaTemplate(): KafkaTemplate<String, UserActionEvent> =
KafkaTemplate(producerFactory()).apply {
setObservationEnabled(true)
}
}
}
Loading