Spring Batch 활용하기

pitseleh·2025년 2월 26일
post-thumbnail

랭킹 페이지 구현을 위해 매일 자정에 모든 사용자들의 전날 대비 랭킹 변화 데이터를 DB에 기록하기로 했다. 이 기회에 Spring Batch를 사용해보자! 라는 생각이 바로 들었다. 프로세스는 아래와 같다.

Trigger : 매일 자정 실행 (Spring Scheduler 사용)
Step 1 : User 테이블에서 현재 랭킹 정보를 가져오기
Step 2 : RankingHistory 테이블에서 전날 랭킹 정보 가져오기
Step 3 : 랭킹 변화 계산 및 RankingHistory, TodayRanking 저장
Step 4 : 배치 실행 로그 남기기

Spring Batch 5를 기준으로 작성한 코드이다.

1. Config

@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 에 대해 알고 있어야 한다.

🔍 JobParameters란?

  • Spring Batch에서 배치 Job 실행 시 전달되는 입력값
  • 배치 Job의 실행 이력을 JobRepository에 저장할 때 키 역할을 함
  • 같은 JobParameters로 실행하면 중복 실행을 방지 (같은 JobParameters는 한 번만 실행 가능)

즉, Spring Batch는 같은 JobParameters로 실행하면 중복 실행을 방지하기 때문에 기본적으로 이전에 실행했던 JobParameters와 동일하면 실행할 수 없다.
따라서 RunIdIncrementer 를 활용하여 매번 실행할 때마다 자동으로 run.id 값을 증가시켜, 동일한 JobParameters라도 새로운 실행으로 간주되도록 보장하는 것이다.

그렇다면 JobRepository란 무엇일까?

🔍 JobRepository란?

Spring Batch에서는 모든 Job과 Step 실행 정보를 데이터베이스에 저장한다. 이때, Job과 Step의 실행 상태(STARTED, COMPLETED, FAILED 등)를 관리하는 역할을 하는 것이 JobRepository이다.

  • Job 실행 정보 저장
    • JobInstance, JobExecution, StepExecution 등의 실행 정보를 저장
  • 현재 실행 중인 Job 조회
    • 이미 실행된 Job과 Step의 상태를 확인하고, 동일한 JobParameters로 실행되었는지 검사
  • 실패한 Job 재시작 가능
    • 이전 실행 이력을 활용하여 실패한 Job을 다시 실행할 수 있도록 관리

참고로 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이란?

Step은 Spring Batch의 작업 단위이며, 하나의 Job 안에서 실행된다. 이때 하나의 Job은 여러 개의 Step으로 구성될 수 있고 Step 내부에서 데이터를 읽고(ItemReader), 처리하고(ItemProcessor), 저장(ItemWriter)하는 작업이 진행된다.

  1. ItemReader 설정
  • reader(rankingHistoryItemReader)
  • ItemReader를 사용하여 데이터를 읽어옴
  1. ItemProcessor 설정
  • processor(rankingHistoryItemProcessor)
  • 읽은 데이터를 가공함
  • <User, RankingComposite>
    • 입력 데이터 타입 : User (ItemReader에서 읽는 데이터 타입)
    • 출력 데이터 타입 : RankingComposite (ItemProcessor에서 변환 후 반환하는 타입)
  1. ItemWriter 설정
  • writer(rankingHistoryItemWriter)
  • 가공한 데이터를 DB에 저장하는 역할
  1. Chunk 설정
  • chunk(100, transactionManager)
  • 100개씩 데이터를 처리 (100개 단위로 트랜잭션을 관리)
  • transactionManager를 사용하여 Step 실행을 트랜잭션 단위로 처리

2. Scheduler

@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)에 랭킹 히스토리를 갱신하는 배치 작업을 실행하는 스케줄러이다.

  • Spring Batch의 JobLauncher를 사용하여 rankingHistoryJob 실행
  • 위에서 설명한 것처럼 Spring Batch는 같은 JobParameters로 실행하면 중복 실행을 방지하는데, 이를 피하기 위해 현재 시간(System.currentTimeMillis())을 추가

3. Job

@Getter
@AllArgsConstructor
public class RankingComposite {

	private RankingHistory rankingHistory;
	private TodayRanking todayRanking;
}

RankingComposite은 Spring Batch에서 ItemProcessor와 ItemWriter를 사용할 때 중간 가공 데이터를 저장하는 역할을 한다.
즉, User 데이터를 RankingComposite 객체로 가공한 후, 최종적으로 ItemWriter에서 데이터베이스에 저장하는 구조이다.

RankingComposite을 만든 이유는 다음과 같다.

  • User 객체 하나로 두 개의 랭킹 데이터(RankingHistory, TodayRanking)를 관리해야 함
  • ItemProcessor에서 User를 변환하여 두 개의 데이터 객체를 하나로 묶어 처리하는 구조가 필요 → ItemWriter에서 RankingComposite을 받아서 RankingHistory, TodayRanking을 각각 데이터베이스에 저장

@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 데이터를 조회하고 하나씩 반환하는 구조이다.

  • 처음 read()가 호출될 경우
    • userIterator가 null이므로 userRepository.findAll()로 User 목록을 조회
    • userIterator를 List의 Iterator로 초기화
  • 그다음부터는 userIterator.next()를 반환
    • userIterator.hasNext()가 true인 경우, 다음 User 객체 반환
    • false가 되면 null을 반환하여 Step 종료

이런 식으로 반환된 데이터는 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를 통해 하나의 객체로 관리한 점이 포인트라고 할 수 있다.

0개의 댓글