reference:
[SpringBatch 연재 04] FlatFileItemReader로 단순 파일 읽고, FlatFileItemWriter로 파일에 쓰기
Review

지난 시간에 Spring Batch의 Step을 구성하는 Tasklet Model과 Chunk Model에 대해 알아보았다.
Tasklet은 한 Step을 단일 단계에서 처리한다.
반면 Chunk는 한 Step을 여러 단계에 걸쳐 처리하는데, 전체 데이터를 지정한 chunk size만큼씩 작은 청크 단위로 나누어 ItemReader -> ItemProcessor -> ItemWriter 과정을 반복한다. (ItemProcessor는 생략 가능)
이번 시간에 다룰 FlatFileItemReader와 FlatFileItemWriter는 Chunk model에서 ItemReader 와 ItemWriter 단계에 해당한다.
코드 관점에서 이야기하자면 FlatFileItemReader와 FlatFileItemWriter는 Spring Batch에서 제공하는 ItemReader Interface와 ItemWriter Interface의 구현체의 일종이다.
특히 이들은 비정형 데이터를 저장하는 DB보다는 구조화된 텍스트 파일을 처리하는 데 최적화되어 있는 기본 구현체들이다.
ItemReader - FlatFileItemReader 구현체
FlatFileItemReader는 ItemReader Interface의 구현체로, Spring Batch에서 제공된다.
다양한 파일 형식을 지원하며, 텍스트 파일의 각 라인을 정의된 데이터 구조로 변환해준다.
FlatFileItemReader의 주요 구성 요소는 아래와 같다.
- Resource: 읽을 파일을 지정한다.
- LineMapper: 파일의 각 라인을 데이터 객체로 변환한다. 예를 들어, CSV의 각 줄을 Customer 객체로 매핑한다. DefaultLineMapper와 같은 구현체를 사용하여 설정한다.
- LineTokenizer: 라인을 개별 데이터 조각으로 나눈다.
- FieldSetMapper: LineTokenizer로 나뉜 데이터를 객체의 속성에 매핑한다.
아래는 customer.csv 파일을 읽고, 각 라인을 쉼표로 구분해 Customer 객체로 매핑하는 FlatFileItemReader를 생성하는 예제이다.
import lombok.Getter;
import lombok.Setter;
@Getter
@Setter
public class Customer {
private String name;
private int age;
private String gender;
}
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder;
import org.springframework.core.io.ClassPathResource;
@Bean
public FlatFileItemReader<Customer> flatFileItemReader() {
return new FlatFileItemReaderBuilder<Customer>()
.name("FlatFileItemReader")
.resource(new ClassPathResource("./customer.csv"))
.encoding("UTF-8")
.delimited().delimiter(",")
.names("name", "age", "gender")
.targetType(Customer.class)
.build();
}
예제에서 확인할 수 있는 FlatFileItemReader의 구성 요소는 아래와 같다.
- Resource: ClassPathResource("./customer.csv")로 customer.csv 파일을 설정하고 있다.
- LineTokenizer: DelimitedLineTokenizer를 사용해 쉼표(,)를 구분자로 지정하여 CSV 데이터를 나누고, 각 항목을 name, age, gender 필드로 정의하고 있다.
- FieldSetMapper: BeanWrapperFieldSetMapper를 사용해 Customer 객체에 매핑하고 있다.
ItemProcessor 구현하기
FlatFileItemReader로 읽은 데이터를 처리하기 위해ItemProcessor를 구현하는 AggregateCustomerProcessor를 정의해보자.
import lombok.extern.slf4j.Slf4j;
import org.springframework.batch.item.ItemProcessor;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
public class AggregateCustomerProcessor implements ItemProcessor<Customer, Customer> {
ConcurrentHashMap<String, Integer> aggregateCustomers;
public AggregateCustomerProcessor(ConcurrentHashMap<String, Integer> aggregateCustomers) {
this.aggregateCustomers = aggregateCustomers;
}
@Override
public Customer process(Customer item) throws Exception {
aggregateCustomers.putIfAbsent("TOTAL_CUSTOMERS", 0);
aggregateCustomers.putIfAbsent("TOTAL_AGES", 0);
aggregateCustomers.put("TOTAL_CUSTOMERS", aggregateCustomers.get("TOTAL_CUSTOMERS") + 1);
aggregateCustomers.put("TOTAL_AGES", aggregateCustomers.get("TOTAL_AGES") + item.getAge());
return item;
}
}
각 아이템을 하나씩 읽고 아이템의 내용을 집계하도록 하였다.
ItemWriter - FlatFileItemWriter 구현체
FlatFileItemWriter는 ItemWriter Interface의 구현체로, Spring Batch에서 제공된다.
FlatFileItemWriter는 데이터를 파일에 출력하는 역할을 하는데, 다양한 설정을 통해 CSV, TSV 등 원하는 형식으로 데이터를 기록할 수 있으며, 헤더와 푸터도 설정할 수 있다.
FlatFileItemWriter의 주요 구성 요소는 아래와 같다.
- Resource: 출력 파일을 지정한다.
- LineAggregator: 각 데이터 객체를 문자열로 변환한다.
- HeaderCallback: 출력 파일에 헤더를 추가할 수 있도록 설정한다.
- FooterCallback: 출력 파일의 마지막에 푸터를 추가할 수 있도록 설정한다.
아래는 customer_new.csv 파일에 데이터를 탭으로 구분하여 작성하고, 헤더는 CustomerHeader로 설정하며, 푸터는 CustomerFooter를 사용해 총 고객 수와 총 나이를 기록하는 예제이다.
import org.springframework.batch.item.file.transform.LineAggregator;
public class CustomerLineAggregator implements LineAggregator<Customer> {
@Override
public String aggregate(Customer item) {
return item.getName() + "," + item.getAge();
}
}
import org.springframework.batch.item.file.FlatFileHeaderCallback;
import java.io.IOException;
import java.io.Writer;
public class CustomerHeader implements FlatFileHeaderCallback {
@Override
public void writeHeader(Writer writer) throws IOException {
writer.write("ID,AGE");
}
}
import lombok.extern.slf4j.Slf4j;
import org.springframework.batch.item.file.FlatFileFooterCallback;
import java.io.IOException;
import java.io.Writer;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
public class CustomerFooter implements FlatFileFooterCallback {
ConcurrentHashMap<String, Integer> aggregateCustomers;
public CustomerFooter(ConcurrentHashMap<String, Integer> aggregateCustomers) {
this.aggregateCustomers = aggregateCustomers;
}
@Override
public void writeFooter(Writer writer) throws IOException {
writer.write("총 고객 수: " + aggregateCustomers.get("TOTAL_CUSTOMERS"));
writer.write(System.lineSeparator());
writer.write("총 나이: " + aggregateCustomers.get("TOTAL_AGES"));
}
}
import org.springframework.batch.item.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.core.io.FileSystemResource;
@Bean
public FlatFileItemWriter<Customer> flatFileItemWriter() {
return new FlatFileItemWriterBuilder<Customer>()
.name("flatFileItemWriter")
.resource(new FileSystemResource("./output/customer_new.csv"))
.encoding("UTF-8")
.delimited().delimiter(",")
.names("Name", "Age", "Gender")
.headerCallback(new CustomerHeader())
.footerCallback(new CustomerFooter(aggregateInfos))
.build();
}
예제에서 확인할 수 있는 FlatFileItemWriter의 구성 요소는 아래와 같다.
- Resource: FileSystemResource("./output/customer_new.csv")로 파일을 customer_new.csv로 지정한다.
- LineAggregator: CustomerLineAggregator를 구현해 Customer 객체를 문자열로 변환한다.
- HeaderCallback: CustomerHeader를 구현하여 ID, AGE와 같은 헤더를 추가한다.
- FooterCallback: CustomerFooter를 구현하여 총 고객 수와 나이를 출력한다.
Step 및 Job 설정
Step은 Spring Batch의 작업 단위로, ItemReader, (ItemProcessor,) ItemWriter로 구성된다.
앞서 작성한 flatFileItemReader로 데이터를 읽고, flatFileItemWriter로 데이터를 기록하는 Step을 만들어보자.
Job은 Spring Batch의 전체 작업 단위로, 여러 Step의 조합으로 구성된다.
flatFileItemReader와 flatFileItemWriter로 이루어진 Step 하나만을 포함하는 Job을 구성해보자.
import org.springframework.batch.core.Step;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.transaction.PlatformTransactionManager;
@Bean
public Step flatFileStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) {
return new StepBuilder("flatFileStep", jobRepository)
.<Customer, Customer>chunk(100, transactionManager)
.reader(flatFileItemReader())
.writer(flatFileItemWriter())
.build();
}
import org.springframework.batch.core.Job;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.launch.support.RunIdIncrementer;
import org.springframework.batch.core.repository.JobRepository;
@Bean
public Job flatFileJob(Step flatFileStep, JobRepository jobRepository) {
return new JobBuilder("FLAT_FILE_CHUNK_JOB", jobRepository)
.incrementer(new RunIdIncrementer())
.start(flatFileStep)
.build();
}
flatFileStep이라는 Step을 정의하고, 이를 flatFileJob이라는 Job에 포함시켰다.
Job이 실행되면 flatFileStep이 실행되어 파일을 읽고 쓰는 작업이 수행된다.
전체 소스 코드
package com.spring_batch.batch_sample.jobs.flatfilereader;
import com.spring_batch.batch_sample.jobs.models.Customer;
import lombok.extern.slf4j.Slf4j;
import org.springframework.batch.item.ItemProcessor;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
public class AggregateCustomerProcessor implements ItemProcessor<Customer, Customer> {
ConcurrentHashMap<String, Integer> aggregateCustomers;
public AggregateCustomerProcessor(ConcurrentHashMap<String, Integer> aggregateCustomers) {
this.aggregateCustomers = aggregateCustomers;
}
@Override
public Customer process(Customer item) throws Exception {
aggregateCustomers.putIfAbsent("TOTAL_CUSTOMERS", 0);
aggregateCustomers.putIfAbsent("TOTAL_AGES", 0);
aggregateCustomers.put("TOTAL_CUSTOMERS", aggregateCustomers.get("TOTAL_CUSTOMERS") + 1);
aggregateCustomers.put("TOTAL_AGES", aggregateCustomers.get("TOTAL_AGES") + item.getAge());
return item;
}
}
package com.spring_batch.batch_sample.jobs.flatfilereader;
import lombok.extern.slf4j.Slf4j;
import org.springframework.batch.item.file.FlatFileFooterCallback;
import java.io.IOException;
import java.io.Writer;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
public class CustomerFooter implements FlatFileFooterCallback {
ConcurrentHashMap<String, Integer> aggregateCustomers;
public CustomerFooter(ConcurrentHashMap<String, Integer> aggregateCustomers) {
this.aggregateCustomers = aggregateCustomers;
}
@Override
public void writeFooter(Writer writer) throws IOException {
writer.write("총 고객 수: " + aggregateCustomers.get("TOTAL_CUSTOMERS"));
writer.write(System.lineSeparator());
writer.write("총 나이: " + aggregateCustomers.get("TOTAL_AGES"));
}
}
package com.spring_batch.batch_sample.jobs.flatfilereader;
import org.springframework.batch.item.file.FlatFileHeaderCallback;
import java.io.IOException;
import java.io.Writer;
public class CustomerHeader implements FlatFileHeaderCallback {
@Override
public void writeHeader(Writer writer) throws IOException {
writer.write("ID,AGE");
}
}
package com.spring_batch.batch_sample.jobs.flatfilereader;
import com.spring_batch.batch_sample.jobs.models.Customer;
import org.springframework.batch.item.file.transform.LineAggregator;
public class CustomerLineAggregator implements LineAggregator<Customer> {
@Override
public String aggregate(Customer item) {
return item.getName() + "," + item.getAge();
}
}
package com.spring_batch.batch_sample.jobs.flatfilereader;
import com.spring_batch.batch_sample.jobs.models.Customer;
import lombok.extern.slf4j.Slf4j;
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.launch.support.RunIdIncrementer;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.ClassPathResource;
import org.springframework.core.io.FileSystemResource;
import org.springframework.transaction.PlatformTransactionManager;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
@Configuration
public class FlatFileItemJobConfig {
/**
* CHUNK 크기를 지정한다.
*/
public static final int CHUNK_SIZE = 100;
public static final String ENCODING = "UTF-8";
public static final String FLAT_FILE_WRITER_CHUNK_JOB = "FLAT_FILE_WRITER_CHUNK_JOB";
private ConcurrentHashMap<String, Integer> aggregateInfos = new ConcurrentHashMap<>();
private final ItemProcessor<Customer, Customer> itemProcessor = new AggregateCustomerProcessor(aggregateInfos);
@Bean
public FlatFileItemReader<Customer> flatFileItemReader() {
return new FlatFileItemReaderBuilder<Customer>()
.name("FlatFileItemReader")
.resource(new ClassPathResource("./customer.csv"))
.encoding(ENCODING)
.delimited().delimiter(",")
.names("name", "age", "gender")
.targetType(Customer.class)
.build();
}
@Bean
public FlatFileItemWriter<Customer> flatFileItemWriter() {
return new FlatFileItemWriterBuilder<Customer>()
.name("flatFileItemWriter")
.resource(new FileSystemResource("./output/customer_new.csv"))
.encoding(ENCODING)
.delimited().delimiter(",")
.names("Name", "Age", "Gender")
.append(false)
.lineAggregator(new CustomerLineAggregator())
.headerCallback(new CustomerHeader())
.footerCallback(new CustomerFooter(aggregateInfos))
.build();
}
@Bean
public Step flatFileStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) {
log.info("------------------ Init flatFileStep -----------------");
return new StepBuilder("flatFileStep", jobRepository)
.<Customer, Customer>chunk(CHUNK_SIZE, transactionManager)
.reader(flatFileItemReader())
.processor(itemProcessor)
.writer(flatFileItemWriter())
.build();
}
@Bean
public Job flatFileJob(Step flatFileStep, JobRepository jobRepository) {
log.info("------------------ Init flatFileJob -----------------");
return new JobBuilder(FLAT_FILE_WRITER_CHUNK_JOB, jobRepository)
.incrementer(new RunIdIncrementer())
.start(flatFileStep)
.build();
}
}
package com.spring_batch.batch_sample.jobs.models;
import lombok.Getter;
import lombok.Setter;
@Setter
@Getter
public class Customer {
private String name;
private int age;
private String gender;
}
customer.csv 파일은 src/main/resources/ 경로에 넣으면 된다.

