3절의 ChunkProcessor 는 "N건 읽고 → 가공 → 쓰고 → 버린다"는 한 가지 루프였습니다. 여기서는 그 루프를 청크 크기 실측, 힙 측정, DB 페이징 시뮬레이션, 백프레셔 유무 비교, 실패·인터럽트 처리라는 다섯 가지 각도에서 다시 씁니다. 모두 외부 의존 없이 실행되며, 시간·메모리 숫자는 환경에 따라 달라지지만 경향은 같습니다.
2.4 절의 "100 부터 시작해서 측정하라"를 직접 해 봅니다. 커밋 한 번에 50µs 가 드는 DB 를 흉내 내어(spin 대기) 청크 크기를 1 → 10,000 으로 바꿔 가며 처리량을 잽니다. LockSupport.parkNanos 대신 spin 을 쓴 이유는 Windows 에서 park 의 최소 단위가 약 1ms 라 50µs 를 정확히 흉내 낼 수 없기 때문입니다.
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) 무릎이 오른쪽으로 조금 이동할 뿐 모양은 같습니다.
예제 3 을 Runtime 으로 다시 재되, 측정 시점마다 System.gc() 를 호출해 살아 있는 객체만 셉니다. GC 없이 재면 Eden 에 쌓인 쓰레기까지 포함되어 스트리밍이 더 커 보이는 착시가 납니다. 파싱한 객체 리스트가 원본 String 리스트의 두 배 이상이라는 2.1 절의 계산도 확인합니다.
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 으로 실행하면 전체 로딩 쪽만 죽는 것을 확인할 수 있습니다.
2.6 절의 두 방식을 TreeMap 으로 흉내 낸 가짜 테이블 위에서 비교합니다. offset 방식은 페이지가 뒤로 갈수록 앞 행을 전부 훑고 버리므로 DB 가 읽은 행 수가 폭발하고, 키셋 방식은 tailMap(lastId, false) 로 인덱스 seek 를 하므로 읽은 행 수 = 반환한 행 수입니다. 그리고 배치 도중 온라인 서비스가 앞쪽 행을 삭제했을 때 offset 방식은 건을 건너뛴다는 것도 재현합니다.
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 ms20,000 건을 1,000 건씩 20 페이지로 읽는데 offset 방식은 230,000 행을 스캔합니다(1+2+…+20 페이지 분량, O(N²/page)). 100만 건이면 5억 행입니다. 삭제 실험에서 offset 방식은 100 건을 조용히 건너뛰었습니다. 앞 행이 사라지면서 뒤 행이 당겨져 다음 offset 이 그만큼 넘어갔기 때문입니다. 키셋 방식은 "마지막 id 다음"이라는 조건이 삭제와 무관하므로 정확히 20,000 건입니다.
예제 5 의 Semaphore 가 정말 필요한지 확인합니다. 읽기는 빠르고(메모리) 쓰기는 느린(1ms sleep) 상황에서 무제한 제출과 Semaphore(threads×2) 제출을 같은 워커 풀로 돌려, 동시에 살아 있는 청크의 최댓값을 잽니다. 힙에 남는 청크 수가 곧 메모리이므로 이 숫자가 OOM 여부를 결정합니다.
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 개면 충분하기 때문입니다.
정상 경로만 있던 청크 루프에 세 가지 비정상 상황을 넣습니다. (1) Reader 가 청크 도중 예외를 던지면 그 청크의 부분 읽기는 버려지고 앞 청크 수만 기록됩니다. (2) 한 청크의 모든 항목이 필터되면 빈 리스트로 write 를 호출하지 않습니다(빈 트랜잭션 커밋 방지).
(3) 병렬 모드에서 읽기 스레드가 인터럽트되면 Semaphore.acquire 가 깨어나고 shutdownNow 로 큐의 작업을 회수합니다. 운영자가 배치를 중단시켰을 때 체크포인트를 남기고 깨끗이 내려가는 경로입니다.
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 로 세마포어를 돌려주므로 풀이 정상 종료됩니다.