ChunkProcessor 를 이용해 data/sales.csv 에서 수량(qty)이 5인 판매만 골라 data/qty5.csv 에 쓰는 프로그램을 작성하세요. 청크 크기 500, Writer 는 BufferedWriter 에 APPEND 로 이어 씁니다(파일은 시작 시 새로 만듦). 마지막에 Result 를 출력하세요.
Path out = Path.of("data/qty5.csv");
Files.deleteIfExists(out);
Files.writeString(out, "id,productId,qty,price\n", StandardCharsets.UTF_8);
try (StreamingFileReader r = new StreamingFileReader(Path.of("data/sales.csv"), true);
BufferedWriter w = Files.newBufferedWriter(out, StandardCharsets.UTF_8, StandardOpenOption.APPEND)) {
ChunkProcessor.Result res = new ChunkProcessor<String, String>(500)
.run(r::read,
line -> line.split(",")[2].equals("5") ? line : null, // 필터
chunk -> {
for (String line : chunk) { w.write(line); w.newLine(); }
w.flush(); // 청크 = flush 단위
});
System.out.println(res);
}
// 출력: read=100,000 written=20,0xx chunks=200 60ms (1,666,666 items/s)
// 힙: BufferedWriter 8KB + 청크 500줄. 파일 크기와 무관ChunkProcessor.run() 에 경과 시간 기준 진행률 로그를 추가하세요. logEverySeconds(int) 로 설정하며, 마지막 로그 이후 N 초가 지났으면 청크가 끝날 때 로그를 찍습니다. 기존 logEveryChunks 와 함께 쓸 수 있어야 합니다(둘 중 하나라도 조건을 만족하면 출력).
// 필드 추가
private int logEverySeconds = 0;
private long lastLogMillis;
public ChunkProcessor<I, O> logEverySeconds(int s) { this.logEverySeconds = s; return this; }
// logProgress 교체
private void logProgress(long start, long read, long chunks) {
long now = System.currentTimeMillis();
boolean byChunk = logEveryChunks > 0 && chunks % logEveryChunks == 0;
boolean byTime = logEverySeconds > 0 && now - lastLogMillis >= logEverySeconds * 1000L;
if (!byChunk && !byTime) return;
lastLogMillis = now;
long elapsed = now - 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);
}
// run() 시작부에 lastLogMillis = start; 추가
// 사용: new ChunkProcessor<String, Sale>(1000).logEverySeconds(5).expectedTotal(rows)
// 출력 (10시간짜리 배치라면):
// [progress] chunk#3,120 read=3,120,000 104,000 items/s ETA 28800s (10.4%) ← 5초마다runParallel 에서 워커가 예외를 던지면 어느 청크(몇 번째)에서 실패했는지 알 수 없습니다. 청크 번호를 함께 기록하도록 수정하고, 실패 시 ChunkFailedException(chunkNo, cause) 를 던지세요. 그리고 30번째 청크의 첫 항목에서 일부러 예외를 던지는 Processor 로 테스트해 예외 메시지에 30 이 포함되는지 확인하세요.
public static class ChunkFailedException extends Exception {
public final long chunkNo;
public ChunkFailedException(long chunkNo, Throwable cause) {
super("chunk #" + chunkNo + " failed: " + cause.getMessage(), cause);
this.chunkNo = chunkNo;
}
}
// runParallel 내부 수정
AtomicReference<ChunkFailedException> failure = new AtomicReference<>();
...
final long thisChunk = chunks; // 람다 캡처용
pool.submit(() -> {
try {
List<O> outs = processChunk(chunk, processor);
writer.write(outs);
written.addAndGet(outs.size());
} catch (Exception e) {
failure.compareAndSet(null, new ChunkFailedException(thisChunk, e));
} finally {
inFlight.release();
}
});
...
if (failure.get() != null) throw failure.get();
// 테스트
int[] seen = {0};
try (StreamingFileReader r = new StreamingFileReader(Path.of("data/sales.csv"), true)) {
new ChunkProcessor<String, Sale>(1_000).runParallel(r::read,
line -> {
Sale s = Sale.parse(line);
if (s.id() == 29_001) throw new IllegalStateException("bad row id=" + s.id()); // 30번째 청크 첫 항목
return s;
},
chunk -> {}, 4);
} catch (ChunkFailedException e) {
System.out.println(e.getMessage() + " / chunkNo=" + e.chunkNo);
}
// 출력: chunk #30 failed: bad row id=29001 / chunkNo=30
// 주의: 병렬이므로 30번 청크가 실패해도 31, 32번 청크는 이미 제출되어 처리될 수 있다.
// "실패 청크 이후는 처리 안 됨" 을 보장하려면 순차 처리가 필요하다.