@@ -204,6 +204,7 @@
|
||||
],
|
||||
jsonReader: { repeatitems: false },
|
||||
height: 500,
|
||||
rowNum: 10000,
|
||||
autowidth: true,
|
||||
footerrow: true,
|
||||
userDataOnFooter: true,
|
||||
|
||||
@@ -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<ApiStatus, String> {
|
||||
@Param("delayAvgRespTime") int delayAvgRespTime
|
||||
);
|
||||
|
||||
/**
|
||||
* 현재 저장된 STATUS_CODE가 prevStatusCode와 일치할 때만 newStatusCode로 갱신한다 (CAS).
|
||||
*
|
||||
* <p>{@code ApiStatusMonitorJob}이 둘 이상의 인스턴스에서 겹쳐 실행될 수 있는 경우를 대비한
|
||||
* 방어다. 두 실행이 동시에 같은 API의 상태 전이(event)를 감지해도, 먼저 커밋한 쪽만
|
||||
* 실제로 갱신에 성공(영향 1행)하고 나머지는 0행이 되어 중복 알림 발송을 막을 수 있다.</p>
|
||||
*
|
||||
* @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) 진단 로그 전용.
|
||||
*
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -40,6 +47,19 @@ 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<String, String> param) {
|
||||
|
||||
@@ -64,12 +84,12 @@ public class ApiStatusService {
|
||||
String newStatusCode = resolveStatusCode(event.getEvent());
|
||||
if (newStatusCode == null) continue;
|
||||
|
||||
ApiStatus apiStatus = apiStatusRepository.findById(event.getEaisvcname()).orElse(new ApiStatus());
|
||||
|
||||
apiStatus.setEaisvcname(event.getEaisvcname());
|
||||
apiStatus.setStatusCode(newStatusCode);
|
||||
|
||||
apiStatusRepository.saveAndFlush(apiStatus); // PK 있으면 UPDATE, 없으면 INSERT (즉시 flush)
|
||||
// 인스턴스가 둘 이상 동시에 이 job을 돌릴 수 있으므로, CAS로 선점(claim)에 성공한 경우만 처리한다.
|
||||
// (그렇지 않으면 두 인스턴스가 같은 상태 전이를 각자 감지해 알림을 중복 발송하게 된다)
|
||||
if (!claimStatusChange(event.getEaisvcname(), event.getPrevStatusCode(), newStatusCode)) {
|
||||
log.debug("API 상태 변경 선점 실패(다른 인스턴스가 이미 처리) - {}", event.getEaisvcname());
|
||||
continue;
|
||||
}
|
||||
|
||||
log.debug("API 상태 변경: {}-{} {} → {}", event.getEaisvcname(), event.getEaisvcdesc(), event.getEvent(), newStatusCode);
|
||||
|
||||
@@ -106,6 +126,32 @@ public class ApiStatusService {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 현재 저장된 상태가 prevStatusCode와 일치할 때만 newStatusCode로 갱신한다 (CAS).
|
||||
*
|
||||
* <p>갱신 대상 행이 없으면(해당 API 최초 기록) INSERT를 시도한다. 그 사이 다른 인스턴스가
|
||||
* 먼저 INSERT했다면 PK 제약 위반이 발생하는데, 이를 선점 실패로 처리한다. 이 INSERT 시도는
|
||||
* 실패해도 {@code updateApiStatus}의 바깥 트랜잭션 전체가 rollback-only가 되지 않도록
|
||||
* 별도 트랜잭션(REQUIRES_NEW)으로 격리한다.</p>
|
||||
*
|
||||
* @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 로그로 남긴다.
|
||||
|
||||
@@ -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,7 +82,7 @@ public class UmsDispatchService {
|
||||
public void processMessages() {
|
||||
log.info("Starting to process pending messages");
|
||||
|
||||
// 상태가 PENDING인 메시지들을 조회하고 바로 PROCESSING으로 변경
|
||||
// 상태가 PENDING인 메시지들을 조회
|
||||
List<MessageRequest> pendingMessages = entityManager.createQuery(
|
||||
"SELECT m FROM MessageRequest m WHERE m.requestStatus = :status " +
|
||||
"ORDER BY m.requestDate ASC")
|
||||
@@ -90,24 +91,39 @@ public class UmsDispatchService {
|
||||
.getResultList();
|
||||
|
||||
if ( !pendingMessages.isEmpty()) {
|
||||
log.info("Found {} pending messages, updating status to PROCESSING", pendingMessages.size());
|
||||
log.info("Found {} pending messages, claiming for processing", pendingMessages.size());
|
||||
|
||||
// 조회된 메시지들의 상태를 일괄 PROCESSING으로 변경
|
||||
entityManager.createQuery(
|
||||
// 건별로 상태를 재검증하며 PROCESSING으로 선점한다.
|
||||
// (여러 인스턴스가 동시에 조회하면 같은 PENDING row를 함께 볼 수 있으므로,
|
||||
// PK만으로 일괄 UPDATE하면 두 인스턴스가 같은 메시지를 중복 발송하게 된다.
|
||||
// 건별로 requestStatus='PENDING' 조건을 다시 걸어 선점(claim)에 성공한 건만 처리 대상으로 삼는다.)
|
||||
List<MessageRequest> claimedMessages = new ArrayList<>();
|
||||
for (MessageRequest message : pendingMessages) {
|
||||
int claimed = entityManager.createQuery(
|
||||
"UPDATE MessageRequest m SET m.requestStatus = :newStatus " +
|
||||
"WHERE m IN :messages")
|
||||
"WHERE m.id = :id AND m.requestStatus = :oldStatus")
|
||||
.setParameter("newStatus", "PROCESSING")
|
||||
.setParameter("messages", pendingMessages)
|
||||
.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<MessageRequest> 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());
|
||||
|
||||
+4
@@ -53,6 +53,10 @@ public class StandardMessageInfoDeploy extends DeployResource<StandardMessageInf
|
||||
@JsonProperty("StatusCode")
|
||||
private String modfimgtstusdstcd;
|
||||
|
||||
/*인터페이스URL Full Path*/
|
||||
@JsonProperty("ApiFullPath")
|
||||
private String apifullpath;
|
||||
|
||||
/*최종수정시간*/
|
||||
@JsonProperty("LastDateTime")
|
||||
@JsonFormat(pattern = "yyyyMMddHHmmssSSS")
|
||||
|
||||
Reference in New Issue
Block a user