공공부하자개발 · 영어 학습 노트
자바
실무 배치대용량 데이터 처리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정리‹ 이전다음 ›

4. 응용 변형 예제

3절의 ChunkProcessor 는 "N건 읽고 → 가공 → 쓰고 → 버린다"는 한 가지 루프였습니다. 여기서는 그 루프를 청크 크기 실측, 힙 측정, DB 페이징 시뮬레이션, 백프레셔 유무 비교, 실패·인터럽트 처리라는 다섯 가지 각도에서 다시 씁니다. 모두 외부 의존 없이 실행되며, 시간·메모리 숫자는 환경에 따라 달라지지만 경향은 같습니다.

변형 1: 청크 크기별 처리량 곡선 — 커밋 비용을 넣고 실측

2.4 절의 "100 부터 시작해서 측정하라"를 직접 해 봅니다. 커밋 한 번에 50µs 가 드는 DB 를 흉내 내어(spin 대기) 청크 크기를 1 → 10,000 으로 바꿔 가며 처리량을 잽니다. LockSupport.parkNanos 대신 spin 을 쓴 이유는 Windows 에서 park 의 최소 단위가 약 1ms 라 50µs 를 정확히 흉내 낼 수 없기 때문입니다.

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

public class ChunkSizeBench {
    static final int ROWS = 100_000;
    static final long COMMIT_COST_NANOS = 50_000;               // 커밋 1회 = 50µs (fsync 흉내)

    record Sale(int productId, long total) {}

    /** 최소한의 청크 루프: 읽기 N건 → 가공 → 쓰기(커밋) → 반복 */
    static <I, O> long[] run(Iterator<I> reader, Function<I, O> proc, Consumer<List<O>> writer, int chunkSize) {
        long read = 0, commits = 0;
        while (reader.hasNext()) {
            List<O> outs = new ArrayList<>(chunkSize);
            for (int i = 0; i < chunkSize && reader.hasNext(); i++) { outs.add(proc.apply(reader.next())); read++; }
            writer.accept(outs); commits++;
        }
        return new long[]{read, commits};
    }

    static void simulateCommit() {                              // parkNanos 는 Windows 에서 1ms 단위라 spin 으로 정확히 대기
        long end = System.nanoTime() + COMMIT_COST_NANOS;
        while (System.nanoTime() < end) Thread.onSpinWait();
    }

    public static void main(String[] args) {
        List<String> lines = new ArrayList<>(ROWS);
        for (int i = 0; i < ROWS; i++) lines.add(i + "," + (i % 50) + "," + (i % 7 + 1) + "," + (1000 + i % 900));

        System.out.printf("%-8s %8s %10s %14s%n", "chunk", "commits", "ms", "items/s");
        for (int size : new int[]{1, 10, 100, 1_000, 10_000}) {
            Map<Integer, Long> total = new HashMap<>();
            long t0 = System.nanoTime();
            long[] r = run(lines.iterator(),
                    l -> { String[] f = l.split(","); return new Sale(Integer.parseInt(f[1]), (long) Integer.parseInt(f[2]) * Integer.parseInt(f[3])); },
                    chunk -> { for (Sale s : chunk) total.merge(s.productId(), s.total(), Long::sum); simulateCommit(); },
                    size);
            long ms = (System.nanoTime() - t0) / 1_000_000;
            System.out.printf("%-8d %8d %10d %14s%n", size, r[1], ms, String.format("%,d", ms == 0 ? 0 : r[0] * 1000 / ms));
        }
    }
}
// 출력 (ms·items/s 는 환경에 따라 다름. 100 이후 평탄해지는 모양이 핵심):
// chunk     commits         ms        items/s
// 1          100000       9194         10,876
// 10          10000       1092         91,575
// 100          1000        239        418,410
// 1000          100        245        408,163
// 10000          10        263        380,228

