Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -412,7 +412,7 @@
"title": "Circuit Breaker 상태",
"description": "closed=정상, half_open=복구시도, open=차단됨. AI API 등 외부 의존성 보호.",
"type": "stat",
"targets": [ { "expr": "resilience4j_circuitbreaker_state == 1", "legendFormat": "{{name}} ({{state}})", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ]
"targets": [ { "expr": "sum by (name) (resilience4j_circuitbreaker_state{state=\"closed\"} * 0) + sum by (name) (resilience4j_circuitbreaker_state{state=\"half_open\"} * 1) + sum by (name) (resilience4j_circuitbreaker_state{state=\"open\"} * 2)", "legendFormat": "{{name}}", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

forced_opendisabled 상태 누락으로 인한 오탐지 위험

현재 PromQL 식은 closed (0), half_open (1), open (2) 상태만 고려하여 합산하고 있습니다.

하지만 Resilience4j에는 forced_opendisabled 상태도 존재합니다. 만약 운영 중에 긴급 대응이나 테스트를 위해 서킷 브레이커를 forced_open 상태로 강제 전환할 경우, closed/half_open/open 메트릭은 모두 0이 되고 forced_open 메트릭만 1이 됩니다.

이 경우 현재 식의 결과는 0이 되어 대시보드에는 정상 상태인 CLOSED (초록색)로 표시되는 오탐지(Misleading)가 발생할 수 있습니다.

따라서 forced_open 상태도 open과 동일하게 2로 매핑할 수 있도록 정규식(=~"open|forced_open")을 사용하는 것을 권장합니다.

Suggested change
"targets": [ { "expr": "sum by (name) (resilience4j_circuitbreaker_state{state=\"closed\"} * 0) + sum by (name) (resilience4j_circuitbreaker_state{state=\"half_open\"} * 1) + sum by (name) (resilience4j_circuitbreaker_state{state=\"open\"} * 2)", "legendFormat": "{{name}}", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ]
"targets": [ { "expr": "sum by (name) (resilience4j_circuitbreaker_state{state=\"closed\"} * 0) + sum by (name) (resilience4j_circuitbreaker_state{state=\"half_open\"} * 1) + sum by (name) (resilience4j_circuitbreaker_state{state=~\"open|forced_open\"} * 2)", "legendFormat": "{{name}}", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ]

},
{
"datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" },
Expand Down
19 changes: 15 additions & 4 deletions src/main/java/com/catchtable/global/config/KafkaConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@
import org.springframework.kafka.core.*;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.support.serializer.ErrorHandlingDeserializer;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.util.backoff.FixedBackOff;

import java.util.HashMap;
Expand Down Expand Up @@ -70,10 +72,19 @@ public ConsumerFactory<String, Object> consumerFactory() {
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
config.put(ConsumerConfig.GROUP_ID_CONFIG, "catchtable-notification-group");
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.springframework.kafka.support.serializer.JsonDeserializer");
config.put("spring.json.trusted.packages", "com.catchtable.notification.event,java.lang.String,java.lang.Object");
config.put("spring.json.use.type.headers", true);

// ErrorHandlingDeserializer 로 감싸 deserialize 실패 시 SerializationException 을 발생시키지 않고
// null payload + 헤더에 실패 정보를 담아 record 를 통과시킨다. DefaultErrorHandler 가 그 record 를
// recoverer(DLT publish) 로 넘기고 다음 메시지로 진행하므로 잘못된 메시지 1건이 consumer 를
// 영구 stuck 시키는 문제가 해소된다.
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName());
config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());

config.put(JsonDeserializer.TRUSTED_PACKAGES, "com.catchtable.notification.event,java.lang.String,java.lang.Object");
// 헤더 미포함 메시지(외부 시스템/CLI) 대비 fallback. 헤더가 있으면 우선 적용된다.
config.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, true);
Comment on lines +80 to +87

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

헤더 미포함 메시지 처리를 위한 StringJsonMessageConverter 도입 제안

현재 설정은 JsonDeserializer를 Consumer 레벨에서 직접 사용하고 있습니다. 주석에 언급된 것처럼 외부 시스템이나 CLI 등에서 헤더(__TypeId__) 없이 전송된 메시지를 처리하려 할 때, JsonDeserializer는 역직렬화할 대상 클래스 타입을 알 수 없어 IllegalArgumentException을 발생시키며 실패하게 됩니다. (단일 default type을 지정하려 해도, 현재 컨슈머 그룹 내에서 여러 종류의 이벤트 클래스(ReservationConfirmedEvent, VacancyEvent 등)를 처리하고 있어 설정이 불가능합니다.)

이 문제를 해결하는 가장 표준적이고 견고한 방법은 **StringJsonMessageConverter**를 사용하는 것입니다.

  1. Consumer Factory에서는 단순히 StringDeserializer를 사용하여 메시지를 문자열로 읽어옵니다.
  2. ConcurrentKafkaListenerContainerFactoryStringJsonMessageConverter를 등록하면, Spring Kafka가 @KafkaListener 메서드의 파라미터 타입(예: ReservationConfirmedEvent)을 기반으로 JSON 문자열을 해당 객체로 자동 역직렬화해 줍니다.
  3. 이 방식을 사용하면 __TypeId__ 헤더 유무와 상관없이 모든 외부 메시지를 안전하게 수신할 수 있습니다.

아래는 kafkaListenerContainerFactory 변경 예시입니다.

@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(KafkaTemplate<String, Object> kafkaTemplate) {
    ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    
    // StringJsonMessageConverter 등록 (메서드 파라미터 기반 자동 역직렬화)
    factory.setRecordMessageConverter(new StringJsonMessageConverter());

    DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
            kafkaTemplate,
            (record, ex) -> new TopicPartition(record.topic() + "-dlt", record.partition())
    );
    factory.setCommonErrorHandler(new DefaultErrorHandler(recoverer, new FixedBackOff(2000L, 3L)));

    return factory;
}
Suggested change
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName());
config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());
config.put(JsonDeserializer.TRUSTED_PACKAGES, "com.catchtable.notification.event,java.lang.String,java.lang.Object");
// 헤더 미포함 메시지(외부 시스템/CLI) 대비 fallback. 헤더가 있으면 우선 적용된다.
config.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, true);
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, StringDeserializer.class.getName());


DefaultKafkaConsumerFactory<String, Object> factory = new DefaultKafkaConsumerFactory<>(config);
factory.addListener(new MicrometerConsumerListener<>(meterRegistry));
Expand Down
Loading