List<Result> results = new ArrayList<>();
processor.run(r::read, this::transform, results::addAll); // 결국 전체 로딩
saveAll(results);메모리 사용량이 입력 크기에 비례합니다. 청크 처리의 의미가 없습니다.
✅ Writer 가 청크를 받은 즉시 외부로 내보냅니다. 남기는 것은 입력 크기와 무관한 것(집계 맵, 카운터)만.
processor.run(r::read, this::transform, chunk -> dao.batchInsert(chunk)); // 즉시 저장try (Stream<String> lines = Files.lines(path)) {
Map<String, List<Order>> byCustomer = lines.map(Order::parse)
.collect(Collectors.groupingBy(Order::customer)); // 전체 Order 가 Map 에
}Files.lines 는 지연 스트림이지만 groupingBy 는 모든 요소를 결과 Map 에 담습니다. toList(), sorted(), collect(toMap) 도 마찬가지입니다. 스트림이 스트리밍인지는 종단 연산이 무엇을 보관하는지로 결정됩니다.
✅ 보관하는 것이 입력 크기와 무관하도록 집계합니다. 정렬이 필요하면 DB 나 외부 정렬을 씁니다.
Map<String, Long> countByCustomer = lines.map(Order::parse)
.collect(Collectors.groupingBy(Order::customer, Collectors.counting())); // 고객 수만큼만new ChunkProcessor<>(100_000) // 청크 하나에 100,000 건 × (입력 + 출력)건당 1 KB 면 청크 하나가 200 MB 입니다. 병렬 4스레드 × in-flight 8 이면 1.6 GB. 게다가 실패 시 100,000 건이 롤백되고, DB 트랜잭션이 오래 열려 락을 잡습니다. 그리고 1,000 이상에서는 속도 이득이 거의 없습니다.
✅ 100~1,000 에서 시작하고 측정합니다. 청크 하나의 메모리 = chunkSize × 건당 크기 × 2 를 계산해 둡니다.
Map<Integer, Long> total = new HashMap<>(); // 여러 워커가 동시에 merge
processor.runParallel(r::read, Sale::parse, chunk -> { for (Sale s : chunk) total.merge(...); }, 4);HashMap 은 동시 쓰기에 안전하지 않습니다. 예외 없이 값이 조용히 유실되거나 무한 루프에 빠집니다. 합계가 매번 다르게 나옵니다.
✅ ConcurrentHashMap, AtomicLong, 또는 스레드별 부분 결과를 만들고 마지막에 합칩니다. 파일 쓰기는 synchronized 또는 스레드별 파일.
while ((chunk = readChunk()) != null) pool.submit(() -> process(chunk)); // 제한 없이 제출읽기가 처리보다 빠르면(디스크 순차 읽기는 매우 빠름) 청크가 큐에 무한히 쌓입니다. 100만 건이면 큐에 100만 건 = 전체 로딩 = OOM. 병렬화했더니 오히려 죽는 전형적 사례입니다.
✅ Semaphore 로 in-flight 청크 수를 제한하거나, 큐 크기가 제한된 ThreadPoolExecutor + CallerRunsPolicy 를 씁니다.
Semaphore inFlight = new Semaphore(threads * 2);
inFlight.acquire(); // 한도 초과 시 읽기 스레드가 대기
pool.submit(() -> { try { process(chunk); } finally { inFlight.release(); } });for (...) { process(item); System.out.println("processed " + i); } // 100만 줄 로그콘솔/파일 I/O 가 처리 자체보다 느려집니다. 로그 파일이 수 GB 가 되고, 정작 필요한 오류 로그를 찾을 수 없습니다.
✅ N 청크마다 또는 N 초마다 한 줄. 건수·속도·ETA 를 포함.