청크 1 은 커밋 10만 번(5초)이 전부이고, 100 부터는 커밋 비용이 사라져 파싱 자체가 병목이 됩니다. 1,000 과 10,000 이 100 보다 빠르지 않다는 점이 2.4 절 곡선의 근거입니다. 커밋 비용이 더 크다면(원격 DB, 1ms) 무릎이 오른쪽으로 조금 이동할 뿐 모양은 같습니다.

변형 2: 힙 실측 — 전체 로딩 vs 스트리밍, GC 후 살아 있는 객체만 재기

예제 3 을 Runtime 으로 다시 재되, 측정 시점마다 System.gc() 를 호출해 살아 있는 객체만 셉니다. GC 없이 재면 Eden 에 쌓인 쓰레기까지 포함되어 스트리밍이 더 커 보이는 착시가 납니다. 파싱한 객체 리스트가 원본 String 리스트의 두 배 이상이라는 2.1 절의 계산도 확인합니다.

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

public class HeapCompare {
    static final int ROWS = 300_000;

    static long usedMbAfterGc() {
        System.gc();
        try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
        Runtime rt = Runtime.getRuntime();
        return (rt.totalMemory() - rt.freeMemory()) / (1024 * 1024);
    }

    record Order(long id, String customer, long amount) {
        static Order parse(String l) { String[] f = l.split(","); return new Order(Long.parseLong(f[0]), f[1], Long.parseLong(f[2])); }
    }

    public static void main(String[] args) throws IOException {
        Path dir = Files.createTempDirectory("heap");
        Path file = dir.resolve("orders.csv");
        try (BufferedWriter w = Files.newBufferedWriter(file, StandardCharsets.UTF_8)) {
            for (int i = 0; i < ROWS; i++) { w.write(i + ",customer-" + (i % 1000) + "," + (i * 37 % 100_000)); w.newLine(); }
        }
        System.out.printf("파일 %,d 줄, %,d KB, 최대 힙 %d MB%n", ROWS, Files.size(file) / 1024, Runtime.getRuntime().maxMemory() / (1024 * 1024));

        // 1) 전체 로딩: String 리스트 → 파싱한 객체 리스트
        long base = usedMbAfterGc();
        List<String> all = Files.readAllLines(file, StandardCharsets.UTF_8);
        long afterLines = usedMbAfterGc();
        List<Order> parsed = new ArrayList<>(all.size());
        for (String l : all) parsed.add(Order.parse(l));
        long afterParse = usedMbAfterGc();
        System.out.printf("readAllLines      : +%d MB (String %,d개)%n", afterLines - base, all.size());
        System.out.printf("+ 전부 파싱        : +%d MB (Order %,d개, 필드 String 포함)%n", afterParse - base, parsed.size());
        all = null; parsed = null;
        System.out.printf("참조 해제 + GC     : +%d MB%n", Math.max(0, usedMbAfterGc() - base));

        // 2) 스트리밍: 같은 파싱을 하되 집계 맵만 남긴다
        base = usedMbAfterGc();
        long peak = 0; int n = 0;
        Map<String, Long> byCustomer = new HashMap<>();
        try (BufferedReader r = Files.newBufferedReader(file, StandardCharsets.UTF_8)) {
            String line;
            while ((line = r.readLine()) != null) {
                Order o = Order.parse(line);
                byCustomer.merge(o.customer(), o.amount(), Long::sum);
                if (++n % 50_000 == 0) peak = Math.max(peak, usedMbAfterGc() - base);   // GC 후 = 살아있는 객체만
            }
        }
        System.out.printf("스트리밍 + 집계     : 최대 +%d MB (고객 %,d명 맵만 유지)%n", Math.max(0, peak), byCustomer.size());
        Files.delete(file); Files.delete(dir);
    }
}
// 출력 (MB 값은 JVM·힙 설정에 따라 다름):
// 파일 300,000 줄, 7,736 KB, 최대 힙 4042 MB
// readAllLines      : +21 MB (String 300,000개)
// + 전부 파싱        : +48 MB (Order 300,000개, 필드 String 포함)
// 참조 해제 + GC     : +0 MB
// 스트리밍 + 집계     : 최대 +0 MB (고객 1,000명 맵만 유지)