Wrap Up
FlatFileItemReader와 FlatFileItemWriter를 통해 Spring Batch를 사용하여 CSV 파일을 읽고, 이를 다른 형식의 파일로 변환하는 법을 알아보았다.
FlatFileItemReader와 FlatFileItemWriter 사이에 ItemProcessor를 구현하여 데이터를 처리하는 작업을 추가하였다.
다음 시간에는 JdbcCursorItemReader와 JdbcBatchItemWriter를 사용해서 텍스트 파일이 아니라 비정형 데이터를 저장하는 데이터베이스를 읽고 쓰는 방법을 알아보자.
'Spring Batch' 카테고리의 다른 글
| [Spring Batch - 07] MyBatisPagingItemReader와 MyBatisBatchItemWriter (3) | 2024.11.19 |
|---|---|
| [Spring Batch - 06] JpaPagingItemReader와 JpaItemWriter (1) | 2024.11.12 |
| [Spring Batch - 05] JdbcPagingItemReader와 JdbcBatchItemWriter (4) | 2024.11.05 |
| [Spring Batch - 03] ChunkModel과 TaskletModel (0) | 2024.10.22 |
| [Spring Batch - 02] Tasklet 예제로 스프링 배치의 아키텍처와 동작 알아보기 (0) | 2024.10.15 |