Spring Boot 3/Batch 5 ์ต์ ์ค์ , Job & Step ๊ตฌ์กฐ, Chunk ํ๋ก์ธ์ฑ, Paging Reader/Writer, JobParameters, Skip & Retry, ์น API ์จ๋๋งจ๋ Jar ์คํ ๋ฐ ํํฐ์
๋ ๊ฐ์ด๋
## 1. Spring Batch ์ํคํ
์ฒ & ํต์ฌ ๋๋ฉ์ธ ๊ฐ๋
๋์ฉ๋ ๋ฐ์ดํฐ(์๋ฐฑ๋ง~์์ต ๊ฑด)๋ฅผ ์์ ์ ์ด๊ณ ์ผ๊ด์ ์ผ๋ก ์ฒ๋ฆฌํ๊ธฐ ์ํด ์ค๊ณ๋ ์ํฐํ๋ผ์ด์ฆ ๋ฐฐ์น ํ๋ ์์ํฌ์
๋๋ค.
```text
[Spring Batch ๊ณ์ธต ๊ตฌ์กฐ]
Job (์ ์ฒด ๋ฐฐ์น ์์
๋จ์)
โโโ Step 1 (๋
๋ฆฝ์ ์ธ ๋ฐฐ์น ๋จ๊ณ)
โโโ Chunk-Oriented Task (Reader -> Processor -> Writer)
โโโ Step 2 (ํ์ ๋จ๊ณ ๋๋ Tasklet ๋จ์ผ ์์
)
[ํต์ฌ ๋๋ฉ์ธ ๊ฐ์ฒด]
โข JobInstance : Job ์ด๋ฆ๊ณผ ์๋ณ JobParameters์ ๊ณ ์ ํ ์กฐํฉ (๋
ผ๋ฆฌ์ ์คํ ๋จ์)
โข JobExecution : JobInstance์ 1ํ ์ค์ ์คํ ์๋ ๊ธฐ๋ก (์ฑ๊ณต, ์คํจ, ์งํ ์ค ์ํ ๋ณด์ )
โข StepExecution : ๊ฐ๋ณ Step์ ์คํ ์ํ, ์ฝ๊ธฐ/์ฐ๊ธฐ/์ปค๋ฐ/๋กค๋ฐฑ ์นด์ดํธ ์ ์ฅ
โข ExecutionContext : Step ๋๋ Job ์์ค์์ ์คํจ ๋ณต๊ตฌ๋ฅผ ์ํ ์ํ(์ฌ์์ ์ง์ ๋ฑ)๋ฅผ ๊ณต์ ํ๋ Key-Value ์ ์ฅ์
โข JobRepository : ๋ชจ๋ ๋ฐฐ์น ๋ฉํ๋ฐ์ดํฐ(์์/์ข
๋ฃ์๊ฐ, ์ํ, ์นด์ดํฐ)๋ฅผ DB์ CRUDํ๋ ์ ์ฅ์
```
### ๐๏ธ ๋ฉํ๋ฐ์ดํฐ ํต์ฌ ํ
์ด๋ธ 6๊ฐ
```sql
-- ๋ฐฐ์น ์คํ ์ด๋ ฅ์ ๊ด๋ฆฌํ๋ ํ์ ํ
์ด๋ธ (Spring Boot ์์ ์ ์๋ ์์ฑ ๊ฐ๋ฅ)
BATCH_JOB_INSTANCE -- Job ์ด๋ฆ๊ณผ ํ๋ผ๋ฏธํฐ ํด์ ํค ๋งคํ (๋์ผ ํ๋ผ๋ฏธํฐ ์ฌ์คํ ๋ฐฉ์ง)
BATCH_JOB_EXECUTION -- Job ์คํ ์ํ, ์์/์ข
๋ฃ ์ผ์, ExitCode
BATCH_JOB_EXECUTION_PARAMS -- ์ ๋ฌ๋ ํ๋ผ๋ฏธํฐ ํ์
๋ฐ ๊ฐ ๋ชฉ๋ก
BATCH_JOB_EXECUTION_CONTEXT -- Job ์์ค์ ExecutionContext (์ง๋ ฌํ๋ ์ํ๊ฐ)
BATCH_STEP_EXECUTION -- Step ๋จ์ ์ฑ๊ณต/์คํจ, READ_COUNT, WRITE_COUNT, COMMIT_COUNT
BATCH_STEP_EXECUTION_CONTEXT-- Step ์์ค ์คํจ ๋ณต๊ตฌ ์ง์ ๋ฐ์ดํฐ
```
## 2. Spring Boot 3 (Batch 5) ์ต์ ์ค์ & ๋ง์ด๊ทธ๋ ์ด์
Spring Boot 3.x์์๋ ๊ธฐ์กด Spring Batch 4์ ํฉํ ๋ฆฌ ํด๋์ค(`JobBuilderFactory`, `StepBuilderFactory`)๊ฐ **์์ ์ ๊ฑฐ(Deprecated/Removed)**๋์์ต๋๋ค.
```java
package com.lucky.batch.config;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.transaction.PlatformTransactionManager;
@Configuration
// โ ๏ธ ์ฃผ์: Spring Boot 3์์๋ @EnableBatchProcessing์ ๋ถ์ด์ง ๋ง์ธ์!
// ๋ถ์ด๋ฉด ์๋ ๊ตฌ์ฑ(BatchAutoConfiguration)์ด ๋นํ์ฑํ๋์ด Runner ๋ฑ๋ก ๋ฑ์ด ๋ฒ๊ฑฐ๋ก์์ง๋๋ค.
public class BatchConfig {
// ๐ [Spring Batch 5 ์ต์ ํ์ค ๋ฌธ๋ฒ]
// JobRepository์ PlatformTransactionManager๋ฅผ ์ง์ ์ฃผ์
๋ฐ์ ๋น๋๋ฅผ ์์ฑํฉ๋๋ค.
@Bean
public Job sampleJob(JobRepository jobRepository, Step sampleStep) {
return new JobBuilder("sampleJob", jobRepository)
.start(sampleStep)
.build();
}
@Bean
public Step sampleStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) {
return new StepBuilder("sampleStep", jobRepository)
.tasklet((contribution, chunkContext) -> {
System.out.println("Spring Batch 5 Step ์คํ ์๋ฃ!");
return org.springframework.batch.repeat.RepeatStatus.FINISHED;
}, transactionManager)
.build();
}
}
```
### `application.yml` ํ์ ์ค์
```yaml
spring:
batch:
job:
name: ${job.name:NONE} # ํน์ Job๋ง ๊ณจ๋ผ์ ์คํ (์ง์ ์ ํ๋ฉด ์ ์ฒด Job ์คํ ๋ฐฉ์ง)
jdbc:
initialize-schema: always # DB์ ๋ฉํ๋ฐ์ดํฐ ํ
์ด๋ธ ์๋ ์์ฑ (always | embedded | never)
```
## 3. ์ฒญํฌ ์งํฅ ํ๋ก์ธ์ฑ (Chunk-Oriented Processing)
๋์ฉ๋ ๋ฐ์ดํฐ๋ฅผ ์ฒญํฌ(Chunk, ๋ฉ์ด๋ฆฌ) ๋จ์๋ก ํธ๋์ญ์
์ ๋ถํ ์ปค๋ฐํ์ฌ OOM(Out of Memory)์ ๋ฐฉ์งํ๊ณ ๋น ๋ฅธ ์ฑ๋ฅ์ ๋ณด์ฅํฉ๋๋ค.
```java
@Bean
public Step chunkStep(
JobRepository jobRepository,
PlatformTransactionManager transactionManager,
ItemReader<User> userReader,
ItemProcessor<User, DormantUser> userProcessor,
ItemWriter<DormantUser> userWriter) {
// ๐ฆ ์ฒญํฌ ํฌ๊ธฐ 1000 ์ค์ : 1000๊ฐ ์ฝ๊ณ -> ๊ฐ๊ณต ํ -> ํ ๋ฒ์ ํธ๋์ญ์
์ปค๋ฐ
return new StepBuilder("dormantUserStep", jobRepository)
.<User, DormantUser>chunk(1000, transactionManager)
.reader(userReader)
.processor(userProcessor)
.writer(userWriter)
.build();
}
/*
* ๐ก [ํฉ๊ธ ๊ท์น: ChunkSize == PageSize]
* ItemReader์ pageSize์ ChunkSize๋ ๋ฐ๋์ ๋์ผํ๊ฒ ๋ง์ถ๋ ๊ฒ์ด ์ ์์
๋๋ค.
* PageSize(100) != ChunkSize(1000) ์ผ ๊ฒฝ์ฐ:
* ์ฒญํฌ 1๊ฐ๋ฅผ ์ฑ์ฐ๊ธฐ ์ํด ์ฟผ๋ฆฌ๊ฐ 10๋ฒ ์คํ๋์ด ๋นํจ์จ์ด ๋ฐ์ํฉ๋๋ค.
*/
```
## 4. ๋ค์ํ DB ItemReader ๊ตฌํ์ฒด (Paging Reader)
๋ฉ๋ชจ๋ฆฌ ๋ถํ๋ฅผ ์ค์ด๊ธฐ ์ํด ์ปค์(Cursor) ๋ฐฉ์๋ณด๋ค **ํ์ด์ง(Paging) ๋ฐฉ์**์ ์๋์ ์ผ๋ก ์ ํธํฉ๋๋ค.
```java
// ๐ [1] JpaPagingItemReader (JPA / Hibernate ํ๊ฒฝ)
@Bean
@StepScope
public JpaPagingItemReader<User> jpaUserReader(EntityManagerFactory emf) {
Map<String, Object> params = new HashMap<>();
params.put("status", UserStatus.ACTIVE);
return new JpaPagingItemReaderBuilder<User>()
.name("jpaUserReader")
.entityManagerFactory(emf)
.pageSize(1000)
.queryString("SELECT u FROM User u WHERE u.status = :status ORDER BY u.id ASC")
.parameterValues(params)
.build();
}
// โก [2] JdbcPagingItemReader (์์ JDBC / ์ด๊ณ ์ ๋์ฉ๋)
@Bean
@StepScope
public JdbcPagingItemReader<User> jdbcUserReader(
DataSource dataSource,
PagingQueryProvider queryProvider) {
return new JdbcPagingItemReaderBuilder<User>()
.name("jdbcUserReader")
.dataSource(dataSource)
.pageSize(1000)
.queryProvider(queryProvider)
.rowMapper(new BeanPropertyRowMapper<>(User.class))
.build();
}
// JdbcPaging์ฉ QueryProvider (DB ๋ฐฉ์ธ๋ณ ํ์ด์ง SQL ์๋ ์์ฑ)
@Bean
public SqlPagingQueryProviderFactoryBean queryProvider(DataSource dataSource) {
SqlPagingQueryProviderFactoryBean provider = new SqlPagingQueryProviderFactoryBean();
provider.setDataSource(dataSource);
provider.setSelectClause("id, username, email, status");
provider.setFromClause("from users");
provider.setWhereClause("where status = 'ACTIVE'");
// โ ๏ธ ์ ๋ ฌ ํค(SortKeys)๋ ๊ณ ์ ์๋ณ์(PK)๋ก ๋ฐ๋์ ์ง์ ํด์ผ ๋๋ฝ/์ค๋ณต์ด ์์ต๋๋ค.
Map<String, Order> sortKeys = new HashMap<>();
sortKeys.put("id", Order.ASCENDING);
provider.setSortKeys(sortKeys);
return provider;
}
```
## 5. ItemProcessor ๋ฐ์ดํฐ ๊ฐ๊ณต & ํํฐ๋ง ํจํด
๋น์ฆ๋์ค ๋ก์ง์ ์ ์ฉํ์ฌ ๋ฐ์ดํฐ๋ฅผ ๋ณํํ๊ฑฐ๋, ์กฐ๊ฑด์ ๋ง์ง ์๋ ๋ถํ์ํ ๋ฐ์ดํฐ๋ฅผ ํํฐ๋งํฉ๋๋ค.
```java
// ๐ [1] ๋จ์ผ ๋ฐ์ดํฐ ๋ณํ ๋ฐ ํํฐ๋ง
@Bean
public ItemProcessor<User, DormantUser> dormantUserProcessor() {
return user -> {
// 1. ํน์ ์กฐ๊ฑด ํํฐ๋ง: null์ ๋ฐํํ๋ฉด ItemWriter๋ก ์ ๋ฌ๋์ง ์๊ณ ์คํต๋จ!
if (user.isVip()) {
return null; // VIP ํ์์ ํด๋ฉด ์ฒ๋ฆฌ ๋์์์ ์ ์ธ
}
// 2. ๊ฐ๊ณต ๋ฐ ์ Entity ๋ณํ
return DormantUser.builder()
.userId(user.getId())
.dormantAt(LocalDateTime.now())
.reason("1๋
์ด์ ๋ฏธ์ ์")
.build();
};
}
// ๐ [2] ๋ณตํฉ ํ๋ก์ธ์ (CompositeItemProcessor: ํ์ดํ๋ผ์ธ ์ฒด์ด๋)
@Bean
public CompositeItemProcessor<User, FinalDto> compositeProcessor() {
CompositeItemProcessor<User, FinalDto> processor = new CompositeItemProcessor<>();
processor.setDelegates(Arrays.asList(
new ValidationProcessor(), // 1๋จ๊ณ ์ ํจ์ฑ ๊ฒ์ฌ
new MaskingProcessor(), // 2๋จ๊ณ ๊ฐ์ธ์ ๋ณด ๋ง์คํน
new TransformationProcessor()// 3๋จ๊ณ ์ต์ข
DTO ๋ณํ
));
return processor;
}
```
## 6. ๋ค์ํ ItemWriter ๊ตฌํ์ฒด & ๊ณ ์ ๋ฒํฌ ์ฐ์ฐ
๊ฐ๊ณต๋ ์ฒญํฌ ๋จ์์ ์ปฌ๋ ์
์ ๋ฐ์ดํฐ๋ฒ ์ด์ค์ ํ ๋ฒ์ ์ผ๊ด ์์ํํฉ๋๋ค.
```java
// ๐ [1] JpaItemWriter (JPA ๊ธฐ๋ฐ ์๋ merge)
@Bean
public JpaItemWriter<DormantUser> jpaUserWriter(EntityManagerFactory emf) {
return new JpaItemWriterBuilder<DormantUser>()
.entityManagerFactory(emf)
.build();
}
// โก [2] JdbcBatchItemWriter (์ด๊ณ ์ Batch Insert / Update - ๊ฐ์ฅ ๊ฐ๋ ฅ ์ถ์ฒ!)
@Bean
public JdbcBatchItemWriter<DormantUser> jdbcBatchWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<DormantUser>()
.dataSource(dataSource)
.sql("INSERT INTO dormant_users (user_id, dormant_at, reason) VALUES (:userId, :dormantAt, :reason)")
.beanMapped() // DTO ๊ฐ์ฒด์ getter ์ด๋ฆ์ ๋ค์๋ ํ๋ผ๋ฏธํฐ(:userId)์ ์๋ ๋ฐ์ธ๋ฉ
.build();
}
// ๐ [3] FlatFileItemWriter (CSV ํ์ผ ์์ฑ)
@Bean
public FlatFileItemWriter<ReportDto> csvWriter() {
return new FlatFileItemWriterBuilder<ReportDto>()
.name("csvWriter")
.resource(new FileSystemResource("output/report.csv"))
.delimited()
.delimiter(",")
.names("id", "username", "amount")
.build();
}
```
## 7. JobParameters & @StepScope ์ง์ฐ ๋ก๋ฉ
๋ฐฐ์น ์คํ ์ ๋์ ํ๋ผ๋ฏธํฐ๋ฅผ ๋๊ธฐ๊ณ , ์คํ ์์ ์ Bean์ ์์ฑํ์ฌ ์์ ํ๊ฒ ํ๋ผ๋ฏธํฐ๋ฅผ ์ฃผ์
๋ฐ์ต๋๋ค.
```java
// ๐ฏ [1] SpEL์ ํ์ฉํ JobParameters ํ๋ผ๋ฏธํฐ ๋ฐ์ธ๋ฉ
// โ ๏ธ ์ฃผ์: JobParameters๋ฅผ ์ฃผ์
๋ฐ์ผ๋ ค๋ฉด ํด๋น Bean์ ๋ฐ๋์ @StepScope๊ฐ ์ ์ธ๋์ด์ผ ํฉ๋๋ค!
@Bean
@StepScope
public JpaPagingItemReader<Order> orderReader(
EntityManagerFactory emf,
@Value("#{jobParameters['requestDate']}") String requestDate,
@Value("#{jobParameters['amount']}") Long minAmount) {
return new JpaPagingItemReaderBuilder<Order>()
.name("orderReader")
.entityManagerFactory(emf)
.queryString("SELECT o FROM Order o WHERE o.orderDate = :reqDate AND o.amount >= :amount")
.parameterValues(Map.of("reqDate", requestDate, "amount", minAmount))
.pageSize(500)
.build();
}
// ๐ [2] RunIdIncrementer: ๋์ผ ํ๋ผ๋ฏธํฐ๋ก ๋งค๋ฒ ์๋ก์ด JobInstance ์์ฑ ํ์ฉ
@Bean
public Job dailySettlementJob(JobRepository jobRepository, Step settlementStep) {
return new JobBuilder("dailySettlementJob", jobRepository)
.incrementer(new RunIdIncrementer()) // run.id ํ๋ผ๋ฏธํฐ๋ฅผ 1์ฉ ์ฆ๊ฐ์์ผ ์๋ ์ ๋ฌ
.start(settlementStep)
.build();
}
```
## 8. ์์ธ ์ฒ๋ฆฌ, ์ฌ์๋ & ๊ฑด๋๋ฐ๊ธฐ (Fault Tolerance)
์ผ๋ถ ๋ฐ์ดํฐ์ ์๋ฌ๊ฐ ๋ฐ์ํด๋ ์ ์ฒด ๋ฐฐ์น๊ฐ ์ค๋จ๋์ง ์๋๋ก ๊ฑด๋๋ฐ๊ฑฐ๋(Skip) ์ฌ์๋(Retry)ํฉ๋๋ค.
```java
@Bean
public Step tolerantStep(
JobRepository jobRepository,
PlatformTransactionManager tm,
ItemReader<User> reader,
ItemProcessor<User, Target> processor,
ItemWriter<Target> writer) {
return new StepBuilder("tolerantStep", jobRepository)
.<User, Target>chunk(500, tm)
.reader(reader)
.processor(processor)
.writer(writer)
// ๐ก๏ธ ๋ด๊ฒฐํจ์ฑ(Fault Tolerant) ์ค์ ํ์ฑํ
.faultTolerant()
// โญ๏ธ [Skip] ํน์ ์์ธ ๋ฐ์ ์ ์ต๋ 10๊ฑด๊น์ง ๊ฑด๋๋ฐ๊ณ ๊ณ์ ์งํ
.skip(IllegalArgumentException.class)
.skip(DataFormatException.class)
.skipLimit(10)
.noSkip(FatalDatabaseException.class) // ์น๋ช
์ DB ์ค๋ฅ๋ ๊ฑด๋๋ฐ์ง ์๊ณ ์ฆ์ ์ค๋จ
// ๐ [Retry] ์ผ์์ ์ฅ์ (๋ฐ๋๋ฝ, ๋คํธ์ํฌ) ๋ฐ์ ์ ์ต๋ 3ํ ์ฌ์๋
.retry(DeadlockLoserDataAccessException.class)
.retry(TransientDataAccessException.class)
.retryLimit(3)
.build();
}
```
## 9. ๋ค์ค ์คํ
ํ๋ฆ ์ ์ด (Flow & Decider)
Step์ ์ฑ๊ณต/์คํจ ๊ฒฐ๊ณผ ์ฝ๋(ExitStatus)์ ๋ฐ๋ผ ๋ค์ ์์
๊ฒฝ๋ก๋ฅผ ์ ์ฐํ๊ฒ ๋ถ๊ธฐํฉ๋๋ค.
```java
@Bean
public Job conditionalJob(
JobRepository jobRepository,
Step mainStep,
Step successStep,
Step alertStep) {
return new JobBuilder("conditionalJob", jobRepository)
.start(mainStep)
.on("FAILED") // mainStep ์คํจ ์
.to(alertStep) // ๊ด๋ฆฌ์ ์ฌ๋ ์๋ฆผ Step ์คํ
.on("*") // ์๋ฆผ ํ
.end() // ์ข
๋ฃ
.from(mainStep)
.on("*") // mainStep ์ฑ๊ณต ์
.to(successStep) // ์ ์ ํ์ Step ์คํ
.next(completeStep())
.end()
.build();
}
// ๐ ์ปค์คํ
๋ถ๊ธฐ ํ๋ณ์ (JobExecutionDecider)
public class OddEvenDecider implements JobExecutionDecider {
@Override
public FlowExecutionStatus decide(JobExecution jobExecution, StepExecution stepExecution) {
int count = getTodayCount();
return (count % 2 == 0) ? new FlowExecutionStatus("EVEN") : new FlowExecutionStatus("ODD");
}
}
```
## 10. ๋์ฉ๋ ์ฑ๋ฅ ์ต์ ํ: ๋ฉํฐ์ค๋ ๋ & ํํฐ์
๋
๋จ์ผ ์ค๋ ๋์ ์ฒ๋ฆฌ ํ๊ณ๋ฅผ ๋์ด CPU ์์์ ๊ทน๋ํํ๋ ๋ณ๋ ฌ ์ฒ๋ฆฌ ๊ธฐ์ ์
๋๋ค.
```java
// ๐ [๋ฐฉ๋ฒ 1] Multi-Threaded Step (์ฒญํฌ ์ฒ๋ฆฌ๋ฅผ ์ฌ๋ฌ ์ค๋ ๋๊ฐ ๋ถ๋ด)
@Bean
public Step multiThreadedStep(
JobRepository jobRepository,
PlatformTransactionManager tm,
PagingItemReader<User> reader,
ItemWriter<User> writer) {
return new StepBuilder("multiThreadedStep", jobRepository)
.<User, User>chunk(1000, tm)
.reader(reader) // โ ๏ธ ์ฃผ์: PagingReader๋ Thread-Safeํ์ง๋ง CursorReader๋ Thread-Safeํ์ง ์์!
.writer(writer)
.taskExecutor(batchTaskExecutor()) // ๋ฉํฐ์ค๋ ๋ ํ ์ฐ๊ฒฐ
.throttleLimit(8) // ๋์ ์คํ ์ค๋ ๋ ์
.build();
}
@Bean
public TaskExecutor batchTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(8);
executor.setMaxPoolSize(16);
executor.setThreadNamePrefix("batch-thread-");
executor.initialize();
return executor;
}
// ๐งฉ [๋ฐฉ๋ฒ 2] Partitioning (๋ฐ์ดํฐ๋ฅผ ๋ฒ์๋ณ๋ก ๋ฌผ๋ฆฌ ๋ถํ ํ์ฌ ๋
๋ฆฝ Worker Step ๋ณ๋ ฌ ์คํ)
@Bean
public Step masterStep(JobRepository jobRepository, Step workerStep) {
return new StepBuilder("masterStep", jobRepository)
.partitioner("workerStep", new IdRangePartitioner()) // ID 1~10๋ง, 10๋ง~20๋ง ๋ถํ ๊ธฐ
.step(workerStep)
.gridSize(4) // 4๊ฐ ํํฐ์
์ผ๋ก ๋ณ๋ ฌ ์ฒ๋ฆฌ
.taskExecutor(batchTaskExecutor())
.build();
}
```
## 11. ๋ฐฐ์น ๋ฆฌ์ค๋ (Listeners) & ๋ชจ๋ํฐ๋ง
๋ฐฐ์น ์ ํ ๋ก๊น
, ์ฒ๋ฆฌ ์๊ฐ ์ธก์ , ์คํจ ์๋ฆผ(Slack Webhook)์ ๋ฆฌ์ค๋๋ก ์บก์ํํฉ๋๋ค.
```java
// ๐ข [1] Job ์คํ ์๋ช
์ฃผ๊ธฐ ๋ฆฌ์ค๋
public class JobPerformanceListener implements JobExecutionListener {
private static final Logger log = LoggerFactory.getLogger(JobPerformanceListener.class);
@Override
public void beforeJob(JobExecution jobExecution) {
log.info("โถ ๋ฐฐ์น Job [{}] ์์ ์ผ์: {}", jobExecution.getJobInstance().getJobName(), LocalDateTime.now());
}
@Override
public void afterJob(JobExecution jobExecution) {
long duration = Duration.between(jobExecution.getStartTime(), jobExecution.getEndTime()).toMillis();
if (jobExecution.getStatus() == BatchStatus.COMPLETED) {
log.info("โ ๋ฐฐ์น ์๋ฃ! ์์ ์๊ฐ: {}ms", duration);
} else {
log.error("โ ๋ฐฐ์น ์คํจ! ExitStatus: {}", jobExecution.getExitStatus().getExitDescription());
sendSlackAlert(jobExecution); // ์ฌ๋ ์๋ฆผ ๋ฐ์ก
}
}
}
// ๐ฏ Step ๋น๋์ ๋ฆฌ์ค๋ ๋ฑ๋ก
// .listener(new JobPerformanceListener())
// .listener(new CustomItemWriteListener())
```
## 12. ์ค์ ํธ๋ฌ๋ธ์ํ
& ํต์ฌ ์ฃผ์์ฌํญ
์ค๋ฌด์์ ๊ฐ์ฅ ๋น๋ฒํ๊ฒ ๋ฐ์ํ๋ Spring Batch ์ค๋ฅ์ ํด๊ฒฐ์ฑ
์
๋๋ค.
```text
[์์ฃผ ๊ฒช๋ ์ค์ & ํด๊ฒฐ์ฑ
]
1. PagingItemReader Update ์ ๋ฐ์ดํฐ ๋๋ฝ (Off-by-one ๋ฒ๊ทธ)
- ์์ธ: status='WAIT'๋ฅผ ์ฝ์ด์ 'COMPLETE'๋ก UPDATEํ๋ฉด, 2ํ์ด์ง ์กฐํ ์ ์์ ๋ฐ์ดํฐ๊ฐ ๋น ์ ธ ๊ฑด๋๋ ๋ฐ์
- ํด๊ฒฐ์ฑ
: ํ์ด์ง ๋ฒํธ๋ฅผ ํญ์ 0์ผ๋ก ๊ณ ์ ํ๊ฑฐ๋(Zero-Paging), ID ๊ธฐ์ค์ผ๋ก ์ปค์ ๋ฒ์๋ฅผ ์ขํ๊ฐ๋ฉฐ ์กฐํ
2. OOM (Out Of Memory) ๋ฐ์
- ์์ธ: CursorItemReader ์ฌ์ฉ ์ fetchSize ๋ฏธ์ง์ , ๋๋ Chunk ์ฒ๋ฆฌ ์ค List ์ปฌ๋ ์
์ ์ฐธ์กฐ ์ ์ง
- ํด๊ฒฐ์ฑ
: PagingItemReader ์ฌ์ฉ ๋ฐ ChunkSize ์ผ์น, ๋ถํ์ํ Entity 1์ฐจ ์บ์ clear (EntityManager.clear)
3. ๋์ผ JobParameter ์ฌ์คํ ๋ถ๊ฐ (JobInstanceAlreadyCompleteException)
- ์์ธ: ์ฑ๊ณตํ JobInstance๋ ๋์ผ ํ๋ผ๋ฏธํฐ๋ก ๋ค์ ์คํ๋์ง ์์ (๋ฐฐ์น ์์ ์ฑ ์ค๊ณ)
- ํด๊ฒฐ์ฑ
: ์คํ ์๋ง๋ค ๊ณ ์ ํ๋ผ๋ฏธํฐ(์: time=System.currentTimeMillis()) ์ถ๊ฐ ๋๋ RunIdIncrementer ์ ์ฉ
4. Reader/Writer์์ @StepScope ๋๋ฝ
- ์์ธ: JobParameters๋ StepExecutionContext๋ฅผ SpEL(#{...})๋ก ์ฃผ์
๋ฐ์ ๋ @StepScope๊ฐ ์์ผ๋ฉด ์ ํ๋ฆฌ์ผ์ด์
๊ธฐ๋ ์์ ๋ฐ์ธ๋ฉ ์คํจ ์๋ฌ ๋ฐ์
- ํด๊ฒฐ์ฑ
: ๋์ ํ๋ผ๋ฏธํฐ๊ฐ ํ์ํ ๋ชจ๋ Reader, Processor, Writer Bean์ @StepScope ๋ช
์
```
## 13. ์น ์๋น์ค(REST API)์์ ์จ๋๋งจ๋ ๋ฐฐ์น Jar ๋์ ์คํ ํจํด
์น ์๋ฒ(Spring Boot Web)๋ ์์ ๊ฐ๋ํ๊ณ , HTTP ์์ฒญ์ด ์ธ์
๋ ๋๋ง **์ธ๋ถ ๋ฐฐ์น Jar๋ฅผ ๋
๋ฆฝ OS ํ๋ก์ธ์ค(์ JVM)๋ก ์คํ**ํ ๋ค ์ข
๋ฃ ์ ์์์ 100% ๋ฐํํ๋ ์ค๋ฌด ํ์ค ์ํคํ
์ฒ์
๋๋ค.
```text
[์น ์๋ฒ] โโ(POST /api/batch/run)โโ> [BatchLauncherService (@Async)]
โ
โผ ProcessBuilder
[java -jar my-batch.jar] (๋
๋ฆฝ JVM)
โ
โผ ์์
์๋ฃ ์
[ํ๋ก์ธ์ค ์๋ ์ข
๋ฃ & ๋ฉ๋ชจ๋ฆฌ ๋ฐํ]
```
```java
// ๐ [1] ๋น๋๊ธฐ ๋ฐฐ์น Jar ๋ฐ์ฒ ์๋น์ค (ProcessBuilder)
package com.lucky.web.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import java.io.BufferedReader;
import java.io.File;
import java.io.InputStreamReader;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
@Service
public class ExternalBatchLauncherService {
private static final Logger log = LoggerFactory.getLogger(ExternalBatchLauncherService.class);
// ๋์ ์คํ ๊ฐ๋ฅํ ์ต๋ ๋ฐฐ์น ํ๋ก์ธ์ค ์ ์ ํ (์๋ฒ CPU/RAM ๊ณ ๊ฐ ๋ฐฉ์ง)
private final Semaphore semaphore = new Semaphore(2);
@Async("batchTaskExecutor") // ๋ณ๋ ์ ์ฉ ์ค๋ ๋ ํ
public CompletableFuture<Integer> runBatchJar(String jobName, String requestDate) {
if (!semaphore.tryAcquire()) {
log.warn("์ด๋ฏธ ์ต๋ ์คํ ๊ฐ๋ฅํ ๋ฐฐ์น ์๊ฐ ๊ฐ๋ ์ค์
๋๋ค.");
return CompletableFuture.failedFuture(new IllegalStateException("์๋ฒ ์ฉ๋ ์ด๊ณผ: ์ ์ ํ ๋ค์ ์๋ํด์ฃผ์ธ์."));
}
try {
List<String> command = new ArrayList<>();
command.add("java");
command.add("-Xms512m");
command.add("-Xmx2g"); // ๋ฐฐ์น ์ ์ฉ ๋ฉ๋ชจ๋ฆฌ ํ ๋น (์น ์๋ฒ์ ์์ ๊ฒฉ๋ฆฌ)
command.add("-jar");
command.add("/app/batch/my-batch-app.jar"); // ์ธ๋ถ JAR ๊ฒฝ๋ก
// Spring Batch ์ ์ฉ ์คํ ์ธ์ ์ ๋ฌ
command.add("--spring.batch.job.name=" + jobName);
command.add("requestDate=" + requestDate);
command.add("run.id=" + System.currentTimeMillis()); // ์ค๋ณต ์คํ ๋ฐฉ์ง
ProcessBuilder pb = new ProcessBuilder(command);
pb.directory(new File("/app/batch"));
pb.redirectErrorStream(true); // ์๋ฌ ์คํธ๋ฆผ์ ํ์ค ์ถ๋ ฅ๊ณผ ๋ณํฉ
Process process = pb.start();
long pid = process.pid();
log.info("๋ฐฐ์น ์๋ธ ํ๋ก์ธ์ค ์์ฑ ์๋ฃ (PID: {})", pid);
// ์ค์๊ฐ ๋ก๊ทธ ์์ง
try (BufferedReader reader = new BufferedReader(new InputStreamReader(process.getInputStream()))) {
String line;
while ((line = reader.readLine()) != null) {
log.info("[BATCH-LOG:{}] {}", pid, line);
}
}
// ์ต๋ 30๋ถ ๋๊ธฐ ํ ์ด๊ณผ ์ ํ๋ก์ธ์ค ๊ฐ์ ์ข
๋ฃ
boolean finished = process.waitFor(30, TimeUnit.MINUTES);
if (!finished) {
process.destroyForcibly();
log.error("๋ฐฐ์น ์คํ ์๊ฐ ์ด๊ณผ๋ก ๊ฐ์ ์ข
๋ฃ (PID: {})", pid);
return CompletableFuture.completedFuture(-1);
}
int exitCode = process.exitValue();
log.info("๋ฐฐ์น ํ๋ก์ธ์ค ์ ์ ์ข
๋ฃ (PID: {}, ExitCode: {})", pid, exitCode);
return CompletableFuture.completedFuture(exitCode);
} catch (Exception e) {
log.error("๋ฐฐ์น ์คํ ์ค ์์ธ ๋ฐ์", e);
return CompletableFuture.failedFuture(e);
} finally {
semaphore.release();
}
}
}
```
```java
// ๐ [2] REST ์ปจํธ๋กค๋ฌ (ํด๋ผ์ด์ธํธ ์ฆ์ ์๋ต ๋ฐํ)
@RestController
@RequestMapping("/api/batch")
public class BatchTriggerController {
private final ExternalBatchLauncherService launcherService;
public BatchTriggerController(ExternalBatchLauncherService launcherService) {
this.launcherService = launcherService;
}
@PostMapping("/run")
public ResponseEntity<?> trigger(@RequestParam String jobName, @RequestParam String date) {
// HTTP ์์ฒญ ์ค๋ ๋๋ฅผ ๋ธ๋กํนํ์ง ์๊ณ ๋น๋๊ธฐ ์ ์ ํ ์ฆ์ ์๋ต
launcherService.runBatchJar(jobName, date);
return ResponseEntity.ok(Map.of(
"status", "ACCEPTED",
"message", "๋ฐฐ์น ์คํ ์์
์ด ๋ฐฑ๊ทธ๋ผ์ด๋์ ๋ฑ๋ก๋์์ต๋๋ค.",
"jobName", jobName,
"date", date
));
}
}
```
```text
[์คํ ๋ฐฉ์ ๋น๊ต ๊ฐ์ด๋]
โข ProcessBuilder (์ธ๋ถ Jar ์คํ) : ์น/๋ฐฐ์น JVM ์์ ๋ถ๋ฆฌ, OOM ๋ฐ์ ์ ์น ์๋ฒ ์ํฅ ์ ๋ก, ๋ฐฐ์น ์ข
๋ฃ ์ RAM 100% ํ์ (๊ฐ๋ ฅ ๊ถ์ฅ)
โข JobLauncher (๋์ผ JVM ๋ด ์คํ) : ๊ฐ์ ํ๋ก์ ํธ์ผ ๋ ๊ฐํธํ๋ ๋์ฉ๋ ๋ฐฐ์น ์ ์น ์๋ฒ ์ ์ฒด OOM ๋ค์ด ์ํ ์กด์ฌ
โข Kubernetes Job (์ปจํ
์ด๋ ๋ฐฐํฌ) : ํด๋ผ์ฐ๋ ํ๊ฒฝ์์ HTTP ์์ฒญ ์ ์ผํ์ฑ Pod ์์ฑ ํ ์ข
๋ฃํ๋ ์์ ์๋ฒ๋ฆฌ์ค ํจํด
```
์๊ฒฌ ๋ฐ ์ง๋ฌธ
0์์ง ๋ฑ๋ก๋ ์๊ฒฌ์ด ์์ต๋๋ค. ์ฒซ ๋ฒ์งธ ๋๊ธ์ ๋จ๊ฒจ๋ณด์ธ์!
๋๊ธ ์์
๋๊ธ ์ญ์