본문 바로가기
Spring Batch

[Spring Batch - 08] CompositeItemProcessor로 데이터 Transform하기

by cafecortado 2024. 11. 30.
reference :
https://devocean.sk.com/experts/techBoardDetail.do?ID=166950
https://docs.spring.io/spring-batch/reference/processor.html
 

Item processing :: Spring Batch

The ItemReaders and ItemWriters chapter discusses multiple approaches to parsing input. Each major implementation throws an exception if it is not “well formed.” The FixedLengthTokenizer throws an exception if a range of data is missing. Similarly, att

docs.spring.io

 

 

[SpringBatch 연재 08] CompositeItemProcessor 으로 여러단계에 걸쳐 데이터 Transform하기

 

devocean.sk.com

 

CompositeItemProcessor 개념

Chunk Oriented Processing

지난 포스팅에서 언급했다시피 Spring Batch의 한 Step은 ItemReader -> (ItemProcessor) -> ItemWriter 로 이루어져 있다.

이번 포스팅에서는 ItemProcessor의 구현체인 CompositeItemProcessor에 대해 알아보려고 한다.

 

ItemProcessor는 Spring Batch에서 데이터 변환 로직을 담당하는 인터페이스로, 읽어온 데이터를 가공하거나 필터링하는 역할을 한다.

 

아래는 Spring Batch ItemProcessor 인터페이스이다.

public interface ItemProcessor<I, O> {
    O process(I item) throws Exception;
}

I 타입의 데이터를 입력받아 process 함수를 수행하여 O 타입의 데이터를 출력하는 간단한 구조이다.

 

오늘 공부할 CompositeItemProcessor는 여러 개의 ItemProcessor를 순차적으로 연결하여 구성할 수 있게 해주는 클래스이다.

각 프로세서의 출력이 다음 프로세서의 입력으로 전달되며, 이를 통해 복잡한 데이터 변환 과정을 단순화할 수 있다.

 

 

CompositeItemProcessor 예제

지난 시간에 살펴본 MyBatisItemReader에 CompositeItemProcessor를 추가해보자.

먼저 이름과 성별을 소문자로 변경하는 ItemProcessor와 나이에 20을 더하는 ItemProcessor를 구현하고, 이 둘을 연결한 CompositeItemProcessor를 추가하려고 한다.

 

LowerCaseItemProcessor 작성

package com.spring_batch.batch_sample.jobs.mybatis;

import com.spring_batch.batch_sample.jobs.models.Customer;
import org.springframework.batch.item.ItemProcessor;

/**
 * 이름, 성별을 소문자로 변경하는 ItemProcessor
 */
public class LowerCaseItemProcessor implements ItemProcessor<Customer, Customer> {
    @Override
    public Customer process(Customer item) throws Exception {
        item.setName(item.getName().toLowerCase());
        item.setGender(item.getGender().toLowerCase());
        return item;
    }
}

이름과 성별을 소문자로 변경한다.

 

After20YearsItemProcessor 작성

package com.spring_batch.batch_sample.jobs.mybatis;

import com.spring_batch.batch_sample.jobs.models.Customer;
import org.springframework.batch.item.ItemProcessor;

/**
 * 나이에 20년을 더하는 ItemProcessor
 */
public class After20YearsItemProcessor implements ItemProcessor<Customer, Customer> {
    @Override
    public Customer process(Customer item) throws Exception {
        item.setAge(item.getAge() + 20);
        return item;
    }
}

나이에 20씩 더한다.

 

CompositeItemProcessor로 연결

@Bean
    public CompositeItemProcessor<Customer, Customer> compositeItemProcessor() {
        return new CompositeItemProcessorBuilder<Customer, Customer>()
                .delegates(List.of(
                        new LowerCaseItemProcessor(),
                        new After20YearsItemProcessor()
                ))
                .build();
    }

 

CompositeItemPricessor의 delegates 메서드로 CompositeItemProcessor가 사용할 ItemProcessor들의 리스트를 지정한다.

이 때, 리스트에 추가된 순서대로 ItemProcessor가 실행된다.

 

Step에 추가

@Bean
    public Step customerJdbcCursorStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) throws Exception {
        log.info("------------------ Init customerJdbcCursorStep -----------------");

        return new StepBuilder("customerJdbcCursorStep", jobRepository)
                .<Customer, Customer>chunk(CHUNK_SIZE, transactionManager)
                .reader(myBatisItemReader())
                .processor(compositeItemProcessor())
                .writer(customerCursorFlatFileItemWriter())
                .build();
    }

정의한 Step의 processor에 compositeItemProcessor를 넣어준다.

 

전체코드

package com.spring_batch.batch_sample.jobs.mybatis;

import com.spring_batch.batch_sample.jobs.models.Customer;
import lombok.extern.slf4j.Slf4j;
import org.apache.ibatis.session.SqlSessionFactory;
import org.mybatis.spring.batch.MyBatisPagingItemReader;
import org.mybatis.spring.batch.builder.MyBatisPagingItemReaderBuilder;
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.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.batch.item.support.CompositeItemProcessor;
import org.springframework.batch.item.support.builder.CompositeItemProcessorBuilder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.FileSystemResource;
import org.springframework.transaction.PlatformTransactionManager;

import javax.sql.DataSource;
import java.util.List;

@Slf4j
@Configuration
public class MyBatisReaderJobConfig {

     // CHUNK 크기를 지정한다.
    public static final int CHUNK_SIZE = 2;
    public static final String ENCODING = "UTF-8";
    public static final String MYBATIS_CHUNK_JOB = "MYBATIS_CHUNK_JOB";

    @Autowired
    DataSource dataSource;

    @Autowired
    SqlSessionFactory sqlSessionFactory;

    @Bean
    public MyBatisPagingItemReader<Customer> myBatisItemReader() throws Exception {

        return new MyBatisPagingItemReaderBuilder<Customer>()
                .sqlSessionFactory(sqlSessionFactory)
                .pageSize(CHUNK_SIZE)
                .queryId("com.spring_batch.batch_sample.jobs.selectCustomers")
                .build();
    }


    @Bean
    public FlatFileItemWriter<Customer> customerCursorFlatFileItemWriter() {
        return new FlatFileItemWriterBuilder<Customer>()
                .name("customerCursorFlatFileItemWriter")
                .resource(new FileSystemResource("./output/customer_new_v4.csv"))
                .encoding(ENCODING)
                .delimited().delimiter(",")
                .names("Name", "Age", "Gender")
                .build();
    }

    @Bean
    public CompositeItemProcessor<Customer, Customer> compositeItemProcessor() {
        return new CompositeItemProcessorBuilder<Customer, Customer>()
                .delegates(List.of(
                        new LowerCaseItemProcessor(),
                        new After20YearsItemProcessor()
                ))
                .build();
    }

    @Bean
    public Step customerJdbcCursorStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) throws Exception {
        log.info("------------------ Init customerJdbcCursorStep -----------------");

        return new StepBuilder("customerJdbcCursorStep", jobRepository)
                .<Customer, Customer>chunk(CHUNK_SIZE, transactionManager)
                .reader(myBatisItemReader())
                .processor(compositeItemProcessor())
                .writer(customerCursorFlatFileItemWriter())
                .build();
    }

    @Bean
    public Job customerJdbcCursorPagingJob(Step customerJdbcCursorStep, JobRepository jobRepository) {
        log.info("------------------ Init customerJdbcCursorPagingJob -----------------");
        return new JobBuilder(MYBATIS_CHUNK_JOB, jobRepository)
                .incrementer(new RunIdIncrementer())
                .start(customerJdbcCursorStep)
                .build();
    }
}

 

실행결과

입력 파일
변환 결과