From b3243bfc48842812ddc696a86eaa27d1e58374e0 Mon Sep 17 00:00:00 2001 From: kimjb Date: Sun, 31 May 2026 18:47:55 +0900 Subject: [PATCH] =?UTF-8?q?Fix:=20=EC=84=9C=ED=82=B7=EB=B8=8C=EB=A0=88?= =?UTF-8?q?=EC=9D=B4=EC=BB=A4=20=ED=8C=A8=EB=84=90=20=ED=91=9C=EC=8B=9C=20?= =?UTF-8?q?=EC=98=A4=EB=A5=98=20=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../catcheat-monitoring-dashboard.json | 2 +- .../catchtable/global/config/KafkaConfig.java | 19 +++++++++++++++---- 2 files changed, 16 insertions(+), 5 deletions(-) diff --git a/grafana/provisioning/dashboards/catcheat-monitoring-dashboard.json b/grafana/provisioning/dashboards/catcheat-monitoring-dashboard.json index e27cb28..51e7884 100644 --- a/grafana/provisioning/dashboards/catcheat-monitoring-dashboard.json +++ b/grafana/provisioning/dashboards/catcheat-monitoring-dashboard.json @@ -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" } } ] }, { "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" }, diff --git a/src/main/java/com/catchtable/global/config/KafkaConfig.java b/src/main/java/com/catchtable/global/config/KafkaConfig.java index 9b67708..82dd75c 100644 --- a/src/main/java/com/catchtable/global/config/KafkaConfig.java +++ b/src/main/java/com/catchtable/global/config/KafkaConfig.java @@ -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; @@ -70,10 +72,19 @@ public ConsumerFactory 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); DefaultKafkaConsumerFactory factory = new DefaultKafkaConsumerFactory<>(config); factory.addListener(new MicrometerConsumerListener<>(meterRegistry));