diff --git a/services/account-service/src/main/java/com/ntropy/account/domain/Batching.java b/services/account-service/src/main/java/com/ntropy/account/domain/Batching.java new file mode 100644 index 00000000..14723932 --- /dev/null +++ b/services/account-service/src/main/java/com/ntropy/account/domain/Batching.java @@ -0,0 +1,31 @@ +package com.ntropy.account.domain; + +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; + +/** + * MyBatis {@code IN}/다중 {@code VALUES} 절에 넘길 리스트를 안전한 크기로 나눈다 (이슈 #233). + * MySQL packet 크기와 파라미터 수 제한을 고려해 사용자·계좌 수가 많아도 단일 쿼리가 과도하게 + * 커지지 않도록 chunk 단위로 분할한다. + */ +public final class Batching { + + private Batching() { + } + + public static List> chunk(List items, int size) { + Objects.requireNonNull(items, "items"); + if (size <= 0) { + throw new IllegalArgumentException("chunk size는 양수여야 합니다."); + } + if (items.isEmpty()) { + return List.of(); + } + List> chunks = new ArrayList<>(); + for (int i = 0; i < items.size(); i += size) { + chunks.add(items.subList(i, Math.min(i + size, items.size()))); + } + return chunks; + } +} diff --git a/services/account-service/src/main/java/com/ntropy/account/mapper/AccountMapper.java b/services/account-service/src/main/java/com/ntropy/account/mapper/AccountMapper.java index 0b50205c..7d4fd59f 100644 --- a/services/account-service/src/main/java/com/ntropy/account/mapper/AccountMapper.java +++ b/services/account-service/src/main/java/com/ntropy/account/mapper/AccountMapper.java @@ -11,11 +11,18 @@ public interface AccountMapper { void upsert(Account account); + /** {@link #upsert}의 다중 VALUES bulk 버전 (이슈 #233). 계좌 수와 무관하게 쿼리 1회로 저장한다. */ + void upsertAll(@Param("list") List accounts); + void updateAccountDetails(Account account); Account findByConnectionIdAndAccountNoHash(@Param("codefConnectionId") Long codefConnectionId, @Param("accountNoHash") String accountNoHash); + /** {@link #findByConnectionIdAndAccountNoHash}의 일괄 조회 버전 (이슈 #233). */ + List findByConnectionIdAndAccountNoHashes(@Param("codefConnectionId") Long codefConnectionId, + @Param("accountNoHashes") List accountNoHashes); + Account findByIdAndUserIdAndProvider(@Param("id") Long id, @Param("userId") Long userId, @Param("provider") String provider); diff --git a/services/account-service/src/main/java/com/ntropy/account/mapper/CodefConnectionMapper.java b/services/account-service/src/main/java/com/ntropy/account/mapper/CodefConnectionMapper.java index 47042bde..23867a7e 100644 --- a/services/account-service/src/main/java/com/ntropy/account/mapper/CodefConnectionMapper.java +++ b/services/account-service/src/main/java/com/ntropy/account/mapper/CodefConnectionMapper.java @@ -1,5 +1,7 @@ package com.ntropy.account.mapper; +import java.util.List; + import com.ntropy.account.domain.entity.CodefConnection; import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Param; @@ -14,4 +16,8 @@ public interface CodefConnectionMapper { void upsert(CodefConnection codefConnection); CodefConnection findByUserIdAndProvider(@Param("userId") Long userId, @Param("provider") String provider); + + /** 일일 동기화의 사용자별 단건 조회를 대체하는 chunk 단위 일괄 조회 (이슈 #233). */ + List findByUserIdsAndProvider(@Param("userIds") List userIds, + @Param("provider") String provider); } diff --git a/services/account-service/src/main/java/com/ntropy/account/service/AccountCollectionService.java b/services/account-service/src/main/java/com/ntropy/account/service/AccountCollectionService.java index da0fad4a..eec005f7 100644 --- a/services/account-service/src/main/java/com/ntropy/account/service/AccountCollectionService.java +++ b/services/account-service/src/main/java/com/ntropy/account/service/AccountCollectionService.java @@ -3,8 +3,10 @@ import java.time.LocalDate; import java.time.YearMonth; import java.util.ArrayList; +import java.util.LinkedHashMap; import java.util.LinkedHashSet; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.function.BooleanSupplier; @@ -22,6 +24,7 @@ import com.ntropy.account.client.codef.parser.LoanTransactionResponseParser; import com.ntropy.account.client.codef.parser.LoanTransactionResponseParser.ParsedLoan; import com.ntropy.account.domain.AccountGroup; +import com.ntropy.account.domain.Batching; import com.ntropy.account.domain.ConnectionProvider; import com.ntropy.account.domain.PersonalBank; import com.ntropy.account.domain.entity.Account; @@ -113,7 +116,22 @@ public List collectForDailySync(Long userId, PersonalB LocalDate transactionStartDate, LocalDate transactionEndDate, BooleanSupplier heartbeat) { - CodefConnection connection = requireCodefConnection(userId); + return collectForDailySync( + userId, bank, requireCodefConnection(userId), birthDate, transactionStartDate, transactionEndDate, + heartbeat + ); + } + + /** + * {@link #collectForDailySync(Long, PersonalBank, String, LocalDate, LocalDate, BooleanSupplier)}와 같지만, + * 호출자가 이미 조회한 {@link CodefConnection}을 그대로 받아 동일 사용자의 기관별 반복 조회에서 + * {@link #requireCodefConnection}을 다시 호출하지 않는다 (이슈 #233). + */ + public List collectForDailySync(Long userId, PersonalBank bank, + CodefConnection connection, String birthDate, + LocalDate transactionStartDate, + LocalDate transactionEndDate, + BooleanSupplier heartbeat) { String normalizedBirthDate = bank.normalizeBirthDate(birthDate); List savedContexts = fetchAndSaveAccounts(userId, bank, connection, heartbeat); @@ -155,6 +173,9 @@ private CodefConnection requireCodefConnection(Long userId) { return connection; } + /** MySQL packet 크기·MyBatis 파라미터 수를 고려한 계좌 bulk upsert/조회 batch 크기 (이슈 #233). */ + private static final int ACCOUNT_UPSERT_BATCH_SIZE = 200; + /** 실행마다 보유계좌를 재조회해 원문 계좌번호를 이 요청 흐름 안에서만 확보한다(저장하지 않음). */ private List fetchAndSaveAccounts(Long userId, PersonalBank bank, CodefConnection connection, BooleanSupplier heartbeat) { @@ -170,13 +191,36 @@ private List fetchAndSaveAccounts(Long userId, PersonalBank List parsedAccounts = AccountResponseParser.parse( accountListResponse.path("data"), connection.getId(), userId, bank.getOrganizationCode() ); + if (parsedAccounts.isEmpty()) { + requireLease(heartbeat); + return List.of(); + } + + Map savedByHash = new LinkedHashMap<>(); + for (List chunk : Batching.chunk(parsedAccounts, ACCOUNT_UPSERT_BATCH_SIZE)) { + requireLease(heartbeat); + List accountsToUpsert = chunk.stream().map(ParsedAccount::account).toList(); + accountMapper.upsertAll(accountsToUpsert); + List accountNoHashes = accountsToUpsert.stream().map(Account::getAccountNoHash).toList(); + List savedAccounts = + accountMapper.findByConnectionIdAndAccountNoHashes(connection.getId(), accountNoHashes); + Map savedChunkByHash = new LinkedHashMap<>(); + for (Account saved : savedAccounts) { + savedChunkByHash.put(saved.getAccountNoHash(), saved); + savedByHash.put(saved.getAccountNoHash(), saved); + } + if (!savedChunkByHash.keySet().containsAll(accountNoHashes)) { + throw new IllegalStateException("CODEF 계좌 bulk upsert 결과가 누락되었습니다."); + } + requireLease(heartbeat); + } List savedContexts = new ArrayList<>(); for (ParsedAccount parsed : parsedAccounts) { - accountMapper.upsert(parsed.account()); - Account saved = accountMapper.findByConnectionIdAndAccountNoHash( - connection.getId(), parsed.account().getAccountNoHash() - ); + Account saved = savedByHash.get(parsed.account().getAccountNoHash()); + if (saved == null) { + throw new IllegalStateException("CODEF 계좌 bulk upsert 결과를 매핑할 수 없습니다."); + } savedContexts.add(new SavedAccountContext(saved, parsed.rawAccountNo())); } requireLease(heartbeat); diff --git a/services/account-service/src/main/java/com/ntropy/account/service/DailyCodefSyncService.java b/services/account-service/src/main/java/com/ntropy/account/service/DailyCodefSyncService.java index 4868349d..97b79f8d 100644 --- a/services/account-service/src/main/java/com/ntropy/account/service/DailyCodefSyncService.java +++ b/services/account-service/src/main/java/com/ntropy/account/service/DailyCodefSyncService.java @@ -15,6 +15,7 @@ import com.ntropy.account.config.IncrementalSyncPolicy; import com.ntropy.account.domain.AccountSyncStatus; +import com.ntropy.account.domain.Batching; import com.ntropy.account.domain.ConnectionProvider; import com.ntropy.account.domain.IncrementalSyncRangeCalculator; import com.ntropy.account.domain.InstitutionKeys; @@ -46,6 +47,9 @@ public class DailyCodefSyncService { public static final String JOB_NAME = "daily-sync-codef"; + /** 사용자 ID IN 절이 과도하게 커지지 않도록 나누는 chunk 크기 (이슈 #233). */ + private static final int USER_ID_CHUNK_SIZE = 500; + private final CodefConnectionMapper codefConnectionMapper; private final AccountSyncStateMapper accountSyncStateMapper; private final AccountTransactionMapper accountTransactionMapper; @@ -62,9 +66,11 @@ public DailyFinancialSyncResult synchronize(List activeUserIds, LocalDate long processedTransactionCount = 0; boolean leaseLost = false; + Map connectionsByUserId = fetchConnectionsByUserId(activeUserIds); + userLoop: for (Long userId : activeUserIds) { - CodefConnection connection = codefConnectionMapper.findByUserIdAndProvider(userId, ConnectionProvider.CODEF.name()); + CodefConnection connection = connectionsByUserId.get(userId); if (connection == null || connection.getConnectedId() == null || connection.getConnectedId().isBlank()) { continue; // 이 provider의 동기화 대상이 아닌 사용자 } @@ -157,6 +163,17 @@ public DailyFinancialSyncResult synchronize(List activeUserIds, LocalDate ); } + /** 활성 사용자의 CODEF 연결을 chunk 단위로 일괄 조회한다 (이슈 #233). */ + private Map fetchConnectionsByUserId(List userIds) { + Map connectionsByUserId = new LinkedHashMap<>(); + for (List chunk : Batching.chunk(userIds, USER_ID_CHUNK_SIZE)) { + for (CodefConnection connection : codefConnectionMapper.findByUserIdsAndProvider(chunk, ConnectionProvider.CODEF.name())) { + connectionsByUserId.put(connection.getUserId(), connection); + } + } + return connectionsByUserId; + } + private InstitutionSyncOutcome synchronizeInstitution(Long userId, PersonalBank bank, CodefConnection connection, LocalDate businessDate, LeaseHandle lease) { String organizationCode = bank.getOrganizationCode(); @@ -185,7 +202,7 @@ private InstitutionSyncOutcome synchronizeInstitution(Long userId, PersonalBank List outcomes; try { outcomes = accountCollectionService.collectForDailySync( - userId, bank, birthDate, startDate, businessDate, () -> leaseService.heartbeat(lease) + userId, bank, connection, birthDate, startDate, businessDate, () -> leaseService.heartbeat(lease) ); } catch (LeaseLostException leaseLostSignal) { throw leaseLostSignal; // 상위 루프에서 전체 중단 처리 diff --git a/services/account-service/src/main/java/com/ntropy/account/service/DailyNtropySyncService.java b/services/account-service/src/main/java/com/ntropy/account/service/DailyNtropySyncService.java index 3be67102..20e01f86 100644 --- a/services/account-service/src/main/java/com/ntropy/account/service/DailyNtropySyncService.java +++ b/services/account-service/src/main/java/com/ntropy/account/service/DailyNtropySyncService.java @@ -16,6 +16,7 @@ import com.ntropy.account.config.IncrementalSyncPolicy; import com.ntropy.account.domain.AccountGroup; import com.ntropy.account.domain.AccountSyncStatus; +import com.ntropy.account.domain.Batching; import com.ntropy.account.domain.ConnectionProvider; import com.ntropy.account.domain.IncrementalSyncRangeCalculator; import com.ntropy.account.domain.InstitutionKeys; @@ -46,6 +47,9 @@ public class DailyNtropySyncService { public static final String JOB_NAME = "daily-sync-ntropy"; private static final String ORDINARY_DEPOSIT_TYPE_CODE = "11"; + /** 사용자 ID IN 절이 과도하게 커지지 않도록 나누는 chunk 크기 (이슈 #233). */ + private static final int USER_ID_CHUNK_SIZE = 500; + private final CodefConnectionMapper codefConnectionMapper; private final AccountMapper accountMapper; private final AccountTransactionMapper accountTransactionMapper; @@ -62,15 +66,18 @@ public DailyFinancialSyncResult synchronize(List activeUserIds, LocalDate long processedTransactionCount = 0; boolean leaseLost = false; + Map connectionsByUserId = fetchConnectionsByUserId(activeUserIds); + userLoop: for (Long userId : activeUserIds) { - CodefConnection connection = codefConnectionMapper.findByUserIdAndProvider(userId, ConnectionProvider.NTROPY.name()); + CodefConnection connection = connectionsByUserId.get(userId); if (connection == null || connection.getConnectedId() == null || connection.getConnectedId().isBlank()) { continue; // 이 provider의 동기화 대상이 아닌 사용자 } boolean userHasFailure = false; boolean userHasSuccess = false; + Map ordinaryAccountsByOrganization = findOrdinaryAccountsByOrganization(userId); for (String organizationCode : InstitutionKeys.parse(connection.getRegisteredInstitutionKeys())) { if (!leaseService.heartbeat(lease)) { @@ -79,7 +86,10 @@ public DailyFinancialSyncResult synchronize(List activeUserIds, LocalDate } accountSyncStateMapper.insertIfAbsent(pendingSyncState(connection.getId(), organizationCode)); - InstitutionGenerationOutcome outcome = generateForInstitution(userId, connection, organizationCode, businessDate); + InstitutionGenerationOutcome outcome = generateForInstitution( + connection, organizationCode, businessDate, + ordinaryAccountsByOrganization.get(organizationCode) + ); processedTransactionCount += outcome.transactionCount(); institutionResults.add(new InstitutionSyncResult( organizationCode, connection.getId(), outcome.aggregate().status(), @@ -140,14 +150,31 @@ public DailyFinancialSyncResult synchronize(List activeUserIds, LocalDate ); } - private InstitutionGenerationOutcome generateForInstitution(Long userId, CodefConnection connection, - String organizationCode, LocalDate businessDate) { - Account ordinaryAccount = accountMapper.findByUserIdAndProvider(userId, ConnectionProvider.NTROPY.name()).stream() - .filter(account -> organizationCode.equals(account.getOrganizationCode())) - .filter(account -> account.getAccountGroup() == AccountGroup.DEPOSIT_TRUST) - .filter(account -> ORDINARY_DEPOSIT_TYPE_CODE.equals(account.getDepositTypeCode())) - .findFirst() - .orElse(null); + /** 사용자의 NTROPY 수시입출 계좌를 한 번만 조회해 기관코드 기준으로 그룹핑한다 (이슈 #233). */ + private Map findOrdinaryAccountsByOrganization(Long userId) { + Map ordinaryAccountsByOrganization = new LinkedHashMap<>(); + for (Account account : accountMapper.findByUserIdAndProvider(userId, ConnectionProvider.NTROPY.name())) { + if (account.getAccountGroup() == AccountGroup.DEPOSIT_TRUST + && ORDINARY_DEPOSIT_TYPE_CODE.equals(account.getDepositTypeCode())) { + ordinaryAccountsByOrganization.putIfAbsent(account.getOrganizationCode(), account); + } + } + return ordinaryAccountsByOrganization; + } + + /** 활성 사용자의 NTROPY 연결을 chunk 단위로 일괄 조회한다 (이슈 #233). */ + private Map fetchConnectionsByUserId(List userIds) { + Map connectionsByUserId = new LinkedHashMap<>(); + for (List chunk : Batching.chunk(userIds, USER_ID_CHUNK_SIZE)) { + for (CodefConnection connection : codefConnectionMapper.findByUserIdsAndProvider(chunk, ConnectionProvider.NTROPY.name())) { + connectionsByUserId.put(connection.getUserId(), connection); + } + } + return connectionsByUserId; + } + + private InstitutionGenerationOutcome generateForInstitution(CodefConnection connection, String organizationCode, + LocalDate businessDate, Account ordinaryAccount) { if (ordinaryAccount == null) { return InstitutionGenerationOutcome.failed(InstitutionAggregate.failed("ORDINARY_ACCOUNT_NOT_FOUND")); } @@ -160,7 +187,7 @@ private InstitutionGenerationOutcome generateForInstitution(Long userId, CodefCo try { List transactions = transactionGenerator.generate( - userId, ordinaryAccount.getId(), startDate, businessDate + ordinaryAccount.getUserId(), ordinaryAccount.getId(), startDate, businessDate ); if (!transactions.isEmpty()) { accountTransactionMapper.insertAll(transactions); diff --git a/services/account-service/src/main/resources/mapper/account/AccountMapper.xml b/services/account-service/src/main/resources/mapper/account/AccountMapper.xml index ea48a257..0b8c9776 100644 --- a/services/account-service/src/main/resources/mapper/account/AccountMapper.xml +++ b/services/account-service/src/main/resources/mapper/account/AccountMapper.xml @@ -33,6 +33,39 @@ deactivated_at = NULL + + + INSERT INTO ACCOUNT ( + codef_connection_id, user_id, organization_code, account_group, deposit_type_code, + account_no_masked, account_no_hash, account_name, balance, + loan_contract_principal, interest_rate, currency_code, + account_start_date, maturity_date, last_tran_date, overdraft_yn, next_payment_date, + status, deactivated_at + ) VALUES + + (#{account.codefConnectionId}, #{account.userId}, #{account.organizationCode}, #{account.accountGroup}, #{account.depositTypeCode}, + #{account.accountNoMasked}, #{account.accountNoHash}, #{account.accountName}, #{account.balance}, + #{account.loanContractPrincipal}, #{account.interestRate}, #{account.currencyCode}, + #{account.accountStartDate}, #{account.maturityDate}, #{account.lastTranDate}, #{account.overdraftYn}, #{account.nextPaymentDate}, + COALESCE(#{account.status}, 'ACTIVE'), #{account.deactivatedAt}) + + ON DUPLICATE KEY UPDATE + account_no_masked = VALUES(account_no_masked), + account_name = VALUES(account_name), + balance = VALUES(balance), + loan_contract_principal = COALESCE(VALUES(loan_contract_principal), loan_contract_principal), + -- interest_rate는 고정금리 정책상 이미 값이 있으면 재동기화로 자동 덮어쓰지 않는다 + interest_rate = COALESCE(interest_rate, VALUES(interest_rate)), + currency_code = VALUES(currency_code), + account_start_date = VALUES(account_start_date), + maturity_date = COALESCE(VALUES(maturity_date), maturity_date), + last_tran_date = VALUES(last_tran_date), + overdraft_yn = VALUES(overdraft_yn), + next_payment_date = VALUES(next_payment_date), + status = 'ACTIVE', + deactivated_at = NULL + + UPDATE ACCOUNT SET account_start_date = COALESCE(#{accountStartDate}, account_start_date), @@ -73,6 +106,37 @@ AND account_no_hash = #{accountNoHash} + + + + 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 7991a656..4963e7d3 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 @@ -596,6 +596,10 @@ private static class StubAccountMapper implements AccountMapper { public void upsert(Account account) { } + @Override + public void upsertAll(List accounts) { + } + @Override public void updateAccountDetails(Account account) { } @@ -605,6 +609,11 @@ public Account findByConnectionIdAndAccountNoHash(Long codefConnectionId, String return null; } + @Override + public List findByConnectionIdAndAccountNoHashes(Long codefConnectionId, List accountNoHashes) { + return List.of(); + } + @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { return null; @@ -692,5 +701,13 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { } return connection; } + + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } } } diff --git a/services/account-service/src/test/java/com/ntropy/account/client/LocalVirtualSettlementDepositCommandClientTest.java b/services/account-service/src/test/java/com/ntropy/account/client/LocalVirtualSettlementDepositCommandClientTest.java index 4f0a4948..6951e628 100644 --- a/services/account-service/src/test/java/com/ntropy/account/client/LocalVirtualSettlementDepositCommandClientTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/client/LocalVirtualSettlementDepositCommandClientTest.java @@ -134,8 +134,10 @@ private static final class StubAccountMapper implements AccountMapper { private final List accounts = new ArrayList<>(); @Override public void upsert(Account account) { } + @Override public void upsertAll(List accountsToUpsert) { } @Override public void updateAccountDetails(Account account) { } @Override public Account findByConnectionIdAndAccountNoHash(Long connectionId, String hash) { return null; } + @Override public List findByConnectionIdAndAccountNoHashes(Long connectionId, List hashes) { return List.of(); } @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { return null; } @Override public List findByUserIdAndProvider(Long userId, String provider) { return accounts; } @Override public boolean existsAnyByUserIdAndProvider(Long userId, String provider) { return !accounts.isEmpty(); } diff --git a/services/account-service/src/test/java/com/ntropy/account/domain/BatchingTest.java b/services/account-service/src/test/java/com/ntropy/account/domain/BatchingTest.java new file mode 100644 index 00000000..547a9607 --- /dev/null +++ b/services/account-service/src/test/java/com/ntropy/account/domain/BatchingTest.java @@ -0,0 +1,40 @@ +package com.ntropy.account.domain; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import java.util.List; + +import org.junit.jupiter.api.Test; + +class BatchingTest { + + @Test + void returnsEmptyChunksForEmptyInput() { + assertEquals(List.of(), Batching.chunk(List.of(), 3)); + } + + @Test + void preservesOrderAcrossChunks() { + assertEquals( + List.of(List.of(1, 2), List.of(3, 4), List.of(5)), + Batching.chunk(List.of(1, 2, 3, 4, 5), 2) + ); + } + + @Test + void keepsExactSizeInputInOneChunk() { + assertEquals(List.of(List.of(1, 2, 3)), Batching.chunk(List.of(1, 2, 3), 3)); + } + + @Test + void rejectsNonPositiveChunkSize() { + assertThrows(IllegalArgumentException.class, () -> Batching.chunk(List.of(1), 0)); + assertThrows(IllegalArgumentException.class, () -> Batching.chunk(List.of(1), -1)); + } + + @Test + void rejectsNullInput() { + assertThrows(NullPointerException.class, () -> Batching.chunk(null, 1)); + } +} diff --git a/services/account-service/src/test/java/com/ntropy/account/integration/batch/DailySyncQueryAmplificationManualVerificationTest.java b/services/account-service/src/test/java/com/ntropy/account/integration/batch/DailySyncQueryAmplificationManualVerificationTest.java new file mode 100644 index 00000000..1306d135 --- /dev/null +++ b/services/account-service/src/test/java/com/ntropy/account/integration/batch/DailySyncQueryAmplificationManualVerificationTest.java @@ -0,0 +1,639 @@ +package com.ntropy.account.integration.batch; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeTrue; + +import java.math.BigDecimal; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.sql.Connection; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +import javax.sql.DataSource; + +import org.apache.ibatis.executor.Executor; +import org.apache.ibatis.mapping.MappedStatement; +import org.apache.ibatis.plugin.Interceptor; +import org.apache.ibatis.plugin.Intercepts; +import org.apache.ibatis.plugin.Invocation; +import org.apache.ibatis.plugin.Signature; +import org.apache.ibatis.session.ResultHandler; +import org.apache.ibatis.session.RowBounds; +import org.junit.jupiter.api.Test; +import org.mybatis.spring.SqlSessionFactoryBean; +import org.mybatis.spring.annotation.MapperScan; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.io.support.PathMatchingResourcePatternResolver; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.ntropy.account.client.codef.CodefBankTransactionClient; +import com.ntropy.account.domain.AccountGroup; +import com.ntropy.account.domain.AccountNoHash; +import com.ntropy.account.domain.AccountNoMask; +import com.ntropy.account.domain.Batching; +import com.ntropy.account.domain.InstitutionKeys; +import com.ntropy.account.domain.PersonalBank; +import com.ntropy.account.domain.entity.Account; +import com.ntropy.account.domain.entity.CodefConnection; +import com.ntropy.account.mapper.AccountMapper; +import com.ntropy.account.mapper.AccountTransactionMapper; +import com.ntropy.account.mapper.CodefConnectionMapper; +import com.ntropy.account.service.AccountCollectionService; +import com.ntropy.account.service.PersonalBankAccountService; +import com.zaxxer.hikari.HikariConfig; +import com.zaxxer.hikari.HikariDataSource; + +/** + * 이슈 #233 실제 MySQL 검증용 수동 테스트. 쿼리 수는 사용자×기관 서비스 경계로 재현하고, + * 시간은 워밍업 후 Before/After 순서를 번갈아 실행해 median/p95를 출력한다. + * RUN_DAILY_SYNC_QUERY_BENCHMARK_TEST=true일 때만 실행한다. + */ +class DailySyncQueryAmplificationManualVerificationTest { + + private static final long BASE_USER_ID = 970_000_000L; + private static final int USER_COUNT = 100; + private static final int ORG_PER_USER = 3; + private static final int ACCOUNT_PER_ORG = 10; + private static final int TOTAL_ACCOUNTS = USER_COUNT * ORG_PER_USER * ACCOUNT_PER_ORG; + private static final int ACCOUNT_UPSERT_BATCH_SIZE = 200; + private static final int USER_ID_CHUNK_SIZE = 500; + private static final List ORGS = List.of("0004", "0088", "0011"); + private static final Path REPORT_PATH = Path.of( + "build", "reports", "benchmarks", "daily-sync-query-amplification.txt" + ); + + @Test + void verifiesBulkMapperSemanticsAndServicePathAgainstRealMysql() throws Exception { + requireOptIn(); + try (AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(TestConfig.class)) { + DataSource dataSource = context.getBean(DataSource.class); + CodefConnectionMapper connectionMapper = context.getBean(CodefConnectionMapper.class); + AccountMapper accountMapper = context.getBean(AccountMapper.class); + AccountTransactionMapper transactionMapper = context.getBean(AccountTransactionMapper.class); + SqlExecutionCountInterceptor interceptor = context.getBean(SqlExecutionCountInterceptor.class); + cleanUp(dataSource); + try { + CodefConnection connection = newConnection(BASE_USER_ID, "CODEF"); + connectionMapper.upsert(connection); + connection = connectionMapper.findByUserIdAndProvider(BASE_USER_ID, "CODEF"); + assertNotNull(connection); + assertEquals(1, connectionMapper.findByUserIdsAndProvider( + List.of(BASE_USER_ID), "CODEF").size()); + + Account original = newAccount(connection.getId(), BASE_USER_ID, ORGS.get(0), "SEMANTIC-1", "old"); + original.setInterestRate(new BigDecimal("2.50")); + original.setMaturityDate(LocalDate.of(2028, 1, 1)); + Account second = newAccount(connection.getId(), BASE_USER_ID, ORGS.get(0), "SEMANTIC-2", "second"); + accountMapper.upsertAll(List.of(original, second)); + Map firstRead = byHash(accountMapper.findByConnectionIdAndAccountNoHashes( + connection.getId(), List.of(original.getAccountNoHash(), second.getAccountNoHash()))); + assertEquals(2, firstRead.size()); + Account persisted = firstRead.get(original.getAccountNoHash()); + assertNotNull(persisted.getId()); + assertNotNull(persisted.getCreatedAt()); + + Account update = newAccount(connection.getId(), BASE_USER_ID, ORGS.get(0), "SEMANTIC-1", "new"); + update.setBalance(new BigDecimal("99.00")); + update.setInterestRate(new BigDecimal("9.99")); + accountMapper.upsertAll(List.of(update)); + Account updated = accountMapper.findByConnectionIdAndAccountNoHash( + connection.getId(), original.getAccountNoHash()); + assertEquals(persisted.getId(), updated.getId()); + assertEquals("new", updated.getAccountName()); + assertEquals(0, new BigDecimal("99.00").compareTo(updated.getBalance())); + assertEquals(0, new BigDecimal("2.50").compareTo(updated.getInterestRate())); + assertEquals(LocalDate.of(2028, 1, 1), updated.getMaturityDate()); + + ObjectMapper objectMapper = new ObjectMapper(); + StubBankTransactionClient transactionClient = new StubBankTransactionClient( + objectMapper.readTree("{\"result\":{\"code\":\"CF-00000\"},\"data\":[]}")); + AccountCollectionService service = new AccountCollectionService( + new StubPersonalBankAccountService(objectMapper.readTree(smokeAccountListJson())), + connectionMapper, transactionClient, null, null, accountMapper, transactionMapper); + CodefConnection suppliedConnection = connection; + BenchmarkResult smoke = measure(interceptor, () -> { + List outcomes = service.collectForDailySync( + BASE_USER_ID, PersonalBank.SHINHAN_BANK, suppliedConnection, null, + LocalDate.of(2026, 1, 1), LocalDate.of(2026, 1, 31), () -> true); + assertEquals(3, outcomes.size()); + assertTrue(outcomes.stream().allMatch(outcome -> + outcome.status() == AccountCollectionService.AccountCollectionOutcome.Status.SUCCESS)); + }); + assertEquals(2, smoke.sqlCount(), "서비스 경로는 bulk upsert 1회 + bulk 조회 1회여야 합니다"); + assertEquals(3, transactionClient.calls); + } finally { + cleanUp(dataSource); + } + } + } + + @Test + void measuresQueryAmplificationBeforeAndAfterAgainstRealMysql() throws Exception { + requireOptIn(); + try (AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(TestConfig.class)) { + DataSource dataSource = context.getBean(DataSource.class); + SqlExecutionCountInterceptor interceptor = context.getBean(SqlExecutionCountInterceptor.class); + CodefConnectionMapper connectionMapper = context.getBean(CodefConnectionMapper.class); + AccountMapper accountMapper = context.getBean(AccountMapper.class); + cleanUp(dataSource); + try { + List userIds = userIds(); + seedConnections(connectionMapper, userIds); + Map codefConnections = byUser( + connectionMapper.findByUserIdsAndProvider(userIds, "CODEF")); + Map ntropyConnections = byUser( + connectionMapper.findByUserIdsAndProvider(userIds, "NTROPY")); + seedNtropyAccounts(accountMapper, userIds, ntropyConnections); + List beforeAccounts = fixtures(codefConnections, "BEFORE"); + List afterAccounts = fixtures(codefConnections, "AFTER"); + + BenchmarkResult connectionBefore = measure(interceptor, + () -> runConnectionsBefore(connectionMapper, userIds)); + BenchmarkResult connectionAfter = measure(interceptor, + () -> runConnectionsAfter(connectionMapper, userIds)); + BenchmarkResult duplicateBefore = measure(interceptor, + () -> runDuplicateConnectionsBefore(connectionMapper, userIds)); + BenchmarkResult duplicateAfter = measure(interceptor, () -> { }); + BenchmarkResult accountBefore = measure(interceptor, + () -> runAccountsBefore(accountMapper, beforeAccounts)); + BenchmarkResult accountAfter = measure(interceptor, + () -> runAccountsAfter(accountMapper, afterAccounts)); + BenchmarkResult ntropyBefore = measure(interceptor, + () -> runNtropyBefore(accountMapper, userIds)); + BenchmarkResult ntropyAfter = measure(interceptor, + () -> runNtropyAfter(accountMapper, userIds)); + + assertEquals(200, connectionBefore.sqlCount()); + assertEquals(2, connectionAfter.sqlCount()); + assertEquals(300, duplicateBefore.sqlCount()); + assertEquals(0, duplicateAfter.sqlCount()); + assertEquals(6_000, accountBefore.sqlCount()); + assertEquals(600, accountAfter.sqlCount(), + "사용자×기관별 계좌 10개는 각각 upsertAll 1회 + findMany 1회여야 합니다"); + assertEquals(300, ntropyBefore.sqlCount()); + assertEquals(100, ntropyAfter.sqlCount()); + int totalBefore = sum(connectionBefore, duplicateBefore, accountBefore, ntropyBefore); + int totalAfter = sum(connectionAfter, duplicateAfter, accountAfter, ntropyAfter); + assertEquals(6_800, totalBefore); + assertEquals(702, totalAfter); + + Runnable fullBefore = () -> runFullBefore( + connectionMapper, accountMapper, userIds, beforeAccounts); + Runnable fullAfter = () -> runFullAfter( + connectionMapper, accountMapper, userIds, afterAccounts); + int warmups = positiveEnvironmentInteger("DAILY_SYNC_BENCHMARK_WARMUPS", 1); + int iterations = positiveEnvironmentInteger("DAILY_SYNC_BENCHMARK_ITERATIONS", 5); + for (int i = 0; i < warmups; i++) { + assertEquals(totalBefore, measure(interceptor, fullBefore).sqlCount()); + assertEquals(totalAfter, measure(interceptor, fullAfter).sqlCount()); + } + List beforeNanos = new ArrayList<>(); + List afterNanos = new ArrayList<>(); + for (int i = 0; i < iterations; i++) { + if (i % 2 == 0) { + addSample(interceptor, fullBefore, totalBefore, beforeNanos); + addSample(interceptor, fullAfter, totalAfter, afterNanos); + } else { + addSample(interceptor, fullAfter, totalAfter, afterNanos); + addSample(interceptor, fullBefore, totalBefore, beforeNanos); + } + } + TimingSummary beforeTiming = TimingSummary.from(beforeNanos); + TimingSummary afterTiming = TimingSummary.from(afterNanos); + + String connectionPlan = explainAnalyze(dataSource, connectionLookupSql(userIds.subList(0, 2))); + String accountPlan = explainAnalyze(dataSource, accountLookupSql(afterAccounts.get(0))); + assertIndexRangePlan(connectionPlan, "uk_codef_connection_user_provider", "CODEF_CONNECTION"); + assertIndexRangePlan(accountPlan, "uk_account_connection_hash", "ACCOUNT"); + assertFalse(connectionPlan.toLowerCase().contains("covering")); + assertFalse(accountPlan.toLowerCase().contains("covering")); + + String report = buildReport(dataSource, warmups, iterations, + connectionBefore, connectionAfter, duplicateBefore, duplicateAfter, + accountBefore, accountAfter, ntropyBefore, ntropyAfter, + totalBefore, totalAfter, beforeTiming, afterTiming, connectionPlan, accountPlan); + Files.createDirectories(REPORT_PATH.getParent()); + Files.writeString(REPORT_PATH, report, StandardCharsets.UTF_8); + System.out.println(report); + } finally { + cleanUp(dataSource); + } + } + } + + private static void requireOptIn() { + assumeTrue("true".equalsIgnoreCase(System.getenv("RUN_DAILY_SYNC_QUERY_BENCHMARK_TEST")), + "실제 MySQL이 필요한 이슈 #233 검증용 테스트"); + } + + private static void seedConnections(CodefConnectionMapper mapper, List userIds) { + for (Long userId : userIds) { + mapper.upsert(newConnection(userId, "CODEF")); + mapper.upsert(newConnection(userId, "NTROPY")); + } + } + + private static CodefConnection newConnection(Long userId, String provider) { + CodefConnection connection = new CodefConnection(); + connection.setUserId(userId); + connection.setProvider(provider); + connection.setConnectedId(provider.toLowerCase() + "-bench-" + userId); + if ("CODEF".equals(provider)) { + connection.setRegisteredInstitutionKeys(InstitutionKeys.serialize(ORGS)); + } + return connection; + } + + private static Map byUser(List connections) { + Map result = new HashMap<>(); + connections.forEach(connection -> result.put(connection.getUserId(), connection)); + assertEquals(USER_COUNT, result.size()); + return result; + } + + private static void seedNtropyAccounts(AccountMapper mapper, List userIds, + Map connections) { + for (Long userId : userIds) { + List accounts = ORGS.stream().map(org -> newAccount( + connections.get(userId).getId(), userId, org, + "NTROPY-" + userId + "-" + org, "ntropy")).toList(); + mapper.upsertAll(accounts); + } + } + + private static List fixtures(Map connections, String prefix) { + List result = new ArrayList<>(USER_COUNT * ORG_PER_USER); + for (Long userId : userIds()) { + for (String org : ORGS) { + List accounts = new ArrayList<>(ACCOUNT_PER_ORG); + for (int index = 0; index < ACCOUNT_PER_ORG; index++) { + accounts.add(newAccount(connections.get(userId).getId(), userId, org, + prefix + "-" + userId + "-" + org + "-" + index, prefix.toLowerCase())); + } + result.add(new AccountFixture(connections.get(userId).getId(), List.copyOf(accounts))); + } + } + return List.copyOf(result); + } + + private static Account newAccount(Long connectionId, Long userId, String org, + String rawAccountNo, String accountName) { + Account account = new Account(); + account.setCodefConnectionId(connectionId); + account.setUserId(userId); + account.setOrganizationCode(org); + account.setAccountGroup(AccountGroup.DEPOSIT_TRUST); + account.setDepositTypeCode("11"); + account.setAccountNoMasked(AccountNoMask.mask(rawAccountNo)); + account.setAccountNoHash(AccountNoHash.hash(org, rawAccountNo)); + account.setAccountName(accountName); + account.setBalance(BigDecimal.TEN); + account.setCurrencyCode("KRW"); + return account; + } + + private static void runConnectionsBefore(CodefConnectionMapper mapper, List userIds) { + for (Long userId : userIds) { + mapper.findByUserIdAndProvider(userId, "CODEF"); + mapper.findByUserIdAndProvider(userId, "NTROPY"); + } + } + + private static void runConnectionsAfter(CodefConnectionMapper mapper, List userIds) { + for (List chunk : Batching.chunk(userIds, USER_ID_CHUNK_SIZE)) { + mapper.findByUserIdsAndProvider(chunk, "CODEF"); + mapper.findByUserIdsAndProvider(chunk, "NTROPY"); + } + } + + private static void runDuplicateConnectionsBefore(CodefConnectionMapper mapper, List userIds) { + for (Long userId : userIds) { + for (int org = 0; org < ORG_PER_USER; org++) { + mapper.findByUserIdAndProvider(userId, "CODEF"); + } + } + } + + private static void runAccountsBefore(AccountMapper mapper, List fixtures) { + for (AccountFixture fixture : fixtures) { + for (Account account : fixture.accounts()) { + mapper.upsert(account); + mapper.findByConnectionIdAndAccountNoHash(fixture.connectionId(), account.getAccountNoHash()); + } + } + } + + private static void runAccountsAfter(AccountMapper mapper, List fixtures) { + for (AccountFixture fixture : fixtures) { + for (List chunk : Batching.chunk(fixture.accounts(), ACCOUNT_UPSERT_BATCH_SIZE)) { + mapper.upsertAll(chunk); + mapper.findByConnectionIdAndAccountNoHashes(fixture.connectionId(), + chunk.stream().map(Account::getAccountNoHash).toList()); + } + } + } + + private static void runNtropyBefore(AccountMapper mapper, List userIds) { + for (Long userId : userIds) { + for (int org = 0; org < ORG_PER_USER; org++) { + mapper.findByUserIdAndProvider(userId, "NTROPY"); + } + } + } + + private static void runNtropyAfter(AccountMapper mapper, List userIds) { + userIds.forEach(userId -> mapper.findByUserIdAndProvider(userId, "NTROPY")); + } + + private static void runFullBefore(CodefConnectionMapper connectionMapper, AccountMapper accountMapper, + List userIds, List fixtures) { + runConnectionsBefore(connectionMapper, userIds); + runDuplicateConnectionsBefore(connectionMapper, userIds); + runAccountsBefore(accountMapper, fixtures); + runNtropyBefore(accountMapper, userIds); + } + + private static void runFullAfter(CodefConnectionMapper connectionMapper, AccountMapper accountMapper, + List userIds, List fixtures) { + runConnectionsAfter(connectionMapper, userIds); + runAccountsAfter(accountMapper, fixtures); + runNtropyAfter(accountMapper, userIds); + } + + private static int sum(BenchmarkResult... results) { + int total = 0; + for (BenchmarkResult result : results) { + total += result.sqlCount(); + } + return total; + } + + private static void addSample(SqlExecutionCountInterceptor interceptor, Runnable scenario, + int expectedCount, List samples) { + BenchmarkResult result = measure(interceptor, scenario); + assertEquals(expectedCount, result.sqlCount()); + samples.add(result.elapsedNanos()); + } + + private static String connectionLookupSql(List userIds) { + String ids = userIds.stream().map(String::valueOf) + .collect(java.util.stream.Collectors.joining(",")); + return """ + SELECT codef_connection_id AS id, user_id AS userId, provider, + connected_id AS connectedId, registered_institution_keys AS registeredInstitutionKeys, + birth_date_ciphertext AS birthDateCiphertext, birth_date_iv AS birthDateIv, + birth_date_key_version AS birthDateKeyVersion, created_at AS createdAt, updated_at AS updatedAt + FROM CODEF_CONNECTION WHERE provider = 'CODEF' AND user_id IN (%s) + """.formatted(ids); + } + + private static String accountLookupSql(AccountFixture fixture) { + String hashes = fixture.accounts().stream().map(Account::getAccountNoHash) + .map(hash -> "'" + hash + "'").collect(java.util.stream.Collectors.joining(",")); + return """ + SELECT account_id AS id, codef_connection_id AS codefConnectionId, user_id AS userId, + organization_code AS organizationCode, account_group AS accountGroup, + deposit_type_code AS depositTypeCode, account_no_masked AS accountNoMasked, + account_no_hash AS accountNoHash, account_name AS accountName, balance, + loan_contract_principal AS loanContractPrincipal, interest_rate AS interestRate, + currency_code AS currencyCode, account_start_date AS accountStartDate, + maturity_date AS maturityDate, last_tran_date AS lastTranDate, + overdraft_yn AS overdraftYn, next_payment_date AS nextPaymentDate, + status, deactivated_at AS deactivatedAt, created_at AS createdAt, updated_at AS updatedAt + FROM ACCOUNT WHERE codef_connection_id = %d AND account_no_hash IN (%s) + """.formatted(fixture.connectionId(), hashes); + } + + private static String explainAnalyze(DataSource dataSource, String sql) throws SQLException { + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement(); + ResultSet resultSet = statement.executeQuery("EXPLAIN ANALYZE " + sql)) { + StringBuilder plan = new StringBuilder(); + while (resultSet.next()) { + plan.append(resultSet.getString(1)).append('\n'); + } + return plan.toString(); + } + } + + private static void assertIndexRangePlan(String plan, String index, String table) { + String normalized = plan.toLowerCase(); + assertTrue(normalized.contains(index.toLowerCase()), table + "가 기대 인덱스를 사용해야 합니다:\n" + plan); + assertFalse(normalized.contains("table scan on " + table.toLowerCase()), + table + "에 full table scan이 없어야 합니다:\n" + plan); + } + + private static String buildReport(DataSource dataSource, int warmups, int iterations, + BenchmarkResult connectionBefore, BenchmarkResult connectionAfter, + BenchmarkResult duplicateBefore, BenchmarkResult duplicateAfter, + BenchmarkResult accountBefore, BenchmarkResult accountAfter, + BenchmarkResult ntropyBefore, BenchmarkResult ntropyAfter, + int totalBefore, int totalAfter, + TimingSummary beforeTiming, TimingSummary afterTiming, + String connectionPlan, String accountPlan) throws SQLException { + StringBuilder report = new StringBuilder(); + report.append("=== 이슈 #233 실제 MySQL Before/After ===\n") + .append("데이터셋: 활성 사용자 ").append(USER_COUNT).append("명 × 기관 ").append(ORG_PER_USER) + .append("개 × 계좌 ").append(ACCOUNT_PER_ORG).append("개 = CODEF 계좌 ") + .append(TOTAL_ACCOUNTS).append("개\nDB: ").append(databaseVersion(dataSource)).append('\n') + .append("범위: 개선 대상 MyBatis mapper SQL (CODEF HTTP·lease/state SQL 제외)\n\n") + .append(String.format("%-34s %10s %10s %10s%n", "측정 항목", "Before", "After", "개선율")); + appendRow(report, "CODEF+NTROPY 연결 조회 SQL", connectionBefore.sqlCount(), connectionAfter.sqlCount()); + appendRow(report, "기관별 중복 CODEF 연결 조회 SQL", duplicateBefore.sqlCount(), duplicateAfter.sqlCount()); + appendRow(report, "CODEF 계좌 저장/재조회 SQL", accountBefore.sqlCount(), accountAfter.sqlCount()); + appendRow(report, "NTROPY 계좌 조회 SQL", ntropyBefore.sqlCount(), ntropyAfter.sqlCount()); + appendRow(report, "개선 대상 Mapper SQL 합계", totalBefore, totalAfter); + report.append("\n로컬 MySQL mapper 경로 시간 (워밍업 ").append(warmups) + .append("회, 교차 실행 샘플 ").append(iterations).append("회; CODEF HTTP 제외)\n") + .append(String.format("%-16s %12s %12s %10s%n", "통계", "Before", "After", "개선율")); + appendTimingRow(report, "median", beforeTiming.medianMillis(), afterTiming.medianMillis()); + appendTimingRow(report, "p95", beforeTiming.p95Millis(), afterTiming.p95Millis()); + report.append("Before samples(ms): ").append(beforeTiming.sampleMillis()).append('\n') + .append("After samples(ms): ").append(afterTiming.sampleMillis()).append('\n') + .append("\n[EXPLAIN ANALYZE] findByUserIdsAndProvider 전체 projection\n").append(connectionPlan) + .append("\n[EXPLAIN ANALYZE] findByConnectionIdAndAccountNoHashes 전체 projection\n") + .append(accountPlan) + .append("\n판정: 두 SELECT 모두 unique key 기반 index range lookup이며 full table scan이 없습니다.\n") + .append("전체 컬럼을 읽으므로 covering index 조회라고 표현하지 않습니다.\n"); + return report.toString(); + } + + private static void appendRow(StringBuilder report, String label, int before, int after) { + report.append(String.format("%-34s %,10d %,10d %9.1f%%%n", + label, before, after, improvement(before, after))); + } + + private static void appendTimingRow(StringBuilder report, String label, double before, double after) { + report.append(String.format("%-16s %,10.1fms %,10.1fms %9.1f%%%n", + label, before, after, improvement(before, after))); + } + + private static double improvement(double before, double after) { + return before == 0 ? 0 : (1 - after / before) * 100; + } + + private static String databaseVersion(DataSource dataSource) throws SQLException { + try (Connection connection = dataSource.getConnection(); + Statement statement = connection.createStatement(); + ResultSet resultSet = statement.executeQuery("SELECT VERSION()")) { + resultSet.next(); + return resultSet.getString(1); + } + } + + private static Map byHash(List accounts) { + Map result = new HashMap<>(); + accounts.forEach(account -> result.put(account.getAccountNoHash(), account)); + return result; + } + + private static List userIds() { + List result = new ArrayList<>(USER_COUNT); + for (int i = 0; i < USER_COUNT; i++) { + result.add(BASE_USER_ID + i); + } + return List.copyOf(result); + } + + private static int positiveEnvironmentInteger(String name, int defaultValue) { + String raw = System.getenv(name); + if (raw == null || raw.isBlank()) { + return defaultValue; + } + int value = Integer.parseInt(raw); + if (value <= 0) { + throw new IllegalArgumentException(name + " must be positive"); + } + return value; + } + + private static BenchmarkResult measure(SqlExecutionCountInterceptor interceptor, Runnable action) { + interceptor.reset(); + long start = System.nanoTime(); + action.run(); + return new BenchmarkResult(interceptor.count(), System.nanoTime() - start); + } + + private static String smokeAccountListJson() { + return """ + {"result":{"code":"CF-00000"},"data":{"resDepositTrust":[ + {"resAccount":"SMOKE-001","resAccountDeposit":"11","resAccountBalance":"10000","resAccountCurrency":"KRW"}, + {"resAccount":"SMOKE-002","resAccountDeposit":"11","resAccountBalance":"20000","resAccountCurrency":"KRW"}, + {"resAccount":"SMOKE-003","resAccountDeposit":"11","resAccountBalance":"30000","resAccountCurrency":"KRW"} + ]}} + """; + } + + private static void cleanUp(DataSource dataSource) throws SQLException { + try (Connection connection = dataSource.getConnection(); Statement statement = connection.createStatement()) { + statement.executeUpdate("DELETE FROM ACCOUNT_TRANSACTION WHERE account_id IN (SELECT account_id FROM " + + "ACCOUNT WHERE user_id >= " + BASE_USER_ID + " AND user_id < " + (BASE_USER_ID + 10_000) + ")"); + statement.executeUpdate("DELETE FROM ACCOUNT WHERE user_id >= " + BASE_USER_ID + + " AND user_id < " + (BASE_USER_ID + 10_000)); + statement.executeUpdate("DELETE FROM ACCOUNT_SYNC_STATE WHERE codef_connection_id IN (SELECT " + + "codef_connection_id FROM CODEF_CONNECTION WHERE user_id >= " + BASE_USER_ID + + " AND user_id < " + (BASE_USER_ID + 10_000) + ")"); + statement.executeUpdate("DELETE FROM CODEF_CONNECTION WHERE user_id >= " + BASE_USER_ID + + " AND user_id < " + (BASE_USER_ID + 10_000)); + } + } + + private record AccountFixture(Long connectionId, List accounts) { } + private record BenchmarkResult(int sqlCount, long elapsedNanos) { } + + private record TimingSummary(double medianMillis, double p95Millis, List sampleMillis) { + static TimingSummary from(List samples) { + List sorted = samples.stream().sorted(Comparator.naturalOrder()).toList(); + int size = sorted.size(); + double median = size % 2 == 0 + ? (sorted.get(size / 2 - 1) + sorted.get(size / 2)) / 2.0 + : sorted.get(size / 2); + int p95Index = Math.max(0, (int) Math.ceil(size * 0.95) - 1); + return new TimingSummary(median / 1_000_000.0, sorted.get(p95Index) / 1_000_000.0, + samples.stream().map(value -> value / 1_000_000.0).toList()); + } + } + + private static class StubPersonalBankAccountService extends PersonalBankAccountService { + private final JsonNode response; + StubPersonalBankAccountService(JsonNode response) { super(null, null, null); this.response = response; } + @Override public JsonNode getPersonalAccountList(Long userId, PersonalBank bank) { return response; } + } + + private static class StubBankTransactionClient extends CodefBankTransactionClient { + private final JsonNode response; + private int calls; + StubBankTransactionClient(JsonNode response) { super(null); this.response = response; } + @Override + public JsonNode getPersonalTransactionList(String organizationCode, String connectedId, String account, + LocalDate startDate, LocalDate endDate, String birthDate, + boolean includeNullAccountPassword) { + calls++; + return response; + } + } + + @Intercepts({ + @Signature(type = Executor.class, method = "update", args = {MappedStatement.class, Object.class}), + @Signature(type = Executor.class, method = "query", + args = {MappedStatement.class, Object.class, RowBounds.class, ResultHandler.class}) + }) + static class SqlExecutionCountInterceptor implements Interceptor { + private final AtomicInteger count = new AtomicInteger(); + @Override public Object intercept(Invocation invocation) throws Throwable { + count.incrementAndGet(); + return invocation.proceed(); + } + void reset() { count.set(0); } + int count() { return count.get(); } + } + + @Configuration + @MapperScan("com.ntropy.account.mapper") + static class TestConfig { + @Bean + DataSource dataSource() { + HikariConfig config = new HikariConfig(); + config.setJdbcUrl(environment("DAILY_SYNC_BENCHMARK_JDBC_URL", + "jdbc:mysql://localhost:3307/db?serverTimezone=Asia/Seoul&characterEncoding=UTF-8")); + config.setUsername(environment("DAILY_SYNC_BENCHMARK_DB_USERNAME", "root")); + config.setPassword(environment("DAILY_SYNC_BENCHMARK_DB_PASSWORD", "root")); + config.setDriverClassName("com.mysql.cj.jdbc.Driver"); + config.setMaximumPoolSize(5); + return new HikariDataSource(config); + } + @Bean SqlExecutionCountInterceptor sqlExecutionCountInterceptor() { + return new SqlExecutionCountInterceptor(); + } + @Bean + SqlSessionFactoryBean sqlSessionFactory(DataSource dataSource, SqlExecutionCountInterceptor interceptor) + throws Exception { + SqlSessionFactoryBean factory = new SqlSessionFactoryBean(); + factory.setDataSource(dataSource); + factory.setMapperLocations(new PathMatchingResourcePatternResolver() + .getResources("classpath*:mapper/**/*.xml")); + factory.setPlugins(interceptor); + return factory; + } + private static String environment(String name, String defaultValue) { + String value = System.getenv(name); + return value == null || value.isBlank() ? defaultValue : value; + } + } +} diff --git a/services/account-service/src/test/java/com/ntropy/account/mapper/AccountMapperContractTest.java b/services/account-service/src/test/java/com/ntropy/account/mapper/AccountMapperContractTest.java index 1a4ed345..8d406e30 100644 --- a/services/account-service/src/test/java/com/ntropy/account/mapper/AccountMapperContractTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/mapper/AccountMapperContractTest.java @@ -1,5 +1,6 @@ package com.ntropy.account.mapper; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -13,6 +14,32 @@ class AccountMapperContractTest { + @Test + void bulkUpsertKeepsSingleUpsertColumnsAndDuplicateUpdatePolicy() throws IOException { + String mapper = readResource("mapper/account/AccountMapper.xml"); + String single = statement(mapper, "insert", "upsert"); + String bulk = statement(mapper, "insert", "upsertAll"); + + assertEquals(insertColumns(single), insertColumns(bulk), + "단건/bulk upsert의 INSERT 컬럼은 같아야 합니다"); + assertEquals(duplicateUpdateClause(single), duplicateUpdateClause(bulk), + "단건/bulk upsert의 ON DUPLICATE KEY 갱신 정책은 같아야 합니다"); + assertTrue(bulk.contains("")); + } + + @Test + void bulkLookupKeepsSingleLookupProjectionAndUsesCompositeKeyPredicate() throws IOException { + String mapper = readResource("mapper/account/AccountMapper.xml"); + String single = statement(mapper, "select", "findByConnectionIdAndAccountNoHash"); + String bulk = statement(mapper, "select", "findByConnectionIdAndAccountNoHashes"); + + assertEquals(selectProjection(single), selectProjection(bulk), + "단건/bulk 조회가 반환하는 Account 필드는 같아야 합니다"); + assertTrue(bulk.contains("codef_connection_id = #{codefConnectionId}")); + assertTrue(bulk.contains("account_no_hash IN")); + assertTrue(bulk.contains("= 0, id + " statement가 필요합니다"); + int bodyStart = mapper.indexOf('>', start) + 1; + int end = mapper.indexOf("", bodyStart); + assertTrue(end > bodyStart, id + " statement가 닫혀 있어야 합니다"); + return mapper.substring(bodyStart, end); + } + + private static String insertColumns(String statement) { + int start = statement.indexOf('(') + 1; + int end = statement.indexOf(") VALUES", start); + return normalize(statement.substring(start, end)); + } + + private static String duplicateUpdateClause(String statement) { + int start = statement.indexOf("ON DUPLICATE KEY UPDATE"); + return normalize(statement.substring(start)); + } + + private static String selectProjection(String statement) { + int start = statement.indexOf("SELECT") + "SELECT".length(); + int end = statement.indexOf("FROM ACCOUNT", start); + return normalize(statement.substring(start, end)); + } + + private static String normalize(String sql) { + return sql.replaceAll("\\s+", " ").trim(); + } } diff --git a/services/account-service/src/test/java/com/ntropy/account/mapper/CodefConnectionMapperContractTest.java b/services/account-service/src/test/java/com/ntropy/account/mapper/CodefConnectionMapperContractTest.java new file mode 100644 index 00000000..e38a4207 --- /dev/null +++ b/services/account-service/src/test/java/com/ntropy/account/mapper/CodefConnectionMapperContractTest.java @@ -0,0 +1,50 @@ +package com.ntropy.account.mapper; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; + +import org.junit.jupiter.api.Test; + +class CodefConnectionMapperContractTest { + + @Test + void bulkLookupKeepsSingleLookupProjectionAndUsesUserIdCollection() throws IOException { + String mapper = readResource("mapper/account/CodefConnectionMapper.xml"); + String single = statement(mapper, "findByUserIdAndProvider"); + String bulk = statement(mapper, "findByUserIdsAndProvider"); + + assertEquals(selectProjection(single), selectProjection(bulk), + "단건/bulk 조회가 반환하는 CodefConnection 필드는 같아야 합니다"); + assertTrue(bulk.contains("provider = #{provider}")); + assertTrue(bulk.contains("user_id IN")); + assertTrue(bulk.contains("= 0, id + " statement가 필요합니다"); + int bodyStart = mapper.indexOf('>', start) + 1; + int end = mapper.indexOf("", bodyStart); + assertTrue(end > bodyStart, id + " statement가 닫혀 있어야 합니다"); + return mapper.substring(bodyStart, end); + } + + private static String selectProjection(String statement) { + int start = statement.indexOf("SELECT") + "SELECT".length(); + int end = statement.indexOf("FROM CODEF_CONNECTION", start); + return statement.substring(start, end).replaceAll("\\s+", " ").trim(); + } + + private static String readResource(String path) throws IOException { + try (InputStream input = CodefConnectionMapperContractTest.class.getClassLoader().getResourceAsStream(path)) { + assertNotNull(input, path + "가 테스트 classpath에 있어야 합니다"); + return new String(input.readAllBytes(), StandardCharsets.UTF_8).replace("\r\n", "\n"); + } + } +} diff --git a/services/account-service/src/test/java/com/ntropy/account/service/AccountCollectionServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/AccountCollectionServiceTest.java index 0d64a9a4..93328d40 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/AccountCollectionServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/AccountCollectionServiceTest.java @@ -300,6 +300,73 @@ void throwsLeaseLostExceptionAndStopsWhenHeartbeatFailsMidAccountLoop() throws E assertEquals(1, transactionClient.calls.size(), "첫 번째 계좌만 처리되고 두 번째는 처리되면 안 됩니다"); } + @Test + void savesAccountListWithOneBulkUpsertAndOneBulkLookup() throws Exception { + String twoOrdinaryAccounts = """ + { + "result": {"code": "CF-00000"}, + "data": {"resDepositTrust": [ + {"resAccount": "110111111111", "resAccountDeposit": "11", + "resAccountBalance": "10000", "resAccountCurrency": "KRW"}, + {"resAccount": "110222222222", "resAccountDeposit": "11", + "resAccountBalance": "20000", "resAccountCurrency": "KRW"} + ]} + } + """; + FakeAccountMapper accountMapper = new FakeAccountMapper(); + AccountCollectionService service = new AccountCollectionService( + new StubPersonalBankAccountService(objectMapper.readTree(twoOrdinaryAccounts)), + new FakeCodefConnectionMapper(1L, "connected-id"), + new FakeCodefBankTransactionClient(objectMapper.readTree( + "{\"result\":{\"code\":\"CF-00000\"},\"data\":[]}")), + null, null, accountMapper, new FakeAccountTransactionMapper() + ); + + List outcomes = service.collectForDailySync( + 1L, PersonalBank.SHINHAN_BANK, null, + LocalDate.of(2026, 1, 1), LocalDate.of(2026, 1, 31), () -> true + ); + + assertEquals(2, outcomes.size()); + assertEquals(1, accountMapper.bulkUpsertCalls); + assertEquals(1, accountMapper.bulkFindCalls); + } + + @Test + void failsFastWhenBulkLookupOmitsAnUpsertedAccount() throws Exception { + String twoOrdinaryAccounts = """ + { + "result": {"code": "CF-00000"}, + "data": {"resDepositTrust": [ + {"resAccount": "110111111111", "resAccountDeposit": "11", + "resAccountBalance": "10000", "resAccountCurrency": "KRW"}, + {"resAccount": "110222222222", "resAccountDeposit": "11", + "resAccountBalance": "20000", "resAccountCurrency": "KRW"} + ]} + } + """; + FakeAccountMapper accountMapper = new FakeAccountMapper(); + accountMapper.omitLastBulkResult = true; + AccountCollectionService service = new AccountCollectionService( + new StubPersonalBankAccountService(objectMapper.readTree(twoOrdinaryAccounts)), + new FakeCodefConnectionMapper(1L, "connected-id"), + new FakeCodefBankTransactionClient(objectMapper.readTree("{}")), + null, null, accountMapper, new FakeAccountTransactionMapper() + ); + + IllegalStateException exception = assertThrows( + IllegalStateException.class, + () -> service.collectForDailySync( + 1L, PersonalBank.SHINHAN_BANK, null, + LocalDate.of(2026, 1, 1), LocalDate.of(2026, 1, 31), () -> true + ) + ); + + assertTrue(exception.getMessage().contains("bulk upsert 결과가 누락")); + assertEquals(1, accountMapper.bulkUpsertCalls); + assertEquals(1, accountMapper.bulkFindCalls); + } + @Test void collectsInstallmentSavingsForDepositType12() throws Exception { FakeCodefConnectionMapper connectionMapper = new FakeCodefConnectionMapper(1L, "connected-id"); @@ -582,6 +649,14 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { connection.setConnectedId(connectedId); return connection; } + + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } } private static class FakeAccountMapper implements AccountMapper { @@ -589,6 +664,9 @@ private static class FakeAccountMapper implements AccountMapper { private final Map store = new HashMap<>(); private long nextId = 1; private int updatedAccountDetails; + private int bulkUpsertCalls; + private int bulkFindCalls; + private boolean omitLastBulkResult; @Override public void upsert(Account account) { @@ -601,6 +679,12 @@ public void upsert(Account account) { store.put(key, account); } + @Override + public void upsertAll(List accounts) { + bulkUpsertCalls++; + accounts.forEach(this::upsert); + } + @Override public void updateAccountDetails(Account account) { updatedAccountDetails++; @@ -626,6 +710,19 @@ public Account findByConnectionIdAndAccountNoHash(Long codefConnectionId, String return store.get(codefConnectionId + ":" + accountNoHash); } + @Override + public List findByConnectionIdAndAccountNoHashes(Long codefConnectionId, List accountNoHashes) { + bulkFindCalls++; + List results = accountNoHashes.stream() + .map(hash -> findByConnectionIdAndAccountNoHash(codefConnectionId, hash)) + .filter(java.util.Objects::nonNull) + .collect(java.util.stream.Collectors.toCollection(ArrayList::new)); + if (omitLastBulkResult && !results.isEmpty()) { + results.remove(results.size() - 1); + } + return results; + } + @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { return store.values().stream() diff --git a/services/account-service/src/test/java/com/ntropy/account/service/CodefConnectionServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/CodefConnectionServiceTest.java index 0a4f5d6c..48c3a7ee 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/CodefConnectionServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/CodefConnectionServiceTest.java @@ -3,6 +3,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; +import java.util.List; + import org.junit.jupiter.api.Test; import com.ntropy.account.client.codef.CodefConnectionClient; @@ -175,5 +177,13 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return connection != null && userId.equals(connection.getUserId()) && provider.equals(connection.getProvider()) ? connection : null; } + + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } } } diff --git a/services/account-service/src/test/java/com/ntropy/account/service/DailyCodefSyncServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/DailyCodefSyncServiceTest.java index 808756ff..0593502c 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/DailyCodefSyncServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/DailyCodefSyncServiceTest.java @@ -328,6 +328,15 @@ public List collectForDailySync(Long userId, PersonalB } return byKey.getOrDefault(key, List.of()); } + + @Override + public List collectForDailySync(Long userId, PersonalBank bank, + CodefConnection connection, String birthDate, + LocalDate transactionStartDate, + LocalDate transactionEndDate, + BooleanSupplier heartbeat) { + return collectForDailySync(userId, bank, birthDate, transactionStartDate, transactionEndDate, heartbeat); + } } private static class FakeCodefConnectionMapper implements CodefConnectionMapper { @@ -354,6 +363,14 @@ public void upsert(CodefConnection codefConnection) { public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return byUserId.get(userId); } + + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } } /** diff --git a/services/account-service/src/test/java/com/ntropy/account/service/DailyNtropySyncServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/DailyNtropySyncServiceTest.java index e8de6fec..8deeb945 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/DailyNtropySyncServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/DailyNtropySyncServiceTest.java @@ -196,6 +196,14 @@ public void upsert(CodefConnection codefConnection) { public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return byUserId.get(userId); } + + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } } private static class FakeAccountMapper implements AccountMapper { @@ -210,6 +218,10 @@ void put(Long userId, Account account) { public void upsert(Account account) { } + @Override + public void upsertAll(List accounts) { + } + @Override public void updateAccountDetails(Account account) { } @@ -219,6 +231,11 @@ public Account findByConnectionIdAndAccountNoHash(Long codefConnectionId, String return null; } + @Override + public List findByConnectionIdAndAccountNoHashes(Long codefConnectionId, List accountNoHashes) { + return List.of(); + } + @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { return null; diff --git a/services/account-service/src/test/java/com/ntropy/account/service/PersonalBankAccountServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/PersonalBankAccountServiceTest.java index 7533cd66..f92f085e 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/PersonalBankAccountServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/PersonalBankAccountServiceTest.java @@ -4,6 +4,8 @@ import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; +import java.util.List; + import org.junit.jupiter.api.Test; import com.fasterxml.jackson.databind.JsonNode; @@ -206,5 +208,13 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return connection != null && userId.equals(connection.getUserId()) && provider.equals(connection.getProvider()) ? connection : null; } + + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } } } diff --git a/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountRegenerationServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountRegenerationServiceTest.java index 982af8f9..4a8661d9 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountRegenerationServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountRegenerationServiceTest.java @@ -197,6 +197,14 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return store.get(userId + ":" + provider); } + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } + private static String key(CodefConnection connection) { return connection.getUserId() + ":" + connection.getProvider(); } @@ -224,6 +232,11 @@ public void upsert(Account account) { store.put(key, account); } + @Override + public void upsertAll(List accounts) { + accounts.forEach(this::upsert); + } + @Override public void updateAccountDetails(Account account) { } @@ -233,6 +246,14 @@ public Account findByConnectionIdAndAccountNoHash(Long codefConnectionId, String return store.get(key(codefConnectionId, accountNoHash)); } + @Override + public List findByConnectionIdAndAccountNoHashes(Long codefConnectionId, List accountNoHashes) { + return accountNoHashes.stream() + .map(hash -> findByConnectionIdAndAccountNoHash(codefConnectionId, hash)) + .filter(java.util.Objects::nonNull) + .toList(); + } + @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { return store.values().stream() diff --git a/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountServiceTest.java index 4c88ca0f..e346a0d4 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/VirtualAccountServiceTest.java @@ -156,6 +156,14 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return store.get(key(userId, provider)); } + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } + private static String key(Long userId, String provider) { return userId + ":" + provider; } @@ -182,6 +190,11 @@ public void upsert(Account account) { store.put(account.getId(), account); } + @Override + public void upsertAll(List accounts) { + accounts.forEach(this::upsert); + } + @Override public void updateAccountDetails(Account account) { } @@ -194,6 +207,14 @@ public Account findByConnectionIdAndAccountNoHash(Long codefConnectionId, String .findFirst().orElse(null); } + @Override + public List findByConnectionIdAndAccountNoHashes(Long codefConnectionId, List accountNoHashes) { + return accountNoHashes.stream() + .map(hash -> findByConnectionIdAndAccountNoHash(codefConnectionId, hash)) + .filter(java.util.Objects::nonNull) + .toList(); + } + @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { lastDetailProvider = provider; diff --git a/services/account-service/src/test/java/com/ntropy/account/service/VirtualConnectionServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/VirtualConnectionServiceTest.java index b86b5423..a4ad1a60 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/VirtualConnectionServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/VirtualConnectionServiceTest.java @@ -5,6 +5,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import java.util.HashMap; +import java.util.List; import java.util.Map; import org.junit.jupiter.api.Test; @@ -102,6 +103,14 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return store.get(key(userId, provider)); } + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } + private static String key(Long userId, String provider) { return userId + ":" + provider; } diff --git a/services/account-service/src/test/java/com/ntropy/account/service/VirtualFinancialDataServiceTest.java b/services/account-service/src/test/java/com/ntropy/account/service/VirtualFinancialDataServiceTest.java index 1e48d7a7..b1a41008 100644 --- a/services/account-service/src/test/java/com/ntropy/account/service/VirtualFinancialDataServiceTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/service/VirtualFinancialDataServiceTest.java @@ -309,6 +309,14 @@ public CodefConnection findByUserIdAndProvider(Long userId, String provider) { return store.get(userId + ":" + provider); } + @Override + public List findByUserIdsAndProvider(List userIds, String provider) { + return userIds.stream() + .map(userId -> findByUserIdAndProvider(userId, provider)) + .filter(java.util.Objects::nonNull) + .toList(); + } + private static String key(CodefConnection connection) { return connection.getUserId() + ":" + connection.getProvider(); } @@ -331,6 +339,11 @@ public void upsert(Account account) { store.put(key, account); } + @Override + public void upsertAll(List accounts) { + accounts.forEach(this::upsert); + } + @Override public void updateAccountDetails(Account account) { } @@ -340,6 +353,14 @@ public Account findByConnectionIdAndAccountNoHash(Long codefConnectionId, String return store.get(key(codefConnectionId, accountNoHash)); } + @Override + public List findByConnectionIdAndAccountNoHashes(Long codefConnectionId, List accountNoHashes) { + return accountNoHashes.stream() + .map(hash -> findByConnectionIdAndAccountNoHash(codefConnectionId, hash)) + .filter(java.util.Objects::nonNull) + .toList(); + } + @Override public Account findByIdAndUserIdAndProvider(Long id, Long userId, String provider) { return store.values().stream()