From 59af4bf7b1edc7435d964ff31b9b9105d8f388c5 Mon Sep 17 00:00:00 2001 From: bigwaveBigwave Date: Sun, 23 Aug 2026 13:26:58 +0900 Subject: [PATCH] =?UTF-8?q?[Feature]=20=EA=B3=84=EC=A2=8C=20=EC=97=B0?= =?UTF-8?q?=EB=8F=99=20=ED=9B=84=20=EC=82=AC=EC=9A=A9=EC=9E=90=20=EA=B1=B0?= =?UTF-8?q?=EB=9E=98=20=EC=A6=89=EC=8B=9C=20=EC=86=8C=EB=B9=84=20=EB=B6=84?= =?UTF-8?q?=EB=A5=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../AccountTransactionAnalysisClient.java | 8 +- ...ransactionClassificationCommandClient.java | 8 ++ ...LocalAccountTransactionAnalysisClient.java | 8 +- .../LocalFinancialAccountCommandClient.java | 13 ++++ .../mapper/FinancialDataQueryMapper.java | 5 ++ .../account/service/TxnAnalysisService.java | 19 ++++- .../account/FinancialDataQueryMapper.xml | 42 +++++++++++ .../client/LocalAccountQueryClientTest.java | 6 ++ ...ocalFinancialAccountCommandClientTest.java | 73 ++++++++++++++++++- .../service/TxnAnalysisServiceTest.java | 6 ++ ...DailyTransactionClassificationService.java | 27 +++++-- ...yTransactionClassificationServiceTest.java | 23 +++++- 12 files changed, 226 insertions(+), 12 deletions(-) create mode 100644 common/src/main/java/com/ntropy/common/client/TransactionClassificationCommandClient.java diff --git a/common/src/main/java/com/ntropy/common/client/AccountTransactionAnalysisClient.java b/common/src/main/java/com/ntropy/common/client/AccountTransactionAnalysisClient.java index 8a692d54..c9e86af8 100644 --- a/common/src/main/java/com/ntropy/common/client/AccountTransactionAnalysisClient.java +++ b/common/src/main/java/com/ntropy/common/client/AccountTransactionAnalysisClient.java @@ -24,6 +24,12 @@ List findUnanalyzedTransactions( int limit ); + /** 특정 사용자의 아직 분석되지 않은 일간 소비 분석 대상 거래를 조회합니다. */ + List findUnanalyzedTransactionsByUserId( + Long userId, + int limit + ); + /** * 일간 소비 분류 결과를 TXN_ANALYSIS에 저장합니다. * @@ -57,4 +63,4 @@ List findClassificationTargets( void saveTransactionAnalyses( TransactionAnalysisSaveRequest request ); -} \ No newline at end of file +} diff --git a/common/src/main/java/com/ntropy/common/client/TransactionClassificationCommandClient.java b/common/src/main/java/com/ntropy/common/client/TransactionClassificationCommandClient.java new file mode 100644 index 00000000..306c4430 --- /dev/null +++ b/common/src/main/java/com/ntropy/common/client/TransactionClassificationCommandClient.java @@ -0,0 +1,8 @@ +package com.ntropy.common.client; + +/** 계좌 거래 수집 완료 후 사용자 단위 소비 분류를 실행하는 내부 계약입니다. */ +public interface TransactionClassificationCommandClient { + + /** 특정 사용자의 아직 분석되지 않은 모든 거래를 분류합니다. */ + int classifyUnanalyzedTransactions(Long userId); +} diff --git a/services/account-service/src/main/java/com/ntropy/account/client/LocalAccountTransactionAnalysisClient.java b/services/account-service/src/main/java/com/ntropy/account/client/LocalAccountTransactionAnalysisClient.java index c7894469..da41297c 100644 --- a/services/account-service/src/main/java/com/ntropy/account/client/LocalAccountTransactionAnalysisClient.java +++ b/services/account-service/src/main/java/com/ntropy/account/client/LocalAccountTransactionAnalysisClient.java @@ -29,6 +29,12 @@ public class LocalAccountTransactionAnalysisClient return txnAnalysisService.findUnanalyzedTransactions(limit); } + @Override + public List + findUnanalyzedTransactionsByUserId(Long userId, int limit) { + return txnAnalysisService.findUnanalyzedTransactionsByUserId(userId, limit); + } + /** * 일간 배치에서 생성한 소비·비소비 분석 결과를 저장합니다. */ @@ -62,4 +68,4 @@ public void saveTransactionAnalyses( ) { txnAnalysisService.saveAnalyses(request); } -} \ No newline at end of file +} diff --git a/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java b/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java index 724ba0e3..a3296e28 100644 --- a/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java +++ b/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java @@ -20,6 +20,7 @@ import com.ntropy.account.service.VirtualFinancialDataService; import com.ntropy.account.service.VirtualFinancialDataService.GenerationSummary; import com.ntropy.common.client.FinancialAccountCommandClient; +import com.ntropy.common.client.TransactionClassificationCommandClient; import com.ntropy.common.dto.account.AccountRegistrationCommand; import com.ntropy.common.dto.account.AccountRegistrationSummary; import com.ntropy.common.dto.account.BankSummary; @@ -44,6 +45,7 @@ public class LocalFinancialAccountCommandClient implements FinancialAccountComma private final AccountLifecycleMapper accountLifecycleMapper; private final AccountMapper accountMapper; private final CodefConnectionMapper codefConnectionMapper; + private final TransactionClassificationCommandClient transactionClassificationCommandClient; @Override public List findSupportedBanks() { @@ -80,6 +82,7 @@ public AccountRegistrationSummary registerAccount(Long userId, AccountRegistrati if ("VIRTUAL".equals(connectionType)) { GenerationSummary summary = virtualAccountRegenerationService.regenerateForUser(userId, bank); + classifyTransactionsSafely(userId); return new AccountRegistrationSummary(connectionType, bank.getOrganizationCode(), summary.accounts()); } @@ -94,9 +97,19 @@ public AccountRegistrationSummary registerAccount(Long userId, AccountRegistrati ); ensureVirtualDatasetSafely(userId, bank); int accountCount = accountCollectionService.collect(userId, bank, birthDate, startDate, endDate).size(); + classifyTransactionsSafely(userId); return new AccountRegistrationSummary(connectionType, bank.getOrganizationCode(), accountCount); } + private void classifyTransactionsSafely(Long userId) { + try { + int processed = transactionClassificationCommandClient.classifyUnanalyzedTransactions(userId); + log.info("계좌 연동 후 소비 분류 완료: userId={}, processed={}", userId, processed); + } catch (RuntimeException e) { + log.warn("계좌 연동 후 소비 분류 실패: userId={}", userId, e); + } + } + private void ensureVirtualDatasetSafely(Long userId, PersonalBank bank) { try { if (accountMapper.existsAnyByUserIdAndProvider(userId, ConnectionProvider.NTROPY.name())) { diff --git a/services/account-service/src/main/java/com/ntropy/account/mapper/FinancialDataQueryMapper.java b/services/account-service/src/main/java/com/ntropy/account/mapper/FinancialDataQueryMapper.java index 36410595..ea77196d 100644 --- a/services/account-service/src/main/java/com/ntropy/account/mapper/FinancialDataQueryMapper.java +++ b/services/account-service/src/main/java/com/ntropy/account/mapper/FinancialDataQueryMapper.java @@ -57,4 +57,9 @@ List findValidTransactionIds( List findUnanalyzedTransactions( @Param("limit") int limit ); + + List findUnanalyzedTransactionsByUserId( + @Param("userId") Long userId, + @Param("limit") int limit + ); } diff --git a/services/account-service/src/main/java/com/ntropy/account/service/TxnAnalysisService.java b/services/account-service/src/main/java/com/ntropy/account/service/TxnAnalysisService.java index 31e5b900..c39619b4 100644 --- a/services/account-service/src/main/java/com/ntropy/account/service/TxnAnalysisService.java +++ b/services/account-service/src/main/java/com/ntropy/account/service/TxnAnalysisService.java @@ -36,15 +36,28 @@ public class TxnAnalysisService { */ public List findUnanalyzedTransactions(int limit) { + validateDailyQueryLimit(limit); + return financialDataQueryMapper.findUnanalyzedTransactions(limit); + } + + /** 특정 사용자의 아직 분석되지 않은 일간 분석 대상 거래를 조회합니다. */ + public List + findUnanalyzedTransactionsByUserId(Long userId, int limit) { + if (userId == null || userId <= 0) { + throw new ServiceException(AccountErrorCode.INVALID_REQUEST, "userId는 양수여야 합니다."); + } + validateDailyQueryLimit(limit); + return financialDataQueryMapper.findUnanalyzedTransactionsByUserId(userId, limit); + } + + private void validateDailyQueryLimit(int limit) { if (limit <= 0 || limit > MAX_DAILY_QUERY_SIZE) { throw new ServiceException( AccountErrorCode.INVALID_REQUEST, "limit은 1~500이어야 합니다." ); } - - return financialDataQueryMapper.findUnanalyzedTransactions(limit); } /** @@ -207,4 +220,4 @@ public void saveAnalyses( txnAnalysisMapper.upsertAnalyses(request); } -} \ No newline at end of file +} diff --git a/services/account-service/src/main/resources/mapper/account/FinancialDataQueryMapper.xml b/services/account-service/src/main/resources/mapper/account/FinancialDataQueryMapper.xml index d7d35ccc..789c7e03 100644 --- a/services/account-service/src/main/resources/mapper/account/FinancialDataQueryMapper.xml +++ b/services/account-service/src/main/resources/mapper/account/FinancialDataQueryMapper.xml @@ -43,6 +43,48 @@ LIMIT #{limit} + + account_row.account_id AS id, account_row.codef_connection_id AS codefConnectionId, diff --git a/services/account-service/src/test/java/com/ntropy/account/client/LocalAccountQueryClientTest.java b/services/account-service/src/test/java/com/ntropy/account/client/LocalAccountQueryClientTest.java index 00e6a444..3a33f567 100644 --- a/services/account-service/src/test/java/com/ntropy/account/client/LocalAccountQueryClientTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/client/LocalAccountQueryClientTest.java @@ -256,5 +256,11 @@ public List findValidTransactionIds( findUnanalyzedTransactions(int limit) { return List.of(); } + + @Override + public List + findUnanalyzedTransactionsByUserId(Long userId, int limit) { + return List.of(); + } } } diff --git a/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java b/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java index d7f34cba..7991a656 100644 --- a/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java @@ -25,6 +25,7 @@ import com.ntropy.account.service.VirtualFinancialDataService; import com.ntropy.account.service.VirtualFinancialDataService.GenerationSummary; import com.ntropy.common.dto.account.AccountRegistrationCommand; +import com.ntropy.common.client.TransactionClassificationCommandClient; import com.ntropy.common.exception.ServiceException; class LocalFinancialAccountCommandClientTest { @@ -135,6 +136,44 @@ collectionService, new StubVirtualAccountRegenerationService(), assertEquals("19900101", collectionService.lastBirthDate); } + @Test + void classifiesUsersTransactionsAfterCodefCollectionCompletes() { + StubTransactionClassificationCommandClient classificationClient = + new StubTransactionClassificationCommandClient(); + LocalFinancialAccountCommandClient client = newClient( + new StubPersonalBankAccountService(), new StubAccountCollectionService(), + new StubVirtualAccountRegenerationService(), new StubVirtualFinancialDataService(), + new StubAccountLifecycleMapper(1, 1), new StubAccountMapper(false), + new StubCodefConnectionMapper(), classificationClient + ); + + client.registerAccount( + 42L, new AccountRegistrationCommand("CODEF", "0088", "bank-id", "bank-password", null) + ); + + assertEquals(List.of(42L), classificationClient.userIds); + } + + @Test + void keepsAccountRegistrationSuccessfulWhenClassificationFails() { + StubTransactionClassificationCommandClient classificationClient = + new StubTransactionClassificationCommandClient(); + classificationClient.failure = new IllegalStateException("분류 실패"); + LocalFinancialAccountCommandClient client = newClient( + new StubPersonalBankAccountService(), new StubAccountCollectionService(), + new StubVirtualAccountRegenerationService(), new StubVirtualFinancialDataService(), + new StubAccountLifecycleMapper(1, 1), new StubAccountMapper(false), + new StubCodefConnectionMapper(), classificationClient + ); + + var result = client.registerAccount( + 42L, new AccountRegistrationCommand("VIRTUAL", "0088", null, null, null) + ); + + assertEquals("VIRTUAL", result.connectionType()); + assertEquals(List.of(42L), classificationClient.userIds); + } + @Test void registersIndustrialBankWithValidBirthDate() { StubAccountCollectionService collectionService = new StubAccountCollectionService(); @@ -418,13 +457,45 @@ private static LocalFinancialAccountCommandClient newClient( AccountLifecycleMapper lifecycleMapper, AccountMapper accountMapper, CodefConnectionMapper connectionMapper + ) { + return newClient( + personalBankAccountService, collectionService, regenerationService, + virtualFinancialDataService, lifecycleMapper, accountMapper, connectionMapper, + new StubTransactionClassificationCommandClient() + ); + } + + private static LocalFinancialAccountCommandClient newClient( + PersonalBankAccountService personalBankAccountService, + AccountCollectionService collectionService, + VirtualAccountRegenerationService regenerationService, + VirtualFinancialDataService virtualFinancialDataService, + AccountLifecycleMapper lifecycleMapper, + AccountMapper accountMapper, + CodefConnectionMapper connectionMapper, + TransactionClassificationCommandClient classificationClient ) { return new LocalFinancialAccountCommandClient( personalBankAccountService, collectionService, regenerationService, virtualFinancialDataService, - lifecycleMapper, accountMapper, connectionMapper + lifecycleMapper, accountMapper, connectionMapper, classificationClient ); } + private static class StubTransactionClassificationCommandClient + implements TransactionClassificationCommandClient { + private final List userIds = new ArrayList<>(); + private RuntimeException failure; + + @Override + public int classifyUnanalyzedTransactions(Long userId) { + userIds.add(userId); + if (failure != null) { + throw failure; + } + return 0; + } + } + private static class StubPersonalBankAccountService extends PersonalBankAccountService { private final List callOrder; private RuntimeException failure; diff --git a/services/account-service/src/test/java/com/ntropy/account/service/TxnAnalysisServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/TxnAnalysisServiceTest.java index 42ca7ef6..14bb7bf8 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/TxnAnalysisServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/TxnAnalysisServiceTest.java @@ -398,5 +398,11 @@ public List findValidTransactionIds(Long userId, String yearMonth, List + findUnanalyzedTransactionsByUserId(Long userId, int limit) { + return List.of(); + } } } 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 c3cd7e8f..e510fecc 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,6 +7,7 @@ import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.function.Supplier; import org.springframework.stereotype.Service; @@ -15,6 +16,7 @@ import com.ntropy.ai.dto.fastapi.TransactionClassificationResult; import com.ntropy.ai.dto.fastapi.TransactionForClassification; import com.ntropy.common.client.AccountTransactionAnalysisClient; +import com.ntropy.common.client.TransactionClassificationCommandClient; import com.ntropy.common.dto.account.DailyClassificationTargetTransaction; import com.ntropy.common.dto.account.TransactionAnalysisSaveItem; @@ -28,7 +30,7 @@ @Slf4j @Service @RequiredArgsConstructor -public class DailyTransactionClassificationService { +public class DailyTransactionClassificationService implements TransactionClassificationCommandClient { private static final int DB_PAGE_SIZE = 500; private static final int FAST_API_BATCH_SIZE = 100; @@ -70,12 +72,27 @@ public class DailyTransactionClassificationService { * @return 저장한 전체 거래 분석 결과 수 */ public int run() { + return runPages(() -> accountTransactionAnalysisClient + .findUnanalyzedTransactions(DB_PAGE_SIZE)); + } + + /** 계좌 연동을 마친 특정 사용자의 미분류 거래만 즉시 처리합니다. */ + @Override + public int classifyUnanalyzedTransactions(Long userId) { + if (userId == null || userId <= 0) { + throw new IllegalArgumentException("userId는 양수여야 합니다."); + } + return runPages(() -> accountTransactionAnalysisClient + .findUnanalyzedTransactionsByUserId(userId, DB_PAGE_SIZE)); + } + + private int runPages( + Supplier> targetSupplier + ) { int totalProcessed = 0; while (true) { - List targets = - accountTransactionAnalysisClient - .findUnanalyzedTransactions(DB_PAGE_SIZE); + List targets = targetSupplier.get(); if (targets == null || targets.isEmpty()) { return totalProcessed; @@ -313,4 +330,4 @@ private TransactionAnalysisSaveItem fallback( "VARIABLE" ); } -} \ No newline at end of file +} 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 4dc38a27..ba320da1 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 @@ -241,6 +241,19 @@ void continuesUntilAccountServiceReturnsNoMorePages() { ); } + @Test + void classifiesOnlyTheRequestedUsersUnanalyzedTransactions() { + FakeAccountClient accountClient = new FakeAccountClient( + List.of(target(1L, "ORDINARY", "스타벅스")) + ); + DailyTransactionClassificationService service = new DailyTransactionClassificationService( + accountClient, new FakeFastApiClient(), new TransactionPreClassificationService() + ); + + assertEquals(1, service.classifyUnanalyzedTransactions(42L)); + assertEquals(List.of(42L, 42L), accountClient.queriedUserIds); + } + private void assertSaved( List saved, Long transactionId, @@ -334,6 +347,7 @@ private static class FakeAccountClient private int pageIndex; private int queryCalls; + private final List queriedUserIds = new ArrayList<>(); private final List saved = new ArrayList<>(); @@ -367,6 +381,13 @@ private FakeAccountClient( return pages.get(pageIndex++); } + @Override + public List + findUnanalyzedTransactionsByUserId(Long userId, int limit) { + queriedUserIds.add(userId); + return findUnanalyzedTransactions(limit); + } + @Override public void saveDailyTransactionAnalyses( List analyses @@ -390,4 +411,4 @@ public void saveTransactionAnalyses( // 기존 월간 분류 계약은 이 테스트에서 사용하지 않습니다. } } -} \ No newline at end of file +}