공공부하자개발 · 영어 학습 노트
자바
실무 배치대용량 데이터 처리0/6 완료
  • 01배치 프로세스 개념과 아키텍처
  • 02대용량 파일 I/O (NIO, Buffered)
  • 03Chunk 단위 처리와 OOM 방지
  • 04트랜잭션: 커밋과 롤백 시뮬레이션
  • 05Skip과 Retry 로직 구현
  • 06메서드 활용 패턴 (배치 유틸)
사이트 소개개인정보처리방침연락처
© 2026 공부하자
홈 › 실무 배치 › 01 / 6

배치 프로세스 개념과 아키텍처

섹션 7진행 0 / 6
1왜 배우는가2핵심 원리3코드 예제4응용 변형 예제5자주 하는 실수 (Tip)6연습 문제7정리‹ 이전다음 ›

4. 응용 변형 예제

3절의 Step 은 "읽기 기준 청크, 실패하면 Job 중단, 로그는 summary 한 줄" 이라는 하나의 정책을 고정해 두었습니다. 여기서는 같은 Reader-Processor-Writer 골격을 유지한 채 청크 기준·Processor 합성·리스너·실패 정책·청크 크기를 바꿔 보며, 이 구조가 왜 유연한지 확인합니다. 각 프로그램은 필요한 인터페이스를 최소 형태로 다시 선언해 단독 실행됩니다.

변형 1: 청크 기준을 파라미터로 — 읽기 기준 vs 쓰기 기준 나란히

연습 문제 3 에서 WriteSizedStep 을 따로 만들었지만, 청크 경계를 "읽은 건수로 셀지, 쓸 건수로 셀지"는 enum 하나로 한 Step 안에서 고를 수 있습니다. 필터 비율이 높은 Step(미납 추출처럼 75% 가 걸러지는 경우)에서 두 정책의 커밋 횟수와 청크 크기가 어떻게 달라지는지 같은 입력으로 비교합니다. 쓰기 효율이 중요하면 WRITE_SIZED, 실패 시 재처리 범위를 좁히려면 READ_SIZED 입니다.

java
import java.util.*;
import java.util.function.*;
import java.util.stream.*;

public class ReadSizedVsWriteSized {
    interface ItemReader<T> { T read() throws Exception; static <T> ItemReader<T> of(List<T> l) { var it = l.iterator(); return () -> it.hasNext() ? it.next() : null; } }
    interface ItemProcessor<I, O> { O process(I item) throws Exception; }
    interface ItemWriter<T> { void write(List<? extends T> items) throws Exception; }
    record StepExecution(String stepName, long readCount, long filterCount, long writeCount, long commitCount) {}

    enum ChunkMode { READ_SIZED, WRITE_SIZED }

    // 청크 경계를 "읽은 건수" 로 셀지 "쓸 건수" 로 셀지를 파라미터로 받는 Step
    static <I, O> StepExecution execute(String name, ItemReader<I> reader, ItemProcessor<I, O> processor,
                                        ItemWriter<O> writer, int chunkSize, ChunkMode mode) throws Exception {
        long read = 0, filtered = 0, written = 0, commits = 0;
        List<O> chunk = new ArrayList<>(chunkSize);
        int readInChunk = 0;
        I item;
        while ((item = reader.read()) != null) {
            read++; readInChunk++;
            O out = processor.process(item);
            if (out == null) filtered++; else chunk.add(out);
            boolean full = mode == ChunkMode.READ_SIZED ? readInChunk == chunkSize : chunk.size() == chunkSize;
            if (full) {
                if (!chunk.isEmpty()) { writer.write(chunk); written += chunk.size(); commits++; }
                chunk.clear(); readInChunk = 0;
            }
        }
        if (!chunk.isEmpty()) { writer.write(chunk); written += chunk.size(); commits++; }
        return new StepExecution(name, read, filtered, written, commits);
    }

    record Invoice(long id, boolean paid) {}

    public static void main(String[] args) throws Exception {
        List<Invoice> invoices = IntStream.rangeClosed(1, 1000).mapToObj(i -> new Invoice(i, i % 4 != 0)).toList();   // 25% 미납
        ItemProcessor<Invoice, Long> unpaidOnly = inv -> inv.paid() ? null : inv.id();

        for (ChunkMode mode : ChunkMode.values()) {
            List<Integer> chunkSizes = new ArrayList<>();
            StepExecution exec = execute("extractUnpaid", ItemReader.of(invoices), unpaidOnly,
                    chunk -> chunkSizes.add(chunk.size()), 100, mode);
            System.out.println(mode + ": " + exec);
            System.out.println("  청크별 write 건수 = " + chunkSizes);
        }
    }
}
// 출력:
// READ_SIZED: StepExecution[stepName=extractUnpaid, readCount=1000, filterCount=750, writeCount=250, commitCount=10]
//   청크별 write 건수 = [25, 25, 25, 25, 25, 25, 25, 25, 25, 25]
// WRITE_SIZED: StepExecution[stepName=extractUnpaid, readCount=1000, filterCount=750, writeCount=250, commitCount=3]
//   청크별 write 건수 = [100, 100, 50]

변형 2: Processor 합성 — 검증 → 변환 → 필터 → 계산 체인

예제 5 의 Processor 는 람다 하나에 검증·변환·필터가 뒤섞여 있었습니다.

ItemProcessor 에 andThen default 메서드를 추가하면 작은 Processor 를 이어 붙여 파이프라인을 만들 수 있고, 앞 단계가 null(필터)을 돌려주면 뒤 단계는 실행되지 않으므로 "필터 = null" 규약이 체인 전체에 자연스럽게 전파됩니다.

단계마다 타입이 바뀌어도(RawOrder → Order → SettlementLine) 제네릭이 따라갑니다. Spring Batch 의 CompositeItemProcessor 가 바로 이것입니다.

java
import java.util.*;

public class ProcessorChain {
    @FunctionalInterface
    interface ItemProcessor<I, O> {
        O process(I item) throws Exception;

        // 합성: 앞 단계가 null(필터) 이면 뒤 단계는 실행하지 않는다
        default <R> ItemProcessor<I, R> andThen(ItemProcessor<? super O, ? extends R> next) {
            return item -> { O mid = process(item); return mid == null ? null : next.process(mid); };
        }
        static <T> ItemProcessor<T, T> filter(java.util.function.Predicate<? super T> keep) {
            return item -> keep.test(item) ? item : null;
        }
    }

    record RawOrder(String id, String amountText, String status) {}
    record Order(String id, long amount, String status) {}
    record SettlementLine(String orderId, long fee, long payout) {}

    // 단계 1: 검증 — 잘못된 행은 예외가 아니라 null 로 걸러낸다(05 레슨에서 스킵/데드 레터로 발전)
    static final ItemProcessor<RawOrder, RawOrder> VALIDATE = r ->
        r.id().startsWith("ORD-") && r.amountText().matches("\\d+") ? r : null;
    // 단계 2: 변환
    static final ItemProcessor<RawOrder, Order> PARSE = r -> new Order(r.id(), Long.parseLong(r.amountText()), r.status());
    // 단계 3: 필터 — 결제 완료만 정산
    static final ItemProcessor<Order, Order> PAID_ONLY = ItemProcessor.filter(o -> o.status().equals("PAID"));
    // 단계 4: 정산 계산 (수수료 3%)
    static final ItemProcessor<Order, SettlementLine> SETTLE = o -> new SettlementLine(o.id(), o.amount() * 3 / 100, o.amount() * 97 / 100);

    public static void main(String[] args) throws Exception {
        ItemProcessor<RawOrder, SettlementLine> pipeline = VALIDATE.andThen(PARSE).andThen(PAID_ONLY).andThen(SETTLE);

        List<RawOrder> input = List.of(
            new RawOrder("ORD-1", "10000", "PAID"),
            new RawOrder("X-2", "5000", "PAID"),          // 검증 실패
            new RawOrder("ORD-3", "12a00", "PAID"),       // 검증 실패 (숫자 아님)
            new RawOrder("ORD-4", "7000", "CANCELLED"),   // 필터
            new RawOrder("ORD-5", "20000", "PAID"));

        int filtered = 0;
        for (RawOrder r : input) {
            SettlementLine out = pipeline.process(r);
            if (out == null) { filtered++; System.out.println(r.id() + " → (filtered)"); }
            else System.out.println(r.id() + " → " + out);
        }
        System.out.println("filtered=" + filtered + " / " + input.size());
    }
}
// 출력:
// ORD-1 → SettlementLine[orderId=ORD-1, fee=300, payout=9700]
// X-2 → (filtered)
// ORD-3 → (filtered)
// ORD-4 → (filtered)
// ORD-5 → SettlementLine[orderId=ORD-5, fee=600, payout=19400]
// filtered=3 / 5

변형 3: 리스너로 관심사 분리 — beforeChunk / afterChunk / onError

예제 3 의 Step 은 실행이 끝난 뒤 StepExecution 하나만 돌려주므로, "청크마다 금액 합계를 찍고 싶다"거나 "실패한 청크 번호를 알림으로 보내고 싶다"면 Step 코드를 고쳐야 합니다. 청크 생명주기 훅을 StepListener 인터페이스로 빼면 로깅·메트릭·알림을 Step 밖에서 꽂을 수 있습니다.

Spring Batch 의 ChunkListener, StepExecutionListener 와 같은 발상입니다.

java
import java.util.*;
import java.util.stream.*;

public class ListenerStep {
    interface ItemReader<T> extends AutoCloseable {
        T read() throws Exception;
        default void open() throws Exception {}
        @Override default void close() throws Exception {}
    }
    interface ItemProcessor<I, O> { O process(I item) throws Exception; }
    interface ItemWriter<T> { void write(List<? extends T> items) throws Exception; }

    // 청크 생명주기 훅. 로깅·모니터링·메트릭을 Step 코드에 섞지 않고 밖에서 꽂는다
    interface StepListener<O> {
        default void beforeChunk(int chunkNo) {}
        default void afterChunk(int chunkNo, List<? extends O> written) {}
        default void onError(int chunkNo, Exception e) {}
        default void afterStep(String status, long read, long written) {}
    }

    static <I, O> void execute(ItemReader<I> reader, ItemProcessor<I, O> processor, ItemWriter<O> writer,
                               int chunkSize, StepListener<O> listener) {
        long read = 0, written = 0; int chunkNo = 0;
        try {
            reader.open();
            List<O> chunk = new ArrayList<>(chunkSize);
            boolean done = false;
            while (!done) {
                chunk.clear(); chunkNo++;
                listener.beforeChunk(chunkNo);
                for (int i = 0; i < chunkSize; i++) {
                    I item = reader.read();
                    if (item == null) { done = true; break; }
                    read++;
                    O out = processor.process(item);
                    if (out != null) chunk.add(out);
                }
                if (!chunk.isEmpty()) { writer.write(chunk); written += chunk.size(); listener.afterChunk(chunkNo, chunk); }
            }
            listener.afterStep("COMPLETED", read, written);
        } catch (Exception e) {
            listener.onError(chunkNo, e);
            listener.afterStep("FAILED", read, written);
        } finally {
            try { reader.close(); } catch (Exception ignored) {}
        }
    }

    record Payment(int id, long amount) {}

    public static void main(String[] args) {
        // 리스너 1: 청크 로그 + 금액 합계 메트릭
        long[] totalAmount = {0};
        StepListener<Payment> metrics = new StepListener<>() {
            public void afterChunk(int no, List<? extends Payment> w) {
                long sum = w.stream().mapToLong(Payment::amount).sum();
                totalAmount[0] += sum;
                System.out.printf("  chunk#%d 커밋 %d건, 금액 %,d%n", no, w.size(), sum);
            }
            public void onError(int no, Exception e) { System.out.println("  chunk#" + no + " 실패: " + e.getMessage()); }
            public void afterStep(String status, long read, long written) {
                System.out.printf("  step %s read=%d written=%d 누적 금액 %,d%n", status, read, written, totalAmount[0]);
            }
        };

        List<Payment> payments = IntStream.rangeClosed(1, 7).mapToObj(i -> new Payment(i, i * 1000L)).toList();
        Iterator<Payment> it1 = payments.iterator();
        System.out.println("정상 실행:");
        execute(() -> it1.hasNext() ? it1.next() : null, p -> p, c -> {}, 3, metrics);

        totalAmount[0] = 0;
        Iterator<Payment> it2 = payments.iterator();
        System.out.println("id=5 에서 예외:");
        execute(() -> it2.hasNext() ? it2.next() : null,
                p -> { if (p.id() == 5) throw new IllegalStateException("id=5 금액 검증 실패"); return p; },
                c -> {}, 3, metrics);
    }
}
// 출력:
// 정상 실행:
//   chunk#1 커밋 3건, 금액 6,000
//   chunk#2 커밋 3건, 금액 15,000
//   chunk#3 커밋 1건, 금액 7,000
//   step COMPLETED read=7 written=7 누적 금액 28,000
// id=5 에서 예외:
//   chunk#1 커밋 3건, 금액 6,000
//   chunk#2 실패: id=5 금액 검증 실패
//   step FAILED read=5 written=3 누적 금액 6,000

변형 4: Step 별 실패 정책 — STOP vs CONTINUE

예제 4 의 JobRunner 는 어떤 Step 이든 실패하면 무조건 중단합니다. 그러나 "알림 발송 실패"가 "정산 파일 생성"을 막을 이유는 없고, 반대로 "거래 집계 실패"는 반드시 이후를 막아야 합니다. Step 마다 OnFail 정책을 두면 이 차이를 선언적으로 표현할 수 있습니다. 어느 쪽이든 실패가 하나라도 있으면 Job 최종 상태는 FAILED 로 남겨 운영자가 아침에 반드시 보게 합니다.

java
import java.util.*;
import java.util.function.Supplier;

public class FailPolicyJob {
    enum Status { COMPLETED, FAILED, SKIPPED }
    enum OnFail { STOP, CONTINUE }                       // Step 실패 시 Job 을 멈출지, 다음 Step 으로 갈지

    record StepDef(String name, Supplier<Status> body, OnFail onFail) {}
    record StepResult(String name, Status status) {}

    // Step 을 순서대로 실행. STOP 정책의 Step 이 실패하면 남은 Step 은 SKIPPED 로 기록하고 Job 은 FAILED.
    // CONTINUE 정책의 Step 이 실패하면 기록만 남기고 계속 진행하되, Job 최종 상태는 FAILED 로 표시한다.
    static List<StepResult> run(String jobName, List<StepDef> steps) {
        List<StepResult> results = new ArrayList<>();
        boolean stopped = false, anyFailed = false;
        for (StepDef s : steps) {
            if (stopped) { results.add(new StepResult(s.name(), Status.SKIPPED)); continue; }
            Status st;
            try { st = s.body().get(); } catch (RuntimeException e) { st = Status.FAILED; }
            results.add(new StepResult(s.name(), st));
            if (st == Status.FAILED) { anyFailed = true; if (s.onFail() == OnFail.STOP) stopped = true; }
        }
        System.out.println("=== JOB " + jobName + " " + (anyFailed ? "FAILED" : "COMPLETED") + " ===");
        results.forEach(r -> System.out.println("    " + r.name() + " → " + r.status()));
        return results;
    }

    public static void main(String[] args) {
        Supplier<Status> ok = () -> Status.COMPLETED;
        Supplier<Status> boom = () -> { throw new IllegalStateException("집계 오류"); };

        // 시나리오 1: 알림 발송(CONTINUE) 실패는 정산 파일 생성을 막지 않는다
        run("dailyJob-A", List.of(
            new StepDef("aggregateTrades", ok, OnFail.STOP),
            new StepDef("sendNotifications", boom, OnFail.CONTINUE),
            new StepDef("writeSettlementFile", ok, OnFail.STOP)));

        // 시나리오 2: 거래 집계(STOP) 실패는 이후 전부 중단 — 불완전한 데이터로 정산 파일을 만들면 안 된다
        run("dailyJob-B", List.of(
            new StepDef("aggregateTrades", boom, OnFail.STOP),
            new StepDef("sendNotifications", ok, OnFail.CONTINUE),
            new StepDef("writeSettlementFile", ok, OnFail.STOP)));
    }
}
// 출력:
// === JOB dailyJob-A FAILED ===
//     aggregateTrades → COMPLETED
//     sendNotifications → FAILED
//     writeSettlementFile → COMPLETED
// === JOB dailyJob-B FAILED ===
//     aggregateTrades → FAILED
//     sendNotifications → SKIPPED
//     writeSettlementFile → SKIPPED

변형 5: 파일 Reader + 청크 크기별 실측 + 엣지 케이스(빈 입력, 쓰기 예외 시 close)

지금까지 Reader 는 List 였습니다. 실제 배치처럼 10만 행 CSV 를 임시 파일로 만들고 BufferedReader 기반 Reader 로 읽되, Writer 에 커밋 고정 비용(200µs, DB fsync 대체)을 넣어 청크 크기 10/100/1,000/10,000 의 소요 시간을 잽니다.

커밋 비용이 지배적일 때 청크 크기가 100→1,000 에서 효과가 크고 그 이상은 평탄해지는 것이 03 레슨의 곡선 그대로입니다. 마지막으로 "헤더만 있는 빈 파일"과 "6번째 청크에서 Writer 예외"를 넣어 커밋 0 건, 그리고 예외가 나도 close() 가 정확히 한 번 호출되는지 확인합니다.

java
import java.io.*;
import java.nio.charset.StandardCharsets;
import java.nio.file.*;
import java.util.*;

public class ChunkSizeBench {
    interface ItemReader<T> extends AutoCloseable {
        T read() throws Exception;
        @Override default void close() throws Exception {}
    }
    interface ItemWriter<T> { void write(List<? extends T> items) throws Exception; }

    // 파일 Reader: 한 줄씩. close 횟수를 세어 "예외가 나도 항상 닫히는가" 를 검증한다
    static class CsvReader implements ItemReader<String> {
        static int closeCount = 0;
        private final BufferedReader br;
        CsvReader(Path p) throws IOException { br = Files.newBufferedReader(p, StandardCharsets.UTF_8); br.readLine(); }   // 헤더
        public String read() throws IOException { return br.readLine(); }
        @Override public void close() throws IOException { br.close(); closeCount++; }
    }

    // 쓰기 비용 흉내: 커밋마다 고정 비용 200µs (DB fsync 대체) + 건당 파싱
    static long runStep(ItemReader<String> reader, int chunkSize, ItemWriter<String> writer) throws Exception {
        long commits = 0;
        try (reader) {
            List<String> chunk = new ArrayList<>(chunkSize);
            String line;
            while ((line = reader.read()) != null) {
                chunk.add(line.substring(0, line.indexOf(',')));            // "가공"
                if (chunk.size() == chunkSize) { writer.write(chunk); commits++; chunk.clear(); }
            }
            if (!chunk.isEmpty()) { writer.write(chunk); commits++; }
        }
        return commits;
    }

    static void busySleepMicros(long micros) { long end = System.nanoTime() + micros * 1_000; while (System.nanoTime() < end) Thread.onSpinWait(); }

    public static void main(String[] args) throws Exception {
        Path dir = Files.createTempDirectory("batch01");
        Path csv = dir.resolve("orders.csv");
        int rows = 100_000;
        try (BufferedWriter w = Files.newBufferedWriter(csv, StandardCharsets.UTF_8)) {
            w.write("id,customer,amount"); w.newLine();
            for (int i = 1; i <= rows; i++) { w.write(i + ",cust-" + (i % 500) + "," + (i % 90_000)); w.newLine(); }
        }
        ItemWriter<String> writer = chunk -> busySleepMicros(200);          // 커밋 고정 비용

        runStep(new CsvReader(csv), 1000, writer);                          // 워밍업
        System.out.println("chunkSize  commits   elapsed");
        for (int size : new int[]{10, 100, 1_000, 10_000}) {
            long t0 = System.nanoTime();
            long commits = runStep(new CsvReader(csv), size, writer);
            System.out.printf("%,9d  %,7d  %,6d ms%n", size, commits, (System.nanoTime() - t0) / 1_000_000);
        }

        // 엣지 1: 빈 파일(헤더만) → 커밋 0
        Path empty = dir.resolve("empty.csv");
        Files.writeString(empty, "id,customer,amount\n", StandardCharsets.UTF_8);
        System.out.println("빈 입력 commits = " + runStep(new CsvReader(empty), 100, writer));

        // 엣지 2: 특정 청크에서 writer 예외 → 그래도 reader 는 닫힌다
        int before = CsvReader.closeCount;
        try {
            runStep(new CsvReader(csv), 1000, chunk -> { if (chunk.get(0).equals("5001")) throw new IllegalStateException("chunk 6 쓰기 실패"); });
        } catch (IllegalStateException e) {
            System.out.println("예외: " + e.getMessage() + ", reader close 호출 = " + (CsvReader.closeCount - before));
        }
    }
}
// 출력 (시간 값은 환경에 따라 다름):
// chunkSize  commits   elapsed
//        10   10,000   4,306 ms
//       100    1,000     848 ms
//     1,000      100     180 ms
//    10,000       10     169 ms
// 빈 입력 commits = 0
// 예외: chunk 6 쓰기 실패, reader close 호출 = 1
응용 변형 예제
  • 변형 1: 청크 기준을 파라미터로 — 읽기 기준 vs 쓰기 기준 나란히
  • 변형 2: Processor 합성 — 검증 → 변환 → 필터 → 계산 체인
  • 변형 3: 리스너로 관심사 분리 — beforeChunk / afterChunk / onError
  • 변형 4: Step 별 실패 정책 — STOP vs CONTINUE
  • 변형 5: 파일 Reader + 청크 크기별 실측 + 엣지 케이스(빈 입력, 쓰기 예외 시 close)
이전 섹션3 코드 예제4 / 7다음 섹션5 자주 하는 실수 (Tip)