7.7 MB 파일이 String 으로 21 MB, 파싱하면 48 MB 가 됩니다(약 6배). 스트리밍은 300,000 건을 똑같이 파싱했지만 GC 후 남는 것은 고객 1,000 명 맵뿐이라 1 MB 미만입니다. -Xmx40m 으로 실행하면 전체 로딩 쪽만 죽는 것을 확인할 수 있습니다.

변형 3: DB 페이징 시뮬레이션 — offset vs 키셋, 스캔 행 수와 삭제 안전성

2.6 절의 두 방식을 TreeMap 으로 흉내 낸 가짜 테이블 위에서 비교합니다. offset 방식은 페이지가 뒤로 갈수록 앞 행을 전부 훑고 버리므로 DB 가 읽은 행 수가 폭발하고, 키셋 방식은 tailMap(lastId, false) 로 인덱스 seek 를 하므로 읽은 행 수 = 반환한 행 수입니다. 그리고 배치 도중 온라인 서비스가 앞쪽 행을 삭제했을 때 offset 방식은 건을 건너뛴다는 것도 재현합니다.

java
import java.util.*;

public class PagingReaders {
    record Order(long id, long amount) {}

    /** DB 테이블 흉내: PK 순 정렬. offset 조회는 앞 행을 전부 훑고 버려야 한다 */
    static class FakeTable {
        final TreeMap<Long, Order> rows = new TreeMap<>();
        long scanned;                                            // DB 가 실제로 읽은 행 수

        List<Order> findByOffset(int offset, int limit) {       // ORDER BY id LIMIT ? OFFSET ?
            List<Order> page = new ArrayList<>(limit);
            int i = 0;
            for (Order o : rows.values()) {
                scanned++;
                if (i++ < offset) continue;
                page.add(o);
                if (page.size() == limit) break;
            }
            return page;
        }
        List<Order> findAfter(long lastId, int limit) {         // WHERE id > ? ORDER BY id LIMIT ?
            List<Order> page = new ArrayList<>(limit);
            for (Order o : rows.tailMap(lastId, false).values()) {   // 인덱스 seek: lastId 다음부터
                scanned++;
                page.add(o);
                if (page.size() == limit) break;
            }
            return page;
        }
    }

    interface Reader { Order read(); }

    static class OffsetReader implements Reader {
        final FakeTable t; final int pageSize; int offset; Deque<Order> page = new ArrayDeque<>();
        OffsetReader(FakeTable t, int pageSize) { this.t = t; this.pageSize = pageSize; }
        public Order read() {
            if (page.isEmpty()) { page.addAll(t.findByOffset(offset, pageSize)); offset += pageSize; }
            return page.poll();
        }
    }
    static class KeysetReader implements Reader {
        final FakeTable t; final int pageSize; long lastId; Deque<Order> page = new ArrayDeque<>();
        KeysetReader(FakeTable t, int pageSize) { this.t = t; this.pageSize = pageSize; }
        public Order read() {
            if (page.isEmpty()) t.findAfter(lastId, pageSize).forEach(page::add);
            Order o = page.poll();
            if (o != null) lastId = o.id();                      // 체크포인트 = lastId
            return o;
        }
    }

    static FakeTable newTable(int n) {
        FakeTable t = new FakeTable();
        for (long i = 1; i <= n; i++) t.rows.put(i * 3, new Order(i * 3, i));   // id 에 빈틈이 있어도 무관
        return t;
    }

    static void drain(String label, FakeTable t, Reader r, Runnable midway) {
        long t0 = System.nanoTime(); long count = 0;
        Order o;
        while ((o = r.read()) != null) { if (++count == 5_000) midway.run(); }
        System.out.printf("%-8s 읽음 %,d건, DB 스캔 %,d행, %d ms%n", label, count, t.scanned, (System.nanoTime() - t0) / 1_000_000);
    }

