From fd696776439a6ebea9606aeb8f6917800656eee7 Mon Sep 17 00:00:00 2001 From: bigwaveBigwave Date: Mon, 24 Aug 2026 13:25:46 +0900 Subject: [PATCH] =?UTF-8?q?refactor:=20=EC=86=8C=EB=B9=84=20=EB=B6=84?= =?UTF-8?q?=EB=A5=98=20=EB=B3=91=EB=A0=AC=20=EC=B2=98=EB=A6=AC=20=EB=B0=8F?= =?UTF-8?q?=20=EC=84=B1=EB=8A=A5=20=EA=B3=84=EC=B8=A1=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...astApiTransactionClassificationClient.java | 11 +- .../com/ntropy/ai/config/FastApiConfig.java | 19 +++ .../ntropy/ai/config/FastApiProperties.java | 20 +++ ...DailyTransactionClassificationService.java | 129 ++++++++++++++---- .../main/resources/fastapi.properties.example | 3 + .../ntropy/ai/config/FastApiConfigTest.java | 31 +++++ ...yTransactionClassificationServiceTest.java | 86 +++++++++++- 7 files changed, 265 insertions(+), 34 deletions(-) diff --git a/services/ai-service/src/main/java/com/ntropy/ai/client/fastapi/FastApiTransactionClassificationClient.java b/services/ai-service/src/main/java/com/ntropy/ai/client/fastapi/FastApiTransactionClassificationClient.java index 1898c99c..f0b96515 100644 --- a/services/ai-service/src/main/java/com/ntropy/ai/client/fastapi/FastApiTransactionClassificationClient.java +++ b/services/ai-service/src/main/java/com/ntropy/ai/client/fastapi/FastApiTransactionClassificationClient.java @@ -8,6 +8,7 @@ import org.springframework.http.MediaType; import org.springframework.stereotype.Component; import org.springframework.web.client.RestTemplate; +import org.springframework.http.client.SimpleClientHttpRequestFactory; import com.ntropy.ai.dto.fastapi.TransactionClassificationRequest; import com.ntropy.ai.dto.fastapi.TransactionClassificationResponse; @@ -20,17 +21,23 @@ @Component public class FastApiTransactionClassificationClient { - private final RestTemplate restTemplate = new RestTemplate(); + private final RestTemplate restTemplate; private final String fastApiBaseUrl; @Autowired public FastApiTransactionClassificationClient(FastApiProperties properties) { this.fastApiBaseUrl = properties.getBaseUrl(); + SimpleClientHttpRequestFactory requestFactory = + new SimpleClientHttpRequestFactory(); + requestFactory.setConnectTimeout(properties.getConnectTimeoutMillis()); + requestFactory.setReadTimeout(properties.getReadTimeoutMillis()); + this.restTemplate = new RestTemplate(requestFactory); } protected FastApiTransactionClassificationClient() { this.fastApiBaseUrl = null; + this.restTemplate = new RestTemplate(); } public TransactionClassificationResponse classifyTransactions( @@ -53,4 +60,4 @@ public TransactionClassificationResponse classifyTransactions( TransactionClassificationResponse.class ); } -} \ No newline at end of file +} diff --git a/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiConfig.java b/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiConfig.java index 2c4b3908..8c95cf24 100644 --- a/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiConfig.java +++ b/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiConfig.java @@ -5,6 +5,10 @@ import org.springframework.context.annotation.PropertySource; import org.springframework.context.annotation.PropertySources; import org.springframework.context.support.PropertySourcesPlaceholderConfigurer; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +import java.util.concurrent.ThreadPoolExecutor; /** FastAPI 연결 주소를 공통 설정과 외부 설정에서 로드합니다. */ @Configuration @@ -24,4 +28,19 @@ public class FastApiConfig { public static PropertySourcesPlaceholderConfigurer fastApiPropertySourcesPlaceholderConfigurer() { return new PropertySourcesPlaceholderConfigurer(); } + + /** 소비 분류 FastAPI 배치를 제한된 동시성으로 실행한다. */ + @Bean("transactionClassificationExecutor") + public ThreadPoolTaskExecutor transactionClassificationExecutor( + @Value("${fastapi.classification.parallelism:4}") int configuredParallelism + ) { + int parallelism = Math.max(1, Math.min(configuredParallelism, 8)); + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(parallelism); + executor.setMaxPoolSize(parallelism); + executor.setQueueCapacity(100); + executor.setThreadNamePrefix("txn-classification-"); + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); + return executor; + } } diff --git a/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiProperties.java b/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiProperties.java index 1ea92ef8..dda8260c 100644 --- a/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiProperties.java +++ b/services/ai-service/src/main/java/com/ntropy/ai/config/FastApiProperties.java @@ -11,6 +11,8 @@ public class FastApiProperties { private final String baseUrl; + private final int connectTimeoutMillis; + private final int readTimeoutMillis; public FastApiProperties(Environment environment) { String configuredBaseUrl = environment.getProperty("fastapi.base-url"); @@ -24,9 +26,27 @@ public FastApiProperties(Environment environment) { ); } this.baseUrl = configuredBaseUrl.trim(); + this.connectTimeoutMillis = positiveInt( + environment.getProperty("fastapi.connect-timeout-ms"), 5_000 + ); + this.readTimeoutMillis = positiveInt( + environment.getProperty("fastapi.read-timeout-ms"), 120_000 + ); } private boolean isBlank(String value) { return value == null || value.trim().isEmpty(); } + + private int positiveInt(String value, int defaultValue) { + if (isBlank(value)) { + return defaultValue; + } + try { + int parsed = Integer.parseInt(value.trim()); + return parsed > 0 ? parsed : defaultValue; + } catch (NumberFormatException ignored) { + return defaultValue; + } + } } diff --git a/services/ai-service/src/main/java/com/ntropy/ai/service/DailyTransactionClassificationService.java b/services/ai-service/src/main/java/com/ntropy/ai/service/DailyTransactionClassificationService.java index e510fecc..0a880b0a 100644 --- a/services/ai-service/src/main/java/com/ntropy/ai/service/DailyTransactionClassificationService.java +++ b/services/ai-service/src/main/java/com/ntropy/ai/service/DailyTransactionClassificationService.java @@ -7,8 +7,11 @@ import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executor; import java.util.function.Supplier; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import com.ntropy.ai.client.fastapi.FastApiTransactionClassificationClient; @@ -20,7 +23,6 @@ import com.ntropy.common.dto.account.DailyClassificationTargetTransaction; import com.ntropy.common.dto.account.TransactionAnalysisSaveItem; -import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; /** @@ -29,7 +31,6 @@ */ @Slf4j @Service -@RequiredArgsConstructor public class DailyTransactionClassificationService implements TransactionClassificationCommandClient { private static final int DB_PAGE_SIZE = 500; @@ -65,6 +66,20 @@ public class DailyTransactionClassificationService implements TransactionClassif private final TransactionPreClassificationService preClassificationService; + private final Executor classificationExecutor; + + public DailyTransactionClassificationService( + AccountTransactionAnalysisClient accountTransactionAnalysisClient, + FastApiTransactionClassificationClient fastApiClient, + TransactionPreClassificationService preClassificationService, + @Qualifier("transactionClassificationExecutor") Executor classificationExecutor + ) { + this.accountTransactionAnalysisClient = accountTransactionAnalysisClient; + this.fastApiClient = fastApiClient; + this.preClassificationService = preClassificationService; + this.classificationExecutor = classificationExecutor; + } + /** * TXN_ANALYSIS가 없는 거래가 더 이상 없을 때까지 * 최대 500건씩 반복해서 처리합니다. @@ -72,8 +87,11 @@ public class DailyTransactionClassificationService implements TransactionClassif * @return 저장한 전체 거래 분석 결과 수 */ public int run() { - return runPages(() -> accountTransactionAnalysisClient - .findUnanalyzedTransactions(DB_PAGE_SIZE)); + return runPages( + "all-users", + () -> accountTransactionAnalysisClient + .findUnanalyzedTransactions(DB_PAGE_SIZE) + ); } /** 계좌 연동을 마친 특정 사용자의 미분류 거래만 즉시 처리합니다. */ @@ -82,35 +100,54 @@ public int classifyUnanalyzedTransactions(Long userId) { if (userId == null || userId <= 0) { throw new IllegalArgumentException("userId는 양수여야 합니다."); } - return runPages(() -> accountTransactionAnalysisClient - .findUnanalyzedTransactionsByUserId(userId, DB_PAGE_SIZE)); + return runPages( + "userId=" + userId, + () -> accountTransactionAnalysisClient + .findUnanalyzedTransactionsByUserId(userId, DB_PAGE_SIZE) + ); } private int runPages( + String scope, Supplier> targetSupplier ) { + long runStartedAt = System.nanoTime(); int totalProcessed = 0; + int pageNumber = 0; while (true) { + long readStartedAt = System.nanoTime(); List targets = targetSupplier.get(); + long dbReadMillis = elapsedMillis(readStartedAt); if (targets == null || targets.isEmpty()) { + log.info( + "[일간 소비 분류] 실행 완료. scope={}, totalProcessed={}, " + + "pages={}, totalElapsedMs={}", + scope, totalProcessed, pageNumber, elapsedMillis(runStartedAt) + ); return totalProcessed; } + pageNumber++; + long classifyStartedAt = System.nanoTime(); List analyses = - classifyPage(targets); + classifyPage(scope, pageNumber, targets); + long classifyMillis = elapsedMillis(classifyStartedAt); + long saveStartedAt = System.nanoTime(); accountTransactionAnalysisClient .saveDailyTransactionAnalyses(analyses); + long dbSaveMillis = elapsedMillis(saveStartedAt); totalProcessed += analyses.size(); log.info( "[일간 소비 분류] 페이지 저장 완료. " - + "pageSize={}, totalProcessed={}", - analyses.size(), - totalProcessed + + "scope={}, page={}, pageSize={}, totalProcessed={}, dbReadMs={}, " + + "classificationMs={}, dbSaveMs={}", + scope, pageNumber, analyses.size(), totalProcessed, dbReadMillis, + classifyMillis, dbSaveMillis ); } } @@ -120,6 +157,8 @@ private int runPages( * FastAPI 대상 거래로 분리합니다. */ private List classifyPage( + String scope, + int pageNumber, List targets ) { List analyses = @@ -143,23 +182,41 @@ private List classifyPage( * FastAPI #33의 요청 최대 크기가 100건이므로 * 최대 100건씩 나눠서 호출합니다. */ - for ( - int start = 0; - start < fastApiTargets.size(); - start += FAST_API_BATCH_SIZE - ) { + List>> futures = + new ArrayList<>(); + + int batchNumber = 0; + for (int start = 0; start < fastApiTargets.size(); start += FAST_API_BATCH_SIZE) { int end = Math.min( start + FAST_API_BATCH_SIZE, fastApiTargets.size() ); - analyses.addAll( - classifyWithFastApi( - fastApiTargets.subList(start, end) + List batch = + List.copyOf(fastApiTargets.subList(start, end)); + int currentBatchNumber = ++batchNumber; + futures.add( + CompletableFuture.supplyAsync( + () -> classifyWithFastApi( + scope, pageNumber, batch, currentBatchNumber + ), + classificationExecutor ) ); } + for (CompletableFuture> future : futures) { + analyses.addAll(future.join()); + } + + log.info( + "[일간 소비 분류] 페이지 분류 완료. scope={}, page={}, " + + "targets={}, deterministic={}, " + + "fastApiTargets={}, fastApiBatches={}", + scope, pageNumber, targets.size(), analyses.size() - fastApiTargets.size(), + fastApiTargets.size(), futures.size() + ); + return analyses; } @@ -171,8 +228,12 @@ private List classifyPage( * ETC / VARIABLE로 확정합니다. */ private List classifyWithFastApi( - List targets + String scope, + int pageNumber, + List targets, + int batchNumber ) { + long startedAt = System.nanoTime(); Map targetById = new HashMap<>(); @@ -237,31 +298,45 @@ private List classifyWithFastApi( } catch (Exception exception) { log.warn( "[일간 소비 분류] FastAPI 호출 실패. " - + "fallbackCount={}", - targets.size(), + + "scope={}, page={}, batch={}, fallbackCount={}", + scope, pageNumber, batchNumber, targets.size(), exception ); } List completed = new ArrayList<>(); + int fallbackCount = 0; /* * FastAPI 응답 순서와 관계없이 원래 요청 순서대로 저장 결과를 * 생성하고, 결과가 없는 거래는 반드시 fallback 처리합니다. */ for (DailyClassificationTargetTransaction target : targets) { - completed.add( - validResults.getOrDefault( - target.getTransactionId(), - fallback(target.getTransactionId()) - ) - ); + TransactionAnalysisSaveItem result = validResults.get(target.getTransactionId()); + if (result == null) { + fallbackCount++; + result = fallback(target.getTransactionId()); + } + completed.add(result); } + log.info( + "[일간 소비 분류] FastAPI 배치 완료. scope={}, page={}, " + + "batch={}, requestCount={}, " + + "validCount={}, fallbackCount={}, elapsedMs={}", + scope, pageNumber, batchNumber, targets.size(), + validResults.size(), fallbackCount, + elapsedMillis(startedAt) + ); + return completed; } + private static long elapsedMillis(long startedAt) { + return (System.nanoTime() - startedAt) / 1_000_000L; + } + /** * FastAPI 응답이 요청 대상 거래이고 소비 결과 계약을 * 만족하는지 검증합니다. diff --git a/services/ai-service/src/main/resources/fastapi.properties.example b/services/ai-service/src/main/resources/fastapi.properties.example index 961cce85..4117f443 100644 --- a/services/ai-service/src/main/resources/fastapi.properties.example +++ b/services/ai-service/src/main/resources/fastapi.properties.example @@ -2,3 +2,6 @@ # ${NTROPY_CONFIG_DIR}/fastapi.properties로 복사해 값을 변경하세요. # base URL에는 Swagger 경로인 /docs를 포함하지 않습니다. fastapi.base-url=https://your-fastapi-server.example.com +fastapi.classification.parallelism=4 +fastapi.connect-timeout-ms=5000 +fastapi.read-timeout-ms=120000 diff --git a/services/ai-service/src/test/java/com/ntropy/ai/config/FastApiConfigTest.java b/services/ai-service/src/test/java/com/ntropy/ai/config/FastApiConfigTest.java index b95f0382..f4a6d56d 100644 --- a/services/ai-service/src/test/java/com/ntropy/ai/config/FastApiConfigTest.java +++ b/services/ai-service/src/test/java/com/ntropy/ai/config/FastApiConfigTest.java @@ -70,6 +70,37 @@ void propertyValueTakesPriorityOverEnvironmentVariable() { assertEquals("https://properties.example.test", properties.getBaseUrl()); } + @Test + void loadsClassificationHttpTimeoutsWithSafeDefaults() { + StandardEnvironment defaultsEnvironment = new StandardEnvironment(); + defaultsEnvironment.getPropertySources().addFirst( + new MapPropertySource( + "fastApiDefaults", + Map.of("fastapi.base-url", "https://fastapi.example.test") + ) + ); + + FastApiProperties defaults = new FastApiProperties(defaultsEnvironment); + assertEquals(5_000, defaults.getConnectTimeoutMillis()); + assertEquals(120_000, defaults.getReadTimeoutMillis()); + + StandardEnvironment configuredEnvironment = new StandardEnvironment(); + configuredEnvironment.getPropertySources().addFirst( + new MapPropertySource( + "fastApiTimeouts", + Map.of( + "fastapi.base-url", "https://fastapi.example.test", + "fastapi.connect-timeout-ms", "3000", + "fastapi.read-timeout-ms", "60000" + ) + ) + ); + + FastApiProperties configured = new FastApiProperties(configuredEnvironment); + assertEquals(3_000, configured.getConnectTimeoutMillis()); + assertEquals(60_000, configured.getReadTimeoutMillis()); + } + @Test void missingPropertyAndEnvironmentVariableFailsClearly() { StandardEnvironment environment = new StandardEnvironment(); diff --git a/services/ai-service/src/test/java/com/ntropy/ai/service/DailyTransactionClassificationServiceTest.java b/services/ai-service/src/test/java/com/ntropy/ai/service/DailyTransactionClassificationServiceTest.java index ba320da1..5979c434 100644 --- a/services/ai-service/src/test/java/com/ntropy/ai/service/DailyTransactionClassificationServiceTest.java +++ b/services/ai-service/src/test/java/com/ntropy/ai/service/DailyTransactionClassificationServiceTest.java @@ -4,6 +4,12 @@ import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executor; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.Test; @@ -63,7 +69,8 @@ void savesDeterministicAndFastApiResultsAndUsesFallbackForMissingResult() { new DailyTransactionClassificationService( accountClient, fastApiClient, - new TransactionPreClassificationService() + new TransactionPreClassificationService(), + Runnable::run ); assertEquals( @@ -122,12 +129,14 @@ void splitsFastApiRequestsIntoBatchesOfAtMostOneHundred() { FakeFastApiClient fastApiClient = new FakeFastApiClient(); + CountingExecutor executor = new CountingExecutor(); DailyTransactionClassificationService service = new DailyTransactionClassificationService( accountClient, fastApiClient, - new TransactionPreClassificationService() + new TransactionPreClassificationService(), + executor ); assertEquals( @@ -139,6 +148,7 @@ void splitsFastApiRequestsIntoBatchesOfAtMostOneHundred() { List.of(100, 1), fastApiClient.requestSizes ); + assertEquals(2, executor.executions); assertEquals( 101, @@ -146,6 +156,31 @@ void splitsFastApiRequestsIntoBatchesOfAtMostOneHundred() { ); } + @Test + void runsFastApiBatchesConcurrentlyOnTheConfiguredExecutor() { + List targets = new ArrayList<>(); + for (long id = 1; id <= 101; id++) { + targets.add(target(id, "ORDINARY", "알 수 없는 상점 " + id)); + } + + ConcurrentFastApiClient fastApiClient = new ConcurrentFastApiClient(); + ExecutorService executor = Executors.newFixedThreadPool(2); + try { + DailyTransactionClassificationService service = + new DailyTransactionClassificationService( + new FakeAccountClient(targets), + fastApiClient, + new TransactionPreClassificationService(), + executor + ); + + assertEquals(101, service.run()); + assertEquals(2, fastApiClient.maxConcurrentRequests.get()); + } finally { + executor.shutdownNow(); + } + } + @Test void invalidFastApiResultUsesFallback() { FakeAccountClient accountClient = @@ -179,7 +214,8 @@ void invalidFastApiResultUsesFallback() { new DailyTransactionClassificationService( accountClient, fastApiClient, - new TransactionPreClassificationService() + new TransactionPreClassificationService(), + Runnable::run ); service.run(); @@ -219,7 +255,8 @@ void continuesUntilAccountServiceReturnsNoMorePages() { new DailyTransactionClassificationService( accountClient, new FakeFastApiClient(), - new TransactionPreClassificationService() + new TransactionPreClassificationService(), + Runnable::run ); assertEquals( @@ -247,7 +284,7 @@ void classifiesOnlyTheRequestedUsersUnanalyzedTransactions() { List.of(target(1L, "ORDINARY", "스타벅스")) ); DailyTransactionClassificationService service = new DailyTransactionClassificationService( - accountClient, new FakeFastApiClient(), new TransactionPreClassificationService() + accountClient, new FakeFastApiClient(), new TransactionPreClassificationService(), Runnable::run ); assertEquals(1, service.classifyUnanalyzedTransactions(42L)); @@ -338,6 +375,45 @@ public TransactionClassificationResponse classifyTransactions( } } + private static class CountingExecutor implements Executor { + private int executions; + + @Override + public void execute(Runnable command) { + executions++; + command.run(); + } + } + + private static class ConcurrentFastApiClient + extends FastApiTransactionClassificationClient { + private final CountDownLatch bothRequestsStarted = new CountDownLatch(2); + private final AtomicInteger activeRequests = new AtomicInteger(); + private final AtomicInteger maxConcurrentRequests = new AtomicInteger(); + + @Override + public TransactionClassificationResponse classifyTransactions( + List transactions + ) { + int active = activeRequests.incrementAndGet(); + maxConcurrentRequests.accumulateAndGet(active, Math::max); + bothRequestsStarted.countDown(); + try { + if (!bothRequestsStarted.await(2, TimeUnit.SECONDS)) { + throw new IllegalStateException("두 FastAPI 배치가 동시에 시작되지 않았습니다."); + } + return new TransactionClassificationResponse( + true, 200, "ok", new TransactionClassificationData(List.of()) + ); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + throw new IllegalStateException(exception); + } finally { + activeRequests.decrementAndGet(); + } + } + } + private static class FakeAccountClient implements AccountTransactionAnalysisClient {