본문 바로가기
Spring Batch

[Spring Batch - 05] JdbcPagingItemReader와 JdbcBatchItemWriter

by cafecortado 2024. 11. 5.
reference :
[SpringBatch 연재 05] JdbcPagingItemReader로 DB내용을 읽고, JdbcBatchItemWriter로 DB에 쓰기

 

Review

지난 포스팅에서는 Spring Batch의 Chunk Processing에서 사용되는 ItemReader와 ItemWriter의 기본 구현체 중 하나인 FlatFileItemReader와 FlatFileItemWriter에 대해 알아보았다.

FlatFileItemReader와 FlatFileItemWriter는 구조화된 텍스트 파일 데이터를 처리하는 데 최적화되어 있었다.

 

이번 시간에는 데이터베이스의 데이터를 처리하는 데 최적화되어 있는 JdbcPagingItemReader와 JdbcBatchItemWriter에 대해 알아보자.

 

주요  차이점은 아래와 같다.

  FlatFileItemReader
FlatFileItemWriter
JdbcPagingItemReader
JdbcBatchItemWriter
데이터 소스 CSV, 텍스트 파일과 같이 구조화된 파일 데이터베이스
데이터 처리 파일을 한 줄씩 읽고 개별 아이템으로 처리 페이징을 통해 데이터베이스에서 데이터를 읽어와, 대용량 데이터도 메모리 효율적으로 처리
데이터 쓰기 파일에 데이터를 순차적으로 기록하는 방식, I/O 성능에 의존적 데이터베이스에 한 번에 여러 레코드를 삽입하는 Batch 기능을 지원, 쓰기 성능 향상

 

또 JdbcPagingItemReader와 JdbcBatchItemWriter는 페이지 크기와 커밋 간격을 설정하여 메모리 사용을 최적화할 수 있으며, 쿼리 최적화를 통해 성능을 높일 수 있고, 재시작 시 작업 상태를 저장하는 기능도 제공된다는 특징이 있다.

 

 

ItemReader - JdbcPagingItemReader 구현체

JdbcPagingItemReader는 데이터베이스에서 데이터를 페이지 단위로 읽어오는 ItemReader 구현체이다.

대량의 데이터를 한 번에 메모리에 로드하지 않고 페이지 크기만큼 나누어 로드할 수 있어 효율적이라는 장점이 있다.

 

아래는 Customer 클래스와 Query Provider를 작성하고, JdbcPagingItemReader를 통해 데이터베이스에서 Customer 데이터를 페이징 방식으로 읽어오는 예제이다.

 

@Data
public class Customer {
    private String name;
    private int age;
    private String gender;
}

JdbcPagingItemReader에서 데이터베이스의 레코드를 읽어올 때 매핑될 Customer 클래스를 정의한다.

 

    @Bean
    public PagingQueryProvider queryProvider() throws Exception {
        SqlPagingQueryProviderFactoryBean queryProvider = new SqlPagingQueryProviderFactoryBean();
        queryProvider.setDataSource(dataSource);
        queryProvider.setSelectClause("id, name, age, gender");
        queryProvider.setFromClause("from customer");
        queryProvider.setWhereClause("where age >= :age");

        Map<String, Order> sortKeys = new HashMap<>(1);
        sortKeys.put("id", Order.DESCENDING);

        queryProvider.setSortKeys(sortKeys);

        return queryProvider.getObject();
    }
  • SqlPagingQueryProviderFactoryBean: 페이징 쿼리를 생성하는 팩토리 역할을 하는 클래스로, 데이터베이스 종류에 맞는 PagingQueryProvider 구현체를 제공한다.
  • setDataSource: Query Provider에 데이터 소스를 설정하여 데이터베이스 종류에 맞는 적절한 페이징 방법이 적용되도록 한다.
  • setSelectClause: SELECT 절에 포함될 필드를 지정한다.
  • setFromClause: 조회할 테이블을 지정한다.
  • setWhereClause: 조건절을 지정한다.
  • setSortKeys: 결과 정렬 방식을 지정한다.

 

    @Bean
    public JdbcPagingItemReader<Customer> jdbcPagingItemReader() throws Exception {

        Map<String, Object> parameterValue = new HashMap<>();
        parameterValue.put("age", 20);

        return new JdbcPagingItemReaderBuilder<Customer>()
                .name("jdbcPagingItemReader")
                .fetchSize(CHUNK_SIZE)
                .dataSource(dataSource)
                .rowMapper(new BeanPropertyRowMapper<>(Customer.class))
                .queryProvider(queryProvider())
                .parameterValues(parameterValue)
                .build();
    }
  • name: Reader의 이름을 설정한다.
  • fetchSize: 한 번에 가져올 레코드 수(페이지 크기)를 설정한다.
  • dataSource: 데이터베이스에 연결하는 DataSource를 설정한다.
  • rowMapper: BeanPropertyRowMapper를 사용하여 Customer 클래스에 매핑될 수 있도록 한다.
  • queryProvider: 데이터 조회 쿼리를 제공한다.
  • parameterValues: 쿼리에 필요한 파라미터 값을 전달한다.

 

 

ItemWriter - JdbcBatchItemWriter 구현체

JdbcBatchItemWriter는 데이터베이스에 데이터를 Batch로 기록하는 ItemWriter 구현체이다.

대량의 데이터를 한 번에 배치로 처리하여 데이터베이스에 기록하므로, 트랜잭션의 횟수를 줄이고 성능을 높일 수 있어 효율적이다.

또 SQL 쿼리를 설정하여 데이터를 원하는 형식으로 저장할 수 있다.

 

아래는 JdbcBatchItemWriter를 사용해 Customer 객체의 데이터를 customer2 테이블에 저장하는 작업을 수행하는 예제이다.

 

spring:
  datasource:
    url: jdbc:mysql://localhost:3306/testdb?useUnicode=true&characterEncoding=utf8&clusterInstanceHostPattern=?&zeroDateTimeBehavior=CONVERT_TO_NULL&allowMultiQueries=true
    username: root
    password: root1234
  batch:
    job:
      name: JDBC_BATCH_WRITER_CHUNK_JOB

application.yaml 파일을 작성한다.

password는 본인의 것으로 수정해야 한다.

 

create table testdb.customer2
(
    id     int auto_increment primary key,
    name   varchar(100) null,
    age    int          null,
    gender varchar(10)  null
);

Customer 객체를 매핑하여 데이터를 저장할 customer2 테이블을 생성한다.

MySQL에 접속하여 실행하면 된다.

 

  @Bean
    public JdbcBatchItemWriter<Customer> flatFileItemWriter() {
        return new JdbcBatchItemWriterBuilder<Customer>()
                .dataSource(dataSource)
                .sql("INSERT INTO customer2 (name, age, gender) VALUES (?, ?, ?)")
                .itemSqlParameterSourceProvider(new CustomerItemSqlParameterSourceProvider())
                .build();
    }
  • dataSource: 데이터베이스에 연결하기 위한 DataSource를 설정한다.
  • sql: 데이터베이스에 데이터를 저장하는 SQL 쿼리이다.
  • itemSqlParameterSourceProvider: SQL 쿼리에 Customer 객체의 각 필드를 자동으로 매핑할 수 있도록 설정한다. 

 

public class CustomerItemSqlParameterSourceProvider implements ItemSqlParameterSourceProvider<Customer> {
    @Override
    public SqlParameterSource createSqlParameterSource(Customer item) {
        return new BeanPropertySqlParameterSource(item);
    }
}

 

  • createSqlParameterSource: 객체의 필드와 SQL 파라미터를 자동으로 매핑해주는 BeanPropertySqlParameterSource를 사용해 Customer 객체의 각 필드를 SQL 파라미터로 변환한다. 

 

 

전체 소스 코드

package com.spring_batch.batch_sample.jobs.jdbc;

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.database.JdbcPagingItemReader;
import org.springframework.batch.item.database.Order;
import org.springframework.batch.item.database.PagingQueryProvider;
import org.springframework.batch.item.database.builder.JdbcPagingItemReaderBuilder;
import org.springframework.batch.item.database.support.SqlPagingQueryProviderFactoryBean;
import org.springframework.batch.item.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.beans.factory.annotation.Autowired;
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.jdbc.core.BeanPropertyRowMapper;
import org.springframework.transaction.PlatformTransactionManager;

import javax.sql.DataSource;
import java.util.HashMap;
import java.util.Map;

@Slf4j
@Configuration
public class JdbcPagingReaderJobConfig {

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

    @Autowired
    DataSource dataSource;

    @Bean
    public JdbcPagingItemReader<Customer> jdbcPagingItemReader() throws Exception {

        Map<String, Object> parameterValue = new HashMap<>();
        parameterValue.put("age", 20);

        return new JdbcPagingItemReaderBuilder<Customer>()
                .name("jdbcPagingItemReader")
                .fetchSize(CHUNK_SIZE)
                .dataSource(dataSource)
                .rowMapper(new BeanPropertyRowMapper<>(Customer.class))
                .queryProvider(queryProvider())
                .parameterValues(parameterValue)
                .build();
    }

    @Bean
    public PagingQueryProvider queryProvider() throws Exception {
        SqlPagingQueryProviderFactoryBean queryProvider = new SqlPagingQueryProviderFactoryBean();
        queryProvider.setDataSource(dataSource);  // DB 에 맞는 PagingQueryProvider 를 선택하기 위함
        queryProvider.setSelectClause("id, name, age, gender");
        queryProvider.setFromClause("from customer");
        queryProvider.setWhereClause("where age >= :age");

        Map<String, Order> sortKeys = new HashMap<>(1);
        sortKeys.put("id", Order.DESCENDING);

        queryProvider.setSortKeys(sortKeys);

        return queryProvider.getObject();
    }

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


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

        return new StepBuilder("customerJdbcPagingStep", jobRepository)
                .<Customer, Customer>chunk(CHUNK_SIZE, transactionManager)
                .reader(jdbcPagingItemReader())
                .writer(customerFlatFileItemWriter())
                .build();
    }

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

 

package com.spring_batch.batch_sample.jobs.jdbc;


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.database.JdbcBatchItemWriter;
import org.springframework.batch.item.database.builder.JdbcBatchItemWriterBuilder;
import org.springframework.batch.item.file.FlatFileItemReader;
import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.ClassPathResource;
import org.springframework.transaction.PlatformTransactionManager;

import javax.sql.DataSource;

@Slf4j
@Configuration
public class JdbcBatchItemJobConfig {

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

    @Autowired
    DataSource dataSource;

    @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 JdbcBatchItemWriter<Customer> flatFileItemWriter() {

        return new JdbcBatchItemWriterBuilder<Customer>()
                .dataSource(dataSource)
                .sql("INSERT INTO customer2 (name, age, gender) VALUES (:name, :age, :gender)")
                .itemSqlParameterSourceProvider(new CustomerItemSqlParameterSourceProvider())
                .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())
                .writer(flatFileItemWriter())
                .build();
    }

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