    public static void main(String[] args) {
        int n = 20_000, page = 1_000;
        FakeTable a = newTable(n), b = newTable(n);
        drain("offset", a, new OffsetReader(a, page), () -> {});
        drain("keyset", b, new KeysetReader(b, page), () -> {});

        // 배치 도중 앞쪽 행이 삭제되면? (온라인 서비스가 주문을 취소)
        FakeTable c = newTable(n), d = newTable(n);
        System.out.println("-- 5,000건 읽은 시점에 앞쪽 100행 삭제 --");
        drain("offset", c, new OffsetReader(c, page), () -> { for (long i = 1; i <= 100; i++) c.rows.remove(i * 3); });
        drain("keyset", d, new KeysetReader(d, page), () -> { for (long i = 1; i <= 100; i++) d.rows.remove(i * 3); });
    }
}
// 출력 (ms 는 환경에 따라 다름. 스캔 행 수와 읽은 건수는 결정적):
// offset   읽음 20,000건, DB 스캔 230,000행, 313 ms
// keyset   읽음 20,000건, DB 스캔 20,000행, 90 ms
// -- 5,000건 읽은 시점에 앞쪽 100행 삭제 --
// offset   읽음 19,900건, DB 스캔 229,800행, 362 ms
// keyset   읽음 20,000건, DB 스캔 20,000행, 17 ms

20,000 건을 1,000 건씩 20 페이지로 읽는데 offset 방식은 230,000 행을 스캔합니다(1+2+…+20 페이지 분량, O(N²/page)). 100만 건이면 5억 행입니다. 삭제 실험에서 offset 방식은 100 건을 조용히 건너뛰었습니다. 앞 행이 사라지면서 뒤 행이 당겨져 다음 offset 이 그만큼 넘어갔기 때문입니다. 키셋 방식은 "마지막 id 다음"이라는 조건이 삭제와 무관하므로 정확히 20,000 건입니다.

변형 4: 백프레셔 유무 비교 — 최대 in-flight 청크 수 측정

예제 5 의 Semaphore 가 정말 필요한지 확인합니다. 읽기는 빠르고(메모리) 쓰기는 느린(1ms sleep) 상황에서 무제한 제출과 Semaphore(threads×2) 제출을 같은 워커 풀로 돌려, 동시에 살아 있는 청크의 최댓값을 잽니다. 힙에 남는 청크 수가 곧 메모리이므로 이 숫자가 OOM 여부를 결정합니다.

java
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;

public class BackpressureCompare {
    static final int CHUNKS = 2_000, CHUNK_SIZE = 500, THREADS = 4;

    /** 청크를 만들어 워커에 제출. limit 이 null 이면 무제한 제출 */
    static String run(Semaphore limit) throws Exception {
        ExecutorService pool = Executors.newFixedThreadPool(THREADS);
        AtomicInteger inFlight = new AtomicInteger(), maxInFlight = new AtomicInteger();
        AtomicLong written = new AtomicLong();
        long t0 = System.nanoTime();
        for (int c = 0; c < CHUNKS; c++) {
            List<Integer> chunk = new ArrayList<>(CHUNK_SIZE);           // 읽기는 빠르다 (메모리)
            for (int i = 0; i < CHUNK_SIZE; i++) chunk.add(c * CHUNK_SIZE + i);
            if (limit != null) limit.acquire();                         // 백프레셔: 한도 차면 읽기 스레드가 대기
            int now = inFlight.incrementAndGet();
            maxInFlight.accumulateAndGet(now, Math::max);
            pool.submit(() -> {
                try {
                    long sum = 0; for (int v : chunk) sum += v;         // 가공
                    Thread.sleep(1);                                    // 쓰기 = 느린 I/O
                    written.addAndGet(chunk.size());
                } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
                finally { inFlight.decrementAndGet(); if (limit != null) limit.release(); }
            });
        }
        pool.shutdown(); pool.awaitTermination(1, TimeUnit.MINUTES);
        return String.format("written=%,d  최대 in-flight 청크 %,d개 (≈ %,d KB 힙)  %d ms",
                written.get(), maxInFlight.get(), (long) maxInFlight.get() * CHUNK_SIZE * 16 / 1024, (System.nanoTime() - t0) / 1_000_000);
    }

