랭킹 페이지 구현을 위해 매일 자정에 모든 사용자들의 전날 대비 랭킹 변화 데이터를 DB에 기록하기로 했다. 이 기회에 Spring Batch를 사용해보자! 라는 생각이 바로 들었다. 프로세스는 아래와 같다.
Trigger: 매일 자정 실행 (Spring Scheduler 사용)
Step 1: User 테이블에서 현재 랭킹 정보를 가져오기
Step 2: RankingHistory 테이블에서 전날 랭킹 정보 가져오기
Step 3: 랭킹 변화 계산 및 RankingHistory, TodayRanking 저장
Step 4: 배치 실행 로그 남기기
Spring Batch 5를 기준으로 작성한 코드이다.
@Configuration
@RequiredArgsConstructor
public class BatchConfig {
private final RankingHistoryItemWriter rankingHistoryItemWriter;
private final JobRepository jobRepository;
private final PlatformTransactionManager transactionManager;
@Bean
public Job rankingHistoryJob(Step rankingHistoryStep) {
return new JobBuilder("rankingHistoryJob", jobRepository)
.incrementer(new RunIdIncrementer())
.start(rankingHistoryStep)
.build();
}
@Bean
public Step rankingHistoryStep(
ItemReader<User> rankingHistoryItemReader,
ItemProcessor<User, RankingComposite> rankingHistoryItemProcessor,
ItemWriter<RankingComposite> rankingHistoryItemWriter) {
return new StepBuilder("rankingHistoryStep", jobRepository)
.<User, RankingComposite>chunk(100, transactionManager)
.reader(rankingHistoryItemReader)
.processor(rankingHistoryItemProcessor)
.writer(rankingHistoryItemWriter)
.build();
}
}
하나하나 살펴보자.
@Bean
public Job rankingHistoryJob(Step rankingHistoryStep) {
return new JobBuilder("rankingHistoryJob", jobRepository)
.incrementer(new RunIdIncrementer()) // 매 실행마다 새로운 ID 생성
.start(rankingHistoryStep) // 첫 번째 Step 설정
.build();
}
먼저 rankingHistoryJob이라는 Job을 설정한다. rankingHistoryJob은 배치 작업의 최상위 개념으로, Job을 실행하면 Step이 실행된다. incrementer(new RunIdIncrementer()) 는 Spring Batch에서 배치 Job을 생성할 때 실행 ID를 자동 증가시키는 역할을 한다. 이 부분에 대해 이해하려면 JobParameters 에 대해 알고 있어야 한다.
즉, Spring Batch는 같은 JobParameters로 실행하면 중복 실행을 방지하기 때문에 기본적으로 이전에 실행했던 JobParameters와 동일하면 실행할 수 없다.
따라서 RunIdIncrementer 를 활용하여 매번 실행할 때마다 자동으로 run.id 값을 증가시켜, 동일한 JobParameters라도 새로운 실행으로 간주되도록 보장하는 것이다.
그렇다면 JobRepository란 무엇일까?
Spring Batch에서는 모든 Job과 Step 실행 정보를 데이터베이스에 저장한다. 이때, Job과 Step의 실행 상태(STARTED, COMPLETED, FAILED 등)를 관리하는 역할을 하는 것이 JobRepository이다.
참고로 Spring Boot에서는 spring.batch.jdbc.initialize-schema=always를 설정하면 자동으로 테이블이 생성된다. 배치 작업이 실행되면 관련 테이블에 데이터가 자동으로 기록되는 것을 확인할 수 있을 것이다.
@Bean
public Step rankingHistoryStep(
ItemReader<User> rankingHistoryItemReader,
ItemProcessor<User, RankingComposite> rankingHistoryItemProcessor,
ItemWriter<RankingComposite> rankingHistoryItemWriter) {
return new StepBuilder("rankingHistoryStep", jobRepository)
.<User, RankingComposite>chunk(100, transactionManager)
.reader(rankingHistoryItemReader)
.processor(rankingHistoryItemProcessor)
.writer(rankingHistoryItemWriter)
.build();
}
이번에는 Spring Batch의 Step을 정의하는 부분에 대해 알아보자.
Step은 Spring Batch의 작업 단위이며, 하나의 Job 안에서 실행된다. 이때 하나의 Job은 여러 개의 Step으로 구성될 수 있고 Step 내부에서 데이터를 읽고(ItemReader), 처리하고(ItemProcessor), 저장(ItemWriter)하는 작업이 진행된다.
reader(rankingHistoryItemReader)processor(rankingHistoryItemProcessor)<User, RankingComposite>writer(rankingHistoryItemWriter) chunk(100, transactionManager)@Slf4j
@Component
@RequiredArgsConstructor
public class RankingHistoryJobScheduler {
private final JobLauncher jobLauncher;
private final Job rankingHistoryJob;
@Scheduled(cron = "0 0 0 * * ?") // 매일 자정 실행
public void runRankingHistoryJob() {
log.info("Spring Batch : RankingHistory update");
try {
JobParameters jobParameters = new JobParametersBuilder()
.addLong("timestamp", System.currentTimeMillis())
.toJobParameters();
jobLauncher.run(rankingHistoryJob, jobParameters);
} catch(Exception e) {
e.printStackTrace();
}
}
}
Spring Batch를 사용하여 매일 자정(00:00:00)에 랭킹 히스토리를 갱신하는 배치 작업을 실행하는 스케줄러이다.
@Getter
@AllArgsConstructor
public class RankingComposite {
private RankingHistory rankingHistory;
private TodayRanking todayRanking;
}
RankingComposite은 Spring Batch에서 ItemProcessor와 ItemWriter를 사용할 때 중간 가공 데이터를 저장하는 역할을 한다.
즉, User 데이터를 RankingComposite 객체로 가공한 후, 최종적으로 ItemWriter에서 데이터베이스에 저장하는 구조이다.
RankingComposite을 만든 이유는 다음과 같다.
@Component
@RequiredArgsConstructor
public class RankingHistoryItemReader implements ItemReader<User> {
private final UserRepository userRepository;
private Iterator<User> userIterator;
@Override
public User read() {
if (userIterator == null) {
List<User> users = userRepository.findAll();
userIterator = users.iterator();
}
return userIterator.hasNext() ? userIterator.next() : null;
}
}
Spring Batch의 ItemReader 구현체로, 배치 작업에서 User 데이터를 읽어오는 역할을 한다. 즉, UserRepository를 이용하여 DB에서 User 데이터를 조회하고 하나씩 반환하는 구조이다.
이런 식으로 반환된 데이터는 Spring Batch의 실행 흐름에 의해 ItemProcessor가 처리한다.
⭐ Spring Batch의 실행 흐름
1️⃣ ItemReader → 데이터를 한 개씩 읽음
2️⃣ ItemProcessor → 읽은 데이터를 가공함
3️⃣ ItemWriter → 가공된 데이터를 저장함
4️⃣ (1~3 반복) → chunk(100)이면 100개씩 반복 후, 한 번에 저장
@Component
@Slf4j
public class RankingHistoryItemProcessor implements ItemProcessor<User, RankingComposite> {
private final RankingHistoryRepository rankingHistoryRepository;
private final UserRepository userRepository;
private final TodayRankingRepository todayRankingRepository;
public RankingHistoryItemProcessor(
RankingHistoryRepository rankingHistoryRepository,
TodayRankingRepository todayRankingRepository,
UserRepository userRepository) {
this.rankingHistoryRepository = rankingHistoryRepository;
this.userRepository = userRepository;
this.todayRankingRepository = todayRankingRepository;
}
@Override
public RankingComposite process(User user) {
// TodayRanking에서 전날 데이터 삭제
todayRankingRepository.deleteByUserId(user.getUserId());
// 전날 랭킹 정보 조회
RankingHistory yesterdayRanking = rankingHistoryRepository.findLatestByUser(user.getUserId());
List<User> rankedUsers = userRepository.findAllOrderedByRankingPointTierUsername();
// 랭킹 계산
int currentRanking = 1;
for (User rankedUser : rankedUsers) {
if (rankedUser.getUserId().equals(user.getUserId())) {
break;
}
currentRanking++;
}
// 변동된 순위
int rankingChange;
// 처음 가입한 유저의 경우 이전 랭킹이 없으므로 계산된 랭킹 그대로 저장
if (yesterdayRanking == null) {
rankingChange = currentRanking;
} else {
// 이전 랭킹 기록이 있는 유저의 경우 (어제 랭킹 - 현재 랭킹)
rankingChange = yesterdayRanking.getRanking()-currentRanking;
}
RankingHistory rankingHistory = RankingHistory.create(
user,
user.getTier(),
user.getRankingPoint(),
currentRanking,
rankingChange
);
TodayRanking todayRanking = TodayRanking.create(
user,
user.getTier(),
user.getRankingPoint(),
(long) rankingChange,
currentRanking
);
return new RankingComposite(rankingHistory, todayRanking);
}
}
위의 RankingHistoryItemReader에서 읽은 데이터를 가공하여 처리한다.
@Component
@RequiredArgsConstructor
public class RankingHistoryItemWriter implements ItemWriter<RankingComposite> {
private final RankingHistoryRepository rankingHistoryRepository;
private final TodayRankingRepository todayRankingRepository;
@Override
public void write(Chunk<? extends RankingComposite> chunk) throws Exception {
List<RankingHistory> histories = new ArrayList<>();
List<TodayRanking> rankings = new ArrayList<>();
for(RankingComposite composite : chunk.getItems()) {
histories.add(composite.getRankingHistory());
rankings.add(composite.getTodayRanking());
}
rankingHistoryRepository.saveAll(histories);
todayRankingRepository.saveAll(rankings);
}
}
마지막으로 RankingHistoryItemProcessor에서 가공한 데이터를 지정한 chunk 크기에 따라 한번에 데이터베이스에 저장한다. 이때 chunk는 Spring Batch에서 데이터를 묶어서 처리하는 단위이다.
RankingHistory와 TodayRanking을 개별적으로 처리하는 ItemProcessor, ItemWriter를 만들어도 되겠지만 서로 연관된 로직이기 때문에 RankingComposite를 통해 하나의 객체로 관리한 점이 포인트라고 할 수 있다.