From 8a7ceacb32cae11f5016ea08b6c414eb1bbaa63a Mon Sep 17 00:00:00 2001 From: eastargh Date: Thu, 27 Aug 2026 11:35:13 +0900 Subject: [PATCH] =?UTF-8?q?=EA=B8=B0=EB=8A=A5=20=EA=B0=9C=EC=84=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../jsp/onl/kjb/statistics/apiUseStatsMan.jsp | 3 +- .../djb/apistatus/ApiStatusRepository.java | 32 +++++++++ .../ext/djb/apistatus/ApiStatusService.java | 72 +++++++++++++++---- .../service/ums/UmsDispatchService.java | 42 +++++++---- .../deploy/StandardMessageInfoDeploy.java | 4 ++ 5 files changed, 126 insertions(+), 27 deletions(-) diff --git a/WebContent/jsp/onl/kjb/statistics/apiUseStatsMan.jsp b/WebContent/jsp/onl/kjb/statistics/apiUseStatsMan.jsp index 7cb2a86..e6cdf7a 100644 --- a/WebContent/jsp/onl/kjb/statistics/apiUseStatsMan.jsp +++ b/WebContent/jsp/onl/kjb/statistics/apiUseStatsMan.jsp @@ -203,7 +203,8 @@ { name: 'maxRespTime', align: 'right', width: '90', formatter: numberFormatter, sortable: false } ], jsonReader: { repeatitems: false }, - height: 500, + height: 500, + rowNum: 10000, autowidth: true, footerrow: true, userDataOnFooter: true, diff --git a/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusRepository.java b/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusRepository.java index fd9b611..1aa9e1f 100644 --- a/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusRepository.java +++ b/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusRepository.java @@ -1,6 +1,8 @@ package com.eactive.eai.rms.ext.djb.apistatus; +import java.time.LocalDateTime; import java.util.List; +import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; @@ -82,6 +84,36 @@ public interface ApiStatusRepository extends BaseRepository { @Param("delayAvgRespTime") int delayAvgRespTime ); + /** + * 현재 저장된 STATUS_CODE가 prevStatusCode와 일치할 때만 newStatusCode로 갱신한다 (CAS). + * + *

{@code ApiStatusMonitorJob}이 둘 이상의 인스턴스에서 겹쳐 실행될 수 있는 경우를 대비한 + * 방어다. 두 실행이 동시에 같은 API의 상태 전이(event)를 감지해도, 먼저 커밋한 쪽만 + * 실제로 갱신에 성공(영향 1행)하고 나머지는 0행이 되어 중복 알림 발송을 막을 수 있다.

+ * + * @return 갱신된 행 수(0 또는 1). 0이면 이미 다른 실행이 선점했거나 최초 기록(행 없음)이다. + */ + @Modifying + @Query("UPDATE ApiStatus a SET a.statusCode = :newStatusCode, a.modifiedBy = 'SCHEDULER', a.modifiedDate = CURRENT_TIMESTAMP " + + "WHERE a.eaisvcname = :id AND a.statusCode = :prevStatusCode") + int updateStatusCodeIfMatches(@Param("id") String id, + @Param("prevStatusCode") String prevStatusCode, + @Param("newStatusCode") String newStatusCode); + + /** + * API_STATUS 최초 기록용 INSERT. {@link #updateStatusCodeIfMatches}가 0행을 갱신했을 때 + * (해당 API의 상태 row가 아직 없을 때) 시도한다. 그 사이 다른 실행이 먼저 INSERT했다면 + * PK(EAISVCNAME) 제약 위반으로 실패하며, 이는 선점 실패로 처리한다. + */ + @Modifying + @Query(nativeQuery = true, value = + "INSERT INTO API_STATUS (EAISVCNAME, STATUS_CODE, MODIFIED_BY, MODIFIED_DATE) " + + "VALUES (:id, :statusCode, :modifiedBy, :modifiedDate)") + void insertIfAbsent(@Param("id") String id, + @Param("statusCode") String statusCode, + @Param("modifiedBy") String modifiedBy, + @Param("modifiedDate") LocalDateTime modifiedDate); + /** * 판정 구간에 트래픽이 있는 API 의 오류율·지연 판정값. 개발 모드(-Deai.systemmode=D) 진단 로그 전용. * diff --git a/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusService.java b/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusService.java index c2f3152..3ae8864 100644 --- a/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusService.java +++ b/src/main/java/com/eactive/eai/rms/ext/djb/apistatus/ApiStatusService.java @@ -2,22 +2,29 @@ package com.eactive.eai.rms.ext.djb.apistatus; import java.math.BigDecimal; import java.math.RoundingMode; +import java.time.LocalDateTime; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import javax.annotation.PostConstruct; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.dao.DataIntegrityViolationException; import org.springframework.stereotype.Service; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionDefinition; import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionTemplate; import com.eactive.apim.portal.template.entity.MessageCode; import com.eactive.eai.rms.common.context.MonitoringContext; import com.eactive.eai.rms.data.entity.man.role.Role; import com.eactive.eai.rms.data.entity.onl.djb.apistatus.ApiErrorRateRow; -import com.eactive.eai.rms.data.entity.onl.djb.apistatus.ApiStatus; import com.eactive.eai.rms.data.entity.onl.djb.apistatus.ApiStatusEvent; import com.eactive.eai.rms.ext.djb.ums.UmsManager; @@ -39,8 +46,21 @@ public class ApiStatusService { @Autowired private MonitoringContext monitoringContext; - - + + @Autowired + @Qualifier("transactionManagerForEMS") + private PlatformTransactionManager transactionManager; + + /** {@link #claimStatusChange} 의 최초 INSERT 시도를 바깥 트랜잭션과 분리하기 위한 템플릿 (REQUIRES_NEW) */ + private TransactionTemplate insertTransactionTemplate; + + @PostConstruct + private void initTransactionTemplate() { + this.insertTransactionTemplate = new TransactionTemplate(transactionManager); + this.insertTransactionTemplate.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRES_NEW); + } + + public void updateApiStatus(HashMap param) { int errorRate = Integer.parseInt(param.get("errorRate")); @@ -64,16 +84,16 @@ public class ApiStatusService { String newStatusCode = resolveStatusCode(event.getEvent()); if (newStatusCode == null) continue; - ApiStatus apiStatus = apiStatusRepository.findById(event.getEaisvcname()).orElse(new ApiStatus()); + // 인스턴스가 둘 이상 동시에 이 job을 돌릴 수 있으므로, CAS로 선점(claim)에 성공한 경우만 처리한다. + // (그렇지 않으면 두 인스턴스가 같은 상태 전이를 각자 감지해 알림을 중복 발송하게 된다) + if (!claimStatusChange(event.getEaisvcname(), event.getPrevStatusCode(), newStatusCode)) { + log.debug("API 상태 변경 선점 실패(다른 인스턴스가 이미 처리) - {}", event.getEaisvcname()); + continue; + } - apiStatus.setEaisvcname(event.getEaisvcname()); - apiStatus.setStatusCode(newStatusCode); - - apiStatusRepository.saveAndFlush(apiStatus); // PK 있으면 UPDATE, 없으면 INSERT (즉시 flush) - log.debug("API 상태 변경: {}-{} {} → {}", event.getEaisvcname(), event.getEaisvcdesc(), event.getEvent(), newStatusCode); - - + + //이벤트별로 api분리 저장 List apis = (List)eventMap.get(event.getEvent()); if (apis == null) { @@ -105,8 +125,34 @@ public class ApiStatusService { apiStatusDectionService.detect(event, apiIds); } } - - + + /** + * 현재 저장된 상태가 prevStatusCode와 일치할 때만 newStatusCode로 갱신한다 (CAS). + * + *

갱신 대상 행이 없으면(해당 API 최초 기록) INSERT를 시도한다. 그 사이 다른 인스턴스가 + * 먼저 INSERT했다면 PK 제약 위반이 발생하는데, 이를 선점 실패로 처리한다. 이 INSERT 시도는 + * 실패해도 {@code updateApiStatus}의 바깥 트랜잭션 전체가 rollback-only가 되지 않도록 + * 별도 트랜잭션(REQUIRES_NEW)으로 격리한다.

+ * + * @return 이번 실행이 상태 변경을 선점했으면 true + */ + private boolean claimStatusChange(String eaisvcname, String prevStatusCode, String newStatusCode) { + int updated = apiStatusRepository.updateStatusCodeIfMatches(eaisvcname, prevStatusCode, newStatusCode); + if (updated == 1) { + return true; + } + + try { + insertTransactionTemplate.execute(status -> { + apiStatusRepository.insertIfAbsent(eaisvcname, newStatusCode, "SCHEDULER", LocalDateTime.now()); + return null; + }); + return true; + } catch (DataIntegrityViolationException e) { + return false; + } + } + /** * 판정 구간에 거래가 있는 API 의 오류율·지연 판정값을 debug 로그로 남긴다. * 개발 모드(-Deai.systemmode=D) 에서만 동작한다. diff --git a/src/main/java/com/eactive/eai/rms/onl/common/service/ums/UmsDispatchService.java b/src/main/java/com/eactive/eai/rms/onl/common/service/ums/UmsDispatchService.java index 02bc0c0..a729fd4 100644 --- a/src/main/java/com/eactive/eai/rms/onl/common/service/ums/UmsDispatchService.java +++ b/src/main/java/com/eactive/eai/rms/onl/common/service/ums/UmsDispatchService.java @@ -1,6 +1,7 @@ package com.eactive.eai.rms.onl.common.service.ums; import java.time.LocalDateTime; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; @@ -81,33 +82,48 @@ public class UmsDispatchService { public void processMessages() { log.info("Starting to process pending messages"); - // 상태가 PENDING인 메시지들을 조회하고 바로 PROCESSING으로 변경 + // 상태가 PENDING인 메시지들을 조회 List pendingMessages = entityManager.createQuery( "SELECT m FROM MessageRequest m WHERE m.requestStatus = :status " + "ORDER BY m.requestDate ASC") .setParameter("status", "PENDING") .setMaxResults(100) .getResultList(); - - if ( !pendingMessages.isEmpty()) { - log.info("Found {} pending messages, updating status to PROCESSING", pendingMessages.size()); - // 조회된 메시지들의 상태를 일괄 PROCESSING으로 변경 - entityManager.createQuery( - "UPDATE MessageRequest m SET m.requestStatus = :newStatus " + - "WHERE m IN :messages") - .setParameter("newStatus", "PROCESSING") - .setParameter("messages", pendingMessages) - .executeUpdate(); + if ( !pendingMessages.isEmpty()) { + log.info("Found {} pending messages, claiming for processing", pendingMessages.size()); + + // 건별로 상태를 재검증하며 PROCESSING으로 선점한다. + // (여러 인스턴스가 동시에 조회하면 같은 PENDING row를 함께 볼 수 있으므로, + // PK만으로 일괄 UPDATE하면 두 인스턴스가 같은 메시지를 중복 발송하게 된다. + // 건별로 requestStatus='PENDING' 조건을 다시 걸어 선점(claim)에 성공한 건만 처리 대상으로 삼는다.) + List claimedMessages = new ArrayList<>(); + for (MessageRequest message : pendingMessages) { + int claimed = entityManager.createQuery( + "UPDATE MessageRequest m SET m.requestStatus = :newStatus " + + "WHERE m.id = :id AND m.requestStatus = :oldStatus") + .setParameter("newStatus", "PROCESSING") + .setParameter("oldStatus", "PENDING") + .setParameter("id", message.getId()) + .executeUpdate(); + if (claimed == 1) { + claimedMessages.add(message); + } + } + + if (claimedMessages.isEmpty()) { + log.info("No messages claimed (already taken by another instance)"); + return; + } // 영속성 컨텍스트 초기화 (벌크 연산 후) entityManager.flush(); entityManager.clear(); - // 상태가 변경된 메시지들을 다시 조회 + // 선점에 성공한 메시지들을 다시 조회 List messagesToProcess = entityManager.createQuery( "SELECT m FROM MessageRequest m WHERE m IN :messages") - .setParameter("messages", pendingMessages) + .setParameter("messages", claimedMessages) .getResultList(); log.info("Starting to submit {} messages for processing", messagesToProcess.size()); diff --git a/src/main/java/com/eactive/eai/rms/onl/loader/intf/deploy/StandardMessageInfoDeploy.java b/src/main/java/com/eactive/eai/rms/onl/loader/intf/deploy/StandardMessageInfoDeploy.java index 1d825e5..dc058d9 100644 --- a/src/main/java/com/eactive/eai/rms/onl/loader/intf/deploy/StandardMessageInfoDeploy.java +++ b/src/main/java/com/eactive/eai/rms/onl/loader/intf/deploy/StandardMessageInfoDeploy.java @@ -53,6 +53,10 @@ public class StandardMessageInfoDeploy extends DeployResource