    public static void main(String[] args) throws Exception {
        System.out.println("무제한 제출     : " + run(null));
        System.out.println("Semaphore(" + THREADS * 2 + ") : " + run(new Semaphore(THREADS * 2)));
    }
}
// 출력 (무제한 쪽의 in-flight 수와 ms 는 실행마다 다름. Semaphore 쪽은 항상 ≤ 8):
// 무제한 제출     : written=1,000,000  최대 in-flight 청크 1,923개 (≈ 15,023 KB 힙)  2180 ms
// Semaphore(8) : written=1,000,000  최대 in-flight 청크 8개 (≈ 62 KB 힙)  1714 ms

무제한 제출은 2,000 청크 중 1,923 개가 큐에 쌓였습니다. 읽기가 쓰기보다 훨씬 빠르니 사실상 전체 로딩입니다. 청크가 500 건이 아니라 500 KB 짜리 JSON 이었다면 1 GB 가 큐에 앉아 있는 셈입니다. Semaphore 를 걸면 최대 8 개로 고정되고, 처리량은 오히려 손해가 없습니다. 워커 4 개가 쉬지 않고 도는 데는 청크 8 개면 충분하기 때문입니다.

변형 5: 실패와 인터럽트 — Reader 예외, 빈 청크, 병렬 중 취소

정상 경로만 있던 청크 루프에 세 가지 비정상 상황을 넣습니다. (1) Reader 가 청크 도중 예외를 던지면 그 청크의 부분 읽기는 버려지고 앞 청크 수만 기록됩니다. (2) 한 청크의 모든 항목이 필터되면 빈 리스트로 write 를 호출하지 않습니다(빈 트랜잭션 커밋 방지).

(3) 병렬 모드에서 읽기 스레드가 인터럽트되면 Semaphore.acquire 가 깨어나고 shutdownNow 로 큐의 작업을 회수합니다. 운영자가 배치를 중단시켰을 때 체크포인트를 남기고 깨끗이 내려가는 경로입니다.

java
import java.util.*;
import java.util.concurrent.*;
import java.util.function.*;

public class ChunkEdgeCases {
    interface Reader<T> { T read() throws Exception; }
    record Result(long read, long written, long chunks, long writes) {}

    static <I, O> Result run(Reader<I> reader, Function<I, O> proc, Consumer<List<O>> writer, int chunkSize) throws Exception {
        long read = 0, written = 0, chunks = 0, writes = 0;
        while (true) {
            List<I> chunk = new ArrayList<>(chunkSize);
            try {
                for (int i = 0; i < chunkSize; i++) { I item = reader.read(); if (item == null) break; chunk.add(item); }
            } catch (Exception e) {
                throw new IllegalStateException("읽기 실패: 청크 " + (chunks + 1) + " 의 " + chunk.size() + "건은 쓰이지 않음 (커밋된 청크 " + chunks + ")", e);
            }
            if (chunk.isEmpty()) break;
            read += chunk.size(); chunks++;
            List<O> outs = new ArrayList<>();
            for (I item : chunk) { O o = proc.apply(item); if (o != null) outs.add(o); }
            if (!outs.isEmpty()) { writer.accept(outs); written += outs.size(); writes++; }   // 빈 청크는 커밋하지 않음
        }
        return new Result(read, written, chunks, writes);
    }

    static Reader<Integer> counting(int n, int failAt) {
        int[] i = {0};
        return () -> { if (++i[0] == failAt) throw new java.io.IOException("디스크 읽기 오류 at " + failAt); return i[0] <= n ? i[0] : null; };
    }

    public static void main(String[] args) throws Exception {
        // 1) 청크 도중 Reader 예외: 앞 청크는 커밋됨, 실패 청크의 부분 읽기는 버려짐
        try { run(counting(1_000, 250), x -> x, c -> {}, 100); }
        catch (IllegalStateException e) { System.out.println("(1) " + e.getMessage() + " / cause=" + e.getCause().getMessage()); }

        // 2) 어떤 청크는 전부 필터됨 → 빈 write 는 건너뛴다
        Result r = run(counting(1_000, -1), x -> x > 300 && x <= 600 ? null : x, c -> {}, 100);
        System.out.println("(2) " + r + "  ← 청크 10개 중 3개는 빈 청크라 write 7회");

        // 3) 병렬 모드에서 인터럽트: 대기 중인 acquire 가 깨어나고 shutdownNow 로 남은 작업을 취소
        ExecutorService pool = Executors.newFixedThreadPool(2);
        Semaphore inFlight = new Semaphore(4);
        List<Future<?>> submitted = new ArrayList<>();
        Thread runner = new Thread(() -> {
            try {
                for (int c = 1; ; c++) {
                    inFlight.acquire();                                      // 워커가 느리면 여기서 대기
                    submitted.add(pool.submit(() -> { try { Thread.sleep(200); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } finally { inFlight.release(); } }));
                }
            } catch (InterruptedException e) {
                List<Runnable> dropped = pool.shutdownNow();                 // 큐에서 대기 중이던 작업 반환 + 실행 중 작업 인터럽트
                System.out.println("(3) 인터럽트 수신: 제출 " + submitted.size() + "개, 큐에서 취소 " + dropped.size() + "개, 진행 중 " + (submitted.size() - dropped.size()) + "개 → 체크포인트 저장 후 종료");
            }
        }, "reader");
        runner.start();
        Thread.sleep(300);
        runner.interrupt();                                                  // 운영자의 Ctrl+C / kill 시뮬레이션
        runner.join();
        System.out.println("(3) 풀 종료 완료: " + pool.awaitTermination(5, TimeUnit.SECONDS));
    }
}
// 출력 ((3) 의 개수는 타이밍에 따라 달라질 수 있음):
// (1) 읽기 실패: 청크 3 의 49건은 쓰이지 않음 (커밋된 청크 2) / cause=디스크 읽기 오류 at 250
// (2) Result[read=1000, written=700, chunks=10, writes=7]  ← 청크 10개 중 3개는 빈 청크라 write 7회
// (3) 인터럽트 수신: 제출 4개, 큐에서 취소 2개, 진행 중 2개 → 체크포인트 저장 후 종료
// (3) 풀 종료 완료: true

(1) 에서 250 번째 읽기가 실패했을 때 청크 3 의 49 건은 어디에도 쓰이지 않았습니다. 04 레슨의 체크포인트는 "커밋된 청크 2" 를 기록하고, 재시작은 201 번째부터입니다. (3) 에서 shutdownNow 가 돌려준 List<Runnable> 은 아직 시작도 못 한 청크들입니다.

이 목록의 크기를 로그에 남기면 "몇 청크가 미처리로 남았는지"를 재시작 전에 알 수 있습니다. 실행 중이던 워커는 sleep 에서 인터럽트로 깨어나 finally 로 세마포어를 돌려주므로 풀이 정상 종료됩니다.

응용 변형 예제
  • 변형 1: 청크 크기별 처리량 곡선 — 커밋 비용을 넣고 실측
  • 변형 2: 힙 실측 — 전체 로딩 vs 스트리밍, GC 후 살아 있는 객체만 재기
  • 변형 3: DB 페이징 시뮬레이션 — offset vs 키셋, 스캔 행 수와 삭제 안전성
  • 변형 4: 백프레셔 유무 비교 — 최대 in-flight 청크 수 측정
  • 변형 5: 실패와 인터럽트 — Reader 예외, 빈 청크, 병렬 중 취소
이전 섹션3 코드 예제4 / 7다음 섹션5 자주 하는 실수 (Tip)