Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import com.ntropy.account.domain.InstitutionKeys;
import com.ntropy.account.domain.PersonalBank;
import com.ntropy.account.domain.entity.CodefConnection;
import com.ntropy.account.event.AccountTransactionClassificationEventPublisher;
import com.ntropy.account.exception.AccountErrorCode;
import com.ntropy.account.mapper.AccountLifecycleMapper;
import com.ntropy.account.mapper.AccountMapper;
Expand All @@ -20,7 +21,6 @@
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;
Expand All @@ -45,7 +45,7 @@ public class LocalFinancialAccountCommandClient implements FinancialAccountComma
private final AccountLifecycleMapper accountLifecycleMapper;
private final AccountMapper accountMapper;
private final CodefConnectionMapper codefConnectionMapper;
private final TransactionClassificationCommandClient transactionClassificationCommandClient;
private final AccountTransactionClassificationEventPublisher classificationEventPublisher;

@Override
public List<BankSummary> findSupportedBanks() {
Expand Down Expand Up @@ -82,7 +82,7 @@ public AccountRegistrationSummary registerAccount(Long userId, AccountRegistrati

if ("VIRTUAL".equals(connectionType)) {
GenerationSummary summary = virtualAccountRegenerationService.regenerateForUser(userId, bank);
classifyTransactionsSafely(userId);
requestTransactionClassificationSafely(userId);
return new AccountRegistrationSummary(connectionType, bank.getOrganizationCode(), summary.accounts());
}

Expand All @@ -97,16 +97,18 @@ public AccountRegistrationSummary registerAccount(Long userId, AccountRegistrati
);
ensureVirtualDatasetSafely(userId, bank);
int accountCount = accountCollectionService.collect(userId, bank, birthDate, startDate, endDate).size();
classifyTransactionsSafely(userId);
requestTransactionClassificationSafely(userId);
return new AccountRegistrationSummary(connectionType, bank.getOrganizationCode(), accountCount);
}

private void classifyTransactionsSafely(Long userId) {
private void requestTransactionClassificationSafely(Long userId) {
try {
int processed = transactionClassificationCommandClient.classifyUnanalyzedTransactions(userId);
log.info("계좌 연동 후 소비 분류 완료: userId={}, processed={}", userId, processed);
classificationEventPublisher.publishAfterCommit(userId);
} catch (RuntimeException e) {
log.warn("계좌 연동 후 소비 분류 실패: userId={}", userId, e);
log.warn(
"계좌 연동 후 소비 분류 작업 등록 실패: userId={}, errorKind={}",
userId, e.getClass().getSimpleName()
);
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
package com.ntropy.account.config;

import java.util.concurrent.ThreadPoolExecutor;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.PropertySource;
import org.springframework.context.annotation.PropertySources;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

/** 계좌 연동 요청과 소비 분류 작업이 서로의 스레드 풀을 점유하지 않도록 분리합니다. */
@Configuration
@PropertySources({
@PropertySource(
value = "classpath:account-classification.properties",
ignoreResourceNotFound = true
),
@PropertySource(
value = "file:${NTROPY_CONFIG_DIR:./config}/account-classification.properties",
ignoreResourceNotFound = true
)
})
public class AccountClassificationAsyncConfig {

@Bean("accountClassificationJobExecutor")
public ThreadPoolTaskExecutor accountClassificationJobExecutor(
@Value("${account.classification.async.core-pool-size:1}") int configuredCorePoolSize,
@Value("${account.classification.async.max-pool-size:2}") int configuredMaxPoolSize,
@Value("${account.classification.async.queue-capacity:100}") int configuredQueueCapacity
) {
int corePoolSize = Math.max(1, Math.min(configuredCorePoolSize, 8));
int maxPoolSize = Math.max(corePoolSize, Math.min(configuredMaxPoolSize, 8));
int queueCapacity = Math.max(1, Math.min(configuredQueueCapacity, 1000));

ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(corePoolSize);
executor.setMaxPoolSize(maxPoolSize);
executor.setQueueCapacity(queueCapacity);
executor.setThreadNamePrefix("account-classification-job-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());
return executor;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package com.ntropy.account.event;

import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Component;
import org.springframework.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager;

import lombok.RequiredArgsConstructor;

/** 거래 저장 트랜잭션이 있으면 커밋된 뒤에만 소비 분류 이벤트를 발행합니다. */
@Component
@RequiredArgsConstructor
public class AccountTransactionClassificationEventPublisher {

private final ApplicationEventPublisher applicationEventPublisher;

public void publishAfterCommit(Long userId) {
AccountTransactionsCollectedEvent event =
new AccountTransactionsCollectedEvent(userId);

if (TransactionSynchronizationManager.isActualTransactionActive()
&& TransactionSynchronizationManager.isSynchronizationActive()) {
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronization() {
@Override
public void afterCommit() {
applicationEventPublisher.publishEvent(event);
}
}
);
return;
}

/* 현재 계좌 조합 경로처럼 하위 저장 메서드가 이미 커밋된 경우입니다. */
applicationEventPublisher.publishEvent(event);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
package com.ntropy.account.event;

import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;

import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.event.EventListener;
import org.springframework.core.task.TaskExecutor;
import org.springframework.stereotype.Component;

import com.ntropy.common.client.TransactionClassificationCommandClient;

import lombok.extern.slf4j.Slf4j;

/** 계좌 연동 응답과 분리된 전용 executor에서 미분류 거래를 처리합니다. */
@Component
@Slf4j
public class AccountTransactionClassificationListener {

private final TransactionClassificationCommandClient classificationCommandClient;
private final TaskExecutor classificationJobExecutor;
private final Set<Long> runningUserIds = ConcurrentHashMap.newKeySet();

public AccountTransactionClassificationListener(
TransactionClassificationCommandClient classificationCommandClient,
@Qualifier("accountClassificationJobExecutor") TaskExecutor classificationJobExecutor
) {
this.classificationCommandClient = classificationCommandClient;
this.classificationJobExecutor = classificationJobExecutor;
}

@EventListener
public void onTransactionsCollected(AccountTransactionsCollectedEvent event) {
Long userId = event.userId();
if (!runningUserIds.add(userId)) {
log.info("[비동기 소비 분류] 중복 작업 건너뜀. scope=userId={}", userId);
return;
}

try {
classificationJobExecutor.execute(() -> classify(userId));
} catch (RuntimeException e) {
runningUserIds.remove(userId);
log.warn(
"[비동기 소비 분류] 작업 등록 실패. scope=userId={}, errorKind={}",
userId, e.getClass().getSimpleName()
);
}
}

private void classify(Long userId) {
long startedAt = System.nanoTime();
log.info("[비동기 소비 분류] 작업 시작. scope=userId={}", userId);

try {
int processed = classificationCommandClient.classifyUnanalyzedTransactions(userId);
log.info(
"[비동기 소비 분류] 작업 완료. scope=userId={}, totalProcessed={}, elapsedMs={}",
userId, processed, elapsedMillis(startedAt)
);
} catch (RuntimeException e) {
log.warn(
"[비동기 소비 분류] 작업 실패. scope=userId={}, elapsedMs={}, errorKind={}",
userId, elapsedMillis(startedAt), e.getClass().getSimpleName()
);
} finally {
runningUserIds.remove(userId);
}
}

private static long elapsedMillis(long startedAt) {
return (System.nanoTime() - startedAt) / 1_000_000L;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
package com.ntropy.account.event;

/** 계좌 연동으로 거래 저장을 마친 뒤 소비 분류를 요청하는 내부 이벤트입니다. */
public record AccountTransactionsCollectedEvent(Long userId) {
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
# 계좌 연동 응답과 분리해 실행하는 사용자 단위 소비 분류 작업 executor 설정
account.classification.async.core-pool-size=1
account.classification.async.max-pool-size=2
account.classification.async.queue-capacity=100
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,15 @@
import java.util.List;

import org.junit.jupiter.api.Test;
import org.springframework.context.ApplicationEventPublisher;

import com.ntropy.account.domain.ConnectionProvider;
import com.ntropy.account.domain.PersonalBank;
import com.ntropy.account.domain.entity.Account;
import com.ntropy.account.domain.entity.CodefConnection;
import com.ntropy.account.exception.AccountErrorCode;
import com.ntropy.account.event.AccountTransactionClassificationEventPublisher;
import com.ntropy.account.event.AccountTransactionsCollectedEvent;
import com.ntropy.account.mapper.AccountLifecycleMapper;
import com.ntropy.account.mapper.AccountMapper;
import com.ntropy.account.mapper.CodefConnectionMapper;
Expand All @@ -25,7 +28,6 @@
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 {
Expand Down Expand Up @@ -137,41 +139,43 @@ collectionService, new StubVirtualAccountRegenerationService(),
}

@Test
void classifiesUsersTransactionsAfterCodefCollectionCompletes() {
StubTransactionClassificationCommandClient classificationClient =
new StubTransactionClassificationCommandClient();
void requestsUsersTransactionClassificationAfterCodefCollectionCompletes() {
StubApplicationEventPublisher eventPublisher = new StubApplicationEventPublisher();
LocalFinancialAccountCommandClient client = newClient(
new StubPersonalBankAccountService(), new StubAccountCollectionService(),
new StubVirtualAccountRegenerationService(), new StubVirtualFinancialDataService(),
new StubAccountLifecycleMapper(1, 1), new StubAccountMapper(false),
new StubCodefConnectionMapper(), classificationClient
new StubCodefConnectionMapper(),
new AccountTransactionClassificationEventPublisher(eventPublisher)
);

client.registerAccount(
42L, new AccountRegistrationCommand("CODEF", "0088", "bank-id", "bank-password", null)
);

assertEquals(List.of(42L), classificationClient.userIds);
assertEquals(
List.of(new AccountTransactionsCollectedEvent(42L)),
eventPublisher.events
);
}

@Test
void keepsAccountRegistrationSuccessfulWhenClassificationFails() {
StubTransactionClassificationCommandClient classificationClient =
new StubTransactionClassificationCommandClient();
classificationClient.failure = new IllegalStateException("분류 실패");
void keepsAccountRegistrationSuccessfulWhenClassificationRequestFails() {
StubApplicationEventPublisher eventPublisher = new StubApplicationEventPublisher();
eventPublisher.failure = new IllegalStateException("작업 등록 실패");
LocalFinancialAccountCommandClient client = newClient(
new StubPersonalBankAccountService(), new StubAccountCollectionService(),
new StubVirtualAccountRegenerationService(), new StubVirtualFinancialDataService(),
new StubAccountLifecycleMapper(1, 1), new StubAccountMapper(false),
new StubCodefConnectionMapper(), classificationClient
new StubCodefConnectionMapper(),
new AccountTransactionClassificationEventPublisher(eventPublisher)
);

var result = client.registerAccount(
42L, new AccountRegistrationCommand("VIRTUAL", "0088", null, null, null)
);

assertEquals("VIRTUAL", result.connectionType());
assertEquals(List.of(42L), classificationClient.userIds);
}

@Test
Expand Down Expand Up @@ -461,7 +465,9 @@ private static LocalFinancialAccountCommandClient newClient(
return newClient(
personalBankAccountService, collectionService, regenerationService,
virtualFinancialDataService, lifecycleMapper, accountMapper, connectionMapper,
new StubTransactionClassificationCommandClient()
new AccountTransactionClassificationEventPublisher(
new StubApplicationEventPublisher()
)
);
}

Expand All @@ -473,26 +479,24 @@ private static LocalFinancialAccountCommandClient newClient(
AccountLifecycleMapper lifecycleMapper,
AccountMapper accountMapper,
CodefConnectionMapper connectionMapper,
TransactionClassificationCommandClient classificationClient
AccountTransactionClassificationEventPublisher classificationEventPublisher
) {
return new LocalFinancialAccountCommandClient(
personalBankAccountService, collectionService, regenerationService, virtualFinancialDataService,
lifecycleMapper, accountMapper, connectionMapper, classificationClient
lifecycleMapper, accountMapper, connectionMapper, classificationEventPublisher
);
}

private static class StubTransactionClassificationCommandClient
implements TransactionClassificationCommandClient {
private final List<Long> userIds = new ArrayList<>();
private static class StubApplicationEventPublisher implements ApplicationEventPublisher {
private final List<Object> events = new ArrayList<>();
private RuntimeException failure;

@Override
public int classifyUnanalyzedTransactions(Long userId) {
userIds.add(userId);
public void publishEvent(Object event) {
if (failure != null) {
throw failure;
}
return 0;
events.add(event);
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package com.ntropy.account.config;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;

import java.util.concurrent.ThreadPoolExecutor;

import org.junit.jupiter.api.Test;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

class AccountClassificationAsyncConfigTest {

private final AccountClassificationAsyncConfig config =
new AccountClassificationAsyncConfig();

@Test
void usesConfiguredPoolAndQueueSizes() {
ThreadPoolTaskExecutor executor =
config.accountClassificationJobExecutor(2, 4, 25);
executor.initialize();

assertEquals(2, executor.getCorePoolSize());
assertEquals(4, executor.getMaxPoolSize());
assertEquals(25, executor.getQueueCapacity());
assertInstanceOf(
ThreadPoolExecutor.AbortPolicy.class,
executor.getThreadPoolExecutor().getRejectedExecutionHandler()
);

executor.shutdown();
}

@Test
void clampsUnsafeConfigurationValues() {
ThreadPoolTaskExecutor executor =
config.accountClassificationJobExecutor(0, 99, 0);
executor.initialize();

assertEquals(1, executor.getCorePoolSize());
assertEquals(8, executor.getMaxPoolSize());
assertEquals(1, executor.getQueueCapacity());

executor.shutdown();
}
}
Loading
Loading