java-src/batch/03_chunk/ 에서 javac *.java && java Main 100000. OOM 재현은 java -Xmx32m Main 1000000.
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만 줄이어도 사용량 변화 거의 없음Reader/Processor/Writer 를 함수형 인터페이스로 받아 청크 루프를 돌립니다. 진행률과 ETA 를 함께 찍습니다.
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)Runtime 으로 힙을 재면 차이가 숫자로 보입니다.
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, 스트리밍은 정상세 예제 모두 같은 ChunkProcessor 를 쓰고 Reader/Processor/Writer 만 바뀝니다.
// (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)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 (집계는 순서와 무관하므로 병렬 가능)