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

Chunk 단위 처리와 OOM 방지

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

3. 코드 예제

java-src/batch/03_chunk/ 에서 javac *.java && java Main 100000. OOM 재현은 java -Xmx32m Main 1000000.

예제 1: StreamingFileReader — 한 줄씩, Iterator 지원

java
public class StreamingFileReader implements Iterator<String>, AutoCloseable {
    private final BufferedReader reader;
    private final boolean skipHeader;
    private String nextLine;          // look-ahead 버퍼
    private long lineNo;
    private boolean headerSkipped;

    public StreamingFileReader(Path path, boolean skipHeader) throws IOException {
        this.reader = Files.newBufferedReader(path, StandardCharsets.UTF_8);
        this.skipHeader = skipHeader;
    }

    public String read() throws IOException {                 // ItemReader 스타일
        if (skipHeader && !headerSkipped) { reader.readLine(); headerSkipped = true; }
        String line = reader.readLine();
        if (line != null) lineNo++;
        return line;
    }

    public long lineNo() { return lineNo; }

    public void skip(long n) throws IOException {             // 재시작용
        for (long i = 0; i < n; i++) if (read() == null) break;
    }

    @Override public boolean hasNext() {                       // Iterator 스타일
        if (nextLine != null) return true;
        try { nextLine = read(); } catch (IOException e) { throw new IllegalStateException(e); }
        return nextLine != null;
    }
    @Override public String next() {
        if (!hasNext()) throw new NoSuchElementException();
        String line = nextLine; nextLine = null; return line;
    }
    @Override public void close() throws IOException { reader.close(); }
}
// try (var r = new StreamingFileReader(path, true)) { while (r.hasNext()) count++; }
// 출력: 힙에는 한 줄만 존재. 100만 줄이어도 사용량 변화 거의 없음

예제 2: ChunkProcessor — 청크 루프의 일반화

Reader/Processor/Writer 를 함수형 인터페이스로 받아 청크 루프를 돌립니다. 진행률과 ETA 를 함께 찍습니다.

java
public class ChunkProcessor<I, O> {
    @FunctionalInterface public interface Reader<T> { T read() throws Exception; }
    @FunctionalInterface public interface Processor<I, O> { O process(I item) throws Exception; }
    @FunctionalInterface public interface Writer<T> { void write(List<T> chunk) throws Exception; }

    public record Result(long readCount, long writeCount, long chunkCount, long elapsedMs) {
        public long itemsPerSec() { return elapsedMs == 0 ? readCount : readCount * 1000 / elapsedMs; }
    }

    private final int chunkSize;
    private long expectedTotal = -1;
    private int logEveryChunks = 0;

    public ChunkProcessor(int chunkSize) { this.chunkSize = chunkSize; }
    public ChunkProcessor<I, O> expectedTotal(long t) { expectedTotal = t; return this; }
    public ChunkProcessor<I, O> logEveryChunks(int n) { logEveryChunks = n; return this; }

    private List<I> readChunk(Reader<I> reader) throws Exception {
        List<I> chunk = new ArrayList<>(chunkSize);
        for (int i = 0; i < chunkSize; i++) {
            I item = reader.read();
            if (item == null) break;
            chunk.add(item);
        }
        return chunk;
    }

    public Result run(Reader<I> reader, Processor<I, O> processor, Writer<O> writer) throws Exception {
        long start = System.currentTimeMillis();
        long read = 0, written = 0, chunks = 0;
        while (true) {
            List<I> chunk = readChunk(reader);            // 힙: 청크 하나
            if (chunk.isEmpty()) break;
            read += chunk.size();
            List<O> outs = new ArrayList<>(chunk.size());
            for (I item : chunk) {
                O out = processor.process(item);
                if (out != null) outs.add(out);           // null = 필터
            }
            writer.write(outs);                           // 커밋 포인트
            written += outs.size();
            chunks++;
            logProgress(start, read, chunks);
        }                                                 // chunk, outs 회수 가능
        return new Result(read, written, chunks, System.currentTimeMillis() - start);
    }

    private void logProgress(long start, long read, long chunks) {
        if (logEveryChunks <= 0 || chunks % logEveryChunks != 0) return;
        long elapsed = System.currentTimeMillis() - start;
        long rate = elapsed == 0 ? 0 : read * 1000 / elapsed;
        String eta = "";
        if (expectedTotal > 0 && rate > 0) {
            eta = String.format("  ETA %ds  (%.1f%%)", (expectedTotal - read) / rate, read * 100.0 / expectedTotal);
        }
        System.out.printf("  [progress] chunk#%,d  read=%,d  %,d items/s%s%n", chunks, read, rate, eta);
    }
}
// 출력:
//   [progress] chunk#25  read=25,000  1,250,000 items/s  ETA 0s  (25.0%)
//   [progress] chunk#50  read=50,000  1,190,476 items/s  ETA 0s  (50.0%)
//   결과: read=100,000 written=100,000 chunks=100  85ms  (1,176,470 items/s)

예제 3: 전체 로딩 vs 스트리밍 — 힙 사용량 측정

Runtime 으로 힙을 재면 차이가 숫자로 보입니다.

java
public final class MemoryMonitor {
    private static final long MB = 1024 * 1024;
    public static long usedMb() {
        Runtime rt = Runtime.getRuntime();
        return (rt.totalMemory() - rt.freeMemory()) / MB;
    }
    public static long usedMbAfterGc() {                 // GC 유도 후 측정 = 살아있는 객체 크기
        System.gc();
        try { Thread.sleep(50); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
        return usedMb();
    }
    public static void log(String label) {
        System.out.printf("  [MEM] %-28s used=%4d MB  total=%4d MB  max=%4d MB%n",
                label, usedMb(), Runtime.getRuntime().totalMemory() / MB, Runtime.getRuntime().maxMemory() / MB);
    }
}

// 비교
long before = MemoryMonitor.usedMbAfterGc();
List<String> all = Files.readAllLines(sales, StandardCharsets.UTF_8);
System.out.printf("readAllLines: %,d줄 -> 힙 +%d MB%n", all.size(), MemoryMonitor.usedMb() - before);
all = null;                                              // 참조 해제
System.out.printf("참조 해제 + GC 후: %d MB%n", MemoryMonitor.usedMbAfterGc());

long before2 = MemoryMonitor.usedMbAfterGc(), maxUsed = 0, count = 0;
try (StreamingFileReader r = new StreamingFileReader(sales, true)) {
    while (r.read() != null) {
        if (++count % 20_000 == 0) maxUsed = Math.max(maxUsed, MemoryMonitor.usedMb());
    }
}
System.out.printf("스트리밍: %,d줄 -> 힙 최대 +%d MB%n", count, Math.max(0, maxUsed - before2));
// 출력 (100,000행):
//   readAllLines: 100,000줄 로딩 -> 힙 +9 MB (List<String> 전체가 살아있음)
//   참조 해제 + GC 후: 2 MB
//   스트리밍: 100,000줄 순회 -> 힙 최대 +1 MB (한 줄만 살아있음)
// 출력 (1,000,000행): readAllLines +9x MB, 스트리밍 +1~2 MB
// java -Xmx32m Main 1000000: readAllLines 에서 OutOfMemoryError, 스트리밍은 정상

예제 4: 실무 예제 — 판매 집계, 로그 에러 통계, 청크별 파일 분할

세 예제 모두 같은 ChunkProcessor 를 쓰고 Reader/Processor/Writer 만 바뀝니다.

java
// (a) 판매 파일 상품별 매출 집계 — Writer 가 남기는 것은 상품 수만큼의 Map 뿐
record Sale(long id, int productId, int qty, long price) {
    static Sale parse(String line) {
        String[] f = line.split(",");
        return new Sale(Long.parseLong(f[0]), Integer.parseInt(f[1]), Integer.parseInt(f[2]), Long.parseLong(f[3]));
    }
    long total() { return qty * price; }
}
Map<Integer, Long> totalByProduct = new TreeMap<>();
try (StreamingFileReader r = new StreamingFileReader(sales, true)) {
    new ChunkProcessor<String, Sale>(1_000)
            .expectedTotal(rows).logEveryChunks(25)
            .run(r::read, Sale::parse,
                 chunk -> { for (Sale s : chunk) totalByProduct.merge(s.productId(), s.total(), Long::sum); });
}
// 출력: 결과: read=100,000 written=100,000 chunks=100  85ms
//        상품 1 매출 12,345,000원 ...

// (b) 로그 파일 에러 통계 — Processor 가 ERROR 가 아닌 줄을 null 로 필터
Map<String, Long> errorStats = new HashMap<>();
try (StreamingFileReader r = new StreamingFileReader(log, false)) {
    new ChunkProcessor<String, String>(5_000)
            .run(r::read,
                 line -> line.contains(" ERROR ") ? line.substring(line.indexOf(" ERROR ") + 7) : null,
                 chunk -> chunk.forEach(msg -> errorStats.merge(msg, 1L, Long::sum)));
}
// 출력: 결과: read=100,000 written=2,9xx chunks=20  40ms
//        SQLTimeoutException          7xx건
//        NullPointerException         7xx건 ...

// (c) 청크별 출력 파일 분할 — 청크 하나 = 파일 하나
int[] fileNo = {0};
try (StreamingFileReader r = new StreamingFileReader(sales, true)) {
    new ChunkProcessor<String, String>(20_000)
            .run(r::read, line -> line,
                 chunk -> Files.write(outDir.resolve(String.format("sales-%03d.csv", fileNo[0]++)), chunk, UTF_8));
}
// 출력: 결과: read=100,000 written=100,000 chunks=5 -> 5개 파일 생성 (data\chunks)

예제 5: 병렬 청크 처리 — Semaphore 백프레셔

java
public Result runParallel(Reader<I> reader, Processor<I, O> processor, Writer<O> writer, int threads)
        throws Exception {
    long start = System.currentTimeMillis();
    ExecutorService pool = Executors.newFixedThreadPool(threads);
    Semaphore inFlight = new Semaphore(threads * 2);          // 동시에 살아있는 청크 상한
    AtomicLong written = new AtomicLong();
    AtomicReference<Exception> failure = new AtomicReference<>();
    long read = 0, chunks = 0;
    try {
        while (failure.get() == null) {
            List<I> chunk = readChunk(reader);                // 읽기는 이 스레드가 순차로
            if (chunk.isEmpty()) break;
            read += chunk.size(); chunks++;
            inFlight.acquire();                               // 한도 차면 여기서 대기 = 백프레셔
            pool.submit(() -> {
                try {
                    List<O> outs = processChunk(chunk, processor);
                    writer.write(outs);                       // writer 는 스레드 안전해야 함
                    written.addAndGet(outs.size());
                } catch (Exception e) {
                    failure.compareAndSet(null, e);
                } finally {
                    inFlight.release();
                }
            });
        }
    } finally {
        pool.shutdown();
        pool.awaitTermination(1, TimeUnit.HOURS);
    }
    if (failure.get() != null) throw failure.get();
    return new Result(read, written.get(), chunks, System.currentTimeMillis() - start);
}

// 사용: 가공에 CPU 부하가 있을 때 (parseWithCpuWork = 파싱 + 해시 반복 2000회)
Map<Integer, Long> seq = new HashMap<>();
new ChunkProcessor<String, Sale>(2_000).run(r1::read, Main::parseWithCpuWork,
        chunk -> { for (Sale s : chunk) seq.merge(s.productId(), s.total(), Long::sum); });
ConcurrentHashMap<Integer, Long> par = new ConcurrentHashMap<>();       // 스레드 안전 Writer
new ChunkProcessor<String, Sale>(2_000).runParallel(r2::read, Main::parseWithCpuWork,
        chunk -> { for (Sale s : chunk) par.merge(s.productId(), s.total(), Long::sum); }, 4);
System.out.println("결과 동일? " + seq.equals(par));
// 출력:
//   순차   : read=100,000 written=100,000 chunks=50  312ms  (320,512 items/s)
//   4스레드: read=100,000 written=100,000 chunks=50   98ms  (1,020,408 items/s)
//   결과 동일? true (집계는 순서와 무관하므로 병렬 가능)

예제 직접 실행

아래 폴더를 JDK 21 로 컴파일하고 실행합니다.

cd java-src\batch\03_chunk
javac -encoding UTF-8 *.java && java Main
코드 예제
  • 예제 1: StreamingFileReader — 한 줄씩, Iterator 지원
  • 예제 2: ChunkProcessor — 청크 루프의 일반화
  • 예제 3: 전체 로딩 vs 스트리밍 — 힙 사용량 측정
  • 예제 4: 실무 예제 — 판매 집계, 로그 에러 통계, 청크별 파일 분할
  • 예제 5: 병렬 청크 처리 — Semaphore 백프레셔
이전 섹션2 핵심 원리3 / 7다음 섹션4 응용 변형 예제