3절의 Step 은 "읽기 기준 청크, 실패하면 Job 중단, 로그는 summary 한 줄" 이라는 하나의 정책을 고정해 두었습니다. 여기서는 같은 Reader-Processor-Writer 골격을 유지한 채 청크 기준·Processor 합성·리스너·실패 정책·청크 크기를 바꿔 보며, 이 구조가 왜 유연한지 확인합니다. 각 프로그램은 필요한 인터페이스를 최소 형태로 다시 선언해 단독 실행됩니다.
연습 문제 3 에서 WriteSizedStep 을 따로 만들었지만, 청크 경계를 "읽은 건수로 셀지, 쓸 건수로 셀지"는 enum 하나로 한 Step 안에서 고를 수 있습니다. 필터 비율이 높은 Step(미납 추출처럼 75% 가 걸러지는 경우)에서 두 정책의 커밋 횟수와 청크 크기가 어떻게 달라지는지 같은 입력으로 비교합니다. 쓰기 효율이 중요하면 WRITE_SIZED, 실패 시 재처리 범위를 좁히려면 READ_SIZED 입니다.
import java.util.*;
import java.util.function.*;
import java.util.stream.*;
public class ReadSizedVsWriteSized {
interface ItemReader<T> { T read() throws Exception; static <T> ItemReader<T> of(List<T> l) { var it = l.iterator(); return () -> it.hasNext() ? it.next() : null; } }
interface ItemProcessor<I, O> { O process(I item) throws Exception; }
interface ItemWriter<T> { void write(List<? extends T> items) throws Exception; }
record StepExecution(String stepName, long readCount, long filterCount, long writeCount, long commitCount) {}
enum ChunkMode { READ_SIZED, WRITE_SIZED }
// 청크 경계를 "읽은 건수" 로 셀지 "쓸 건수" 로 셀지를 파라미터로 받는 Step
static <I, O> StepExecution execute(String name, ItemReader<I> reader, ItemProcessor<I, O> processor,
ItemWriter<O> writer, int chunkSize, ChunkMode mode) throws Exception {
long read = 0, filtered = 0, written = 0, commits = 0;
List<O> chunk = new ArrayList<>(chunkSize);
int readInChunk = 0;
I item;
while ((item = reader.read()) != null) {
read++; readInChunk++;
O out = processor.process(item);
if (out == null) filtered++; else chunk.add(out);
boolean full = mode == ChunkMode.READ_SIZED ? readInChunk == chunkSize : chunk.size() == chunkSize;
if (full) {
if (!chunk.isEmpty()) { writer.write(chunk); written += chunk.size(); commits++; }
chunk.clear(); readInChunk = 0;
}
}
if (!chunk.isEmpty()) { writer.write(chunk); written += chunk.size(); commits++; }
return new StepExecution(name, read, filtered, written, commits);
}
record Invoice(long id, boolean paid) {}
public static void main(String[] args) throws Exception {
List<Invoice> invoices = IntStream.rangeClosed(1, 1000).mapToObj(i -> new Invoice(i, i % 4 != 0)).toList(); // 25% 미납
ItemProcessor<Invoice, Long> unpaidOnly = inv -> inv.paid() ? null : inv.id();
for (ChunkMode mode : ChunkMode.values()) {
List<Integer> chunkSizes = new ArrayList<>();
StepExecution exec = execute("extractUnpaid", ItemReader.of(invoices), unpaidOnly,
chunk -> chunkSizes.add(chunk.size()), 100, mode);
System.out.println(mode + ": " + exec);
System.out.println(" 청크별 write 건수 = " + chunkSizes);
}
}
}
// 출력:
// READ_SIZED: StepExecution[stepName=extractUnpaid, readCount=1000, filterCount=750, writeCount=250, commitCount=10]
// 청크별 write 건수 = [25, 25, 25, 25, 25, 25, 25, 25, 25, 25]
// WRITE_SIZED: StepExecution[stepName=extractUnpaid, readCount=1000, filterCount=750, writeCount=250, commitCount=3]
// 청크별 write 건수 = [100, 100, 50]예제 5 의 Processor 는 람다 하나에 검증·변환·필터가 뒤섞여 있었습니다.
ItemProcessor 에 andThen default 메서드를 추가하면 작은 Processor 를 이어 붙여 파이프라인을 만들 수 있고, 앞 단계가 null(필터)을 돌려주면 뒤 단계는 실행되지 않으므로 "필터 = null" 규약이 체인 전체에 자연스럽게 전파됩니다.
단계마다 타입이 바뀌어도(RawOrder → Order → SettlementLine) 제네릭이 따라갑니다. Spring Batch 의 CompositeItemProcessor 가 바로 이것입니다.
import java.util.*;
public class ProcessorChain {
@FunctionalInterface
interface ItemProcessor<I, O> {
O process(I item) throws Exception;
// 합성: 앞 단계가 null(필터) 이면 뒤 단계는 실행하지 않는다
default <R> ItemProcessor<I, R> andThen(ItemProcessor<? super O, ? extends R> next) {
return item -> { O mid = process(item); return mid == null ? null : next.process(mid); };
}
static <T> ItemProcessor<T, T> filter(java.util.function.Predicate<? super T> keep) {
return item -> keep.test(item) ? item : null;
}
}
record RawOrder(String id, String amountText, String status) {}
record Order(String id, long amount, String status) {}
record SettlementLine(String orderId, long fee, long payout) {}
// 단계 1: 검증 — 잘못된 행은 예외가 아니라 null 로 걸러낸다(05 레슨에서 스킵/데드 레터로 발전)
static final ItemProcessor<RawOrder, RawOrder> VALIDATE = r ->
r.id().startsWith("ORD-") && r.amountText().matches("\\d+") ? r : null;
// 단계 2: 변환
static final ItemProcessor<RawOrder, Order> PARSE = r -> new Order(r.id(), Long.parseLong(r.amountText()), r.status());
// 단계 3: 필터 — 결제 완료만 정산
static final ItemProcessor<Order, Order> PAID_ONLY = ItemProcessor.filter(o -> o.status().equals("PAID"));
// 단계 4: 정산 계산 (수수료 3%)
static final ItemProcessor<Order, SettlementLine> SETTLE = o -> new SettlementLine(o.id(), o.amount() * 3 / 100, o.amount() * 97 / 100);
public static void main(String[] args) throws Exception {
ItemProcessor<RawOrder, SettlementLine> pipeline = VALIDATE.andThen(PARSE).andThen(PAID_ONLY).andThen(SETTLE);
List<RawOrder> input = List.of(
new RawOrder("ORD-1", "10000", "PAID"),
new RawOrder("X-2", "5000", "PAID"), // 검증 실패
new RawOrder("ORD-3", "12a00", "PAID"), // 검증 실패 (숫자 아님)
new RawOrder("ORD-4", "7000", "CANCELLED"), // 필터
new RawOrder("ORD-5", "20000", "PAID"));
int filtered = 0;
for (RawOrder r : input) {
SettlementLine out = pipeline.process(r);
if (out == null) { filtered++; System.out.println(r.id() + " → (filtered)"); }
else System.out.println(r.id() + " → " + out);
}
System.out.println("filtered=" + filtered + " / " + input.size());
}
}
// 출력:
// ORD-1 → SettlementLine[orderId=ORD-1, fee=300, payout=9700]
// X-2 → (filtered)
// ORD-3 → (filtered)
// ORD-4 → (filtered)
// ORD-5 → SettlementLine[orderId=ORD-5, fee=600, payout=19400]
// filtered=3 / 5예제 3 의 Step 은 실행이 끝난 뒤 StepExecution 하나만 돌려주므로, "청크마다 금액 합계를 찍고 싶다"거나 "실패한 청크 번호를 알림으로 보내고 싶다"면 Step 코드를 고쳐야 합니다. 청크 생명주기 훅을 StepListener 인터페이스로 빼면 로깅·메트릭·알림을 Step 밖에서 꽂을 수 있습니다.
Spring Batch 의 ChunkListener, StepExecutionListener 와 같은 발상입니다.
import java.util.*;
import java.util.stream.*;
public class ListenerStep {
interface ItemReader<T> extends AutoCloseable {
T read() throws Exception;
default void open() throws Exception {}
@Override default void close() throws Exception {}
}
interface ItemProcessor<I, O> { O process(I item) throws Exception; }
interface ItemWriter<T> { void write(List<? extends T> items) throws Exception; }
// 청크 생명주기 훅. 로깅·모니터링·메트릭을 Step 코드에 섞지 않고 밖에서 꽂는다
interface StepListener<O> {
default void beforeChunk(int chunkNo) {}
default void afterChunk(int chunkNo, List<? extends O> written) {}
default void onError(int chunkNo, Exception e) {}
default void afterStep(String status, long read, long written) {}
}
static <I, O> void execute(ItemReader<I> reader, ItemProcessor<I, O> processor, ItemWriter<O> writer,
int chunkSize, StepListener<O> listener) {
long read = 0, written = 0; int chunkNo = 0;
try {
reader.open();
List<O> chunk = new ArrayList<>(chunkSize);
boolean done = false;
while (!done) {
chunk.clear(); chunkNo++;
listener.beforeChunk(chunkNo);
for (int i = 0; i < chunkSize; i++) {
I item = reader.read();
if (item == null) { done = true; break; }
read++;
O out = processor.process(item);
if (out != null) chunk.add(out);
}
if (!chunk.isEmpty()) { writer.write(chunk); written += chunk.size(); listener.afterChunk(chunkNo, chunk); }
}
listener.afterStep("COMPLETED", read, written);
} catch (Exception e) {
listener.onError(chunkNo, e);
listener.afterStep("FAILED", read, written);
} finally {
try { reader.close(); } catch (Exception ignored) {}
}
}
record Payment(int id, long amount) {}
public static void main(String[] args) {
// 리스너 1: 청크 로그 + 금액 합계 메트릭
long[] totalAmount = {0};
StepListener<Payment> metrics = new StepListener<>() {
public void afterChunk(int no, List<? extends Payment> w) {
long sum = w.stream().mapToLong(Payment::amount).sum();
totalAmount[0] += sum;
System.out.printf(" chunk#%d 커밋 %d건, 금액 %,d%n", no, w.size(), sum);
}
public void onError(int no, Exception e) { System.out.println(" chunk#" + no + " 실패: " + e.getMessage()); }
public void afterStep(String status, long read, long written) {
System.out.printf(" step %s read=%d written=%d 누적 금액 %,d%n", status, read, written, totalAmount[0]);
}
};
List<Payment> payments = IntStream.rangeClosed(1, 7).mapToObj(i -> new Payment(i, i * 1000L)).toList();
Iterator<Payment> it1 = payments.iterator();
System.out.println("정상 실행:");
execute(() -> it1.hasNext() ? it1.next() : null, p -> p, c -> {}, 3, metrics);
totalAmount[0] = 0;
Iterator<Payment> it2 = payments.iterator();
System.out.println("id=5 에서 예외:");
execute(() -> it2.hasNext() ? it2.next() : null,
p -> { if (p.id() == 5) throw new IllegalStateException("id=5 금액 검증 실패"); return p; },
c -> {}, 3, metrics);
}
}
// 출력:
// 정상 실행:
// chunk#1 커밋 3건, 금액 6,000
// chunk#2 커밋 3건, 금액 15,000
// chunk#3 커밋 1건, 금액 7,000
// step COMPLETED read=7 written=7 누적 금액 28,000
// id=5 에서 예외:
// chunk#1 커밋 3건, 금액 6,000
// chunk#2 실패: id=5 금액 검증 실패
// step FAILED read=5 written=3 누적 금액 6,000예제 4 의 JobRunner 는 어떤 Step 이든 실패하면 무조건 중단합니다. 그러나 "알림 발송 실패"가 "정산 파일 생성"을 막을 이유는 없고, 반대로 "거래 집계 실패"는 반드시 이후를 막아야 합니다. Step 마다 OnFail 정책을 두면 이 차이를 선언적으로 표현할 수 있습니다. 어느 쪽이든 실패가 하나라도 있으면 Job 최종 상태는 FAILED 로 남겨 운영자가 아침에 반드시 보게 합니다.
import java.util.*;
import java.util.function.Supplier;
public class FailPolicyJob {
enum Status { COMPLETED, FAILED, SKIPPED }
enum OnFail { STOP, CONTINUE } // Step 실패 시 Job 을 멈출지, 다음 Step 으로 갈지
record StepDef(String name, Supplier<Status> body, OnFail onFail) {}
record StepResult(String name, Status status) {}
// Step 을 순서대로 실행. STOP 정책의 Step 이 실패하면 남은 Step 은 SKIPPED 로 기록하고 Job 은 FAILED.
// CONTINUE 정책의 Step 이 실패하면 기록만 남기고 계속 진행하되, Job 최종 상태는 FAILED 로 표시한다.
static List<StepResult> run(String jobName, List<StepDef> steps) {
List<StepResult> results = new ArrayList<>();
boolean stopped = false, anyFailed = false;
for (StepDef s : steps) {
if (stopped) { results.add(new StepResult(s.name(), Status.SKIPPED)); continue; }
Status st;
try { st = s.body().get(); } catch (RuntimeException e) { st = Status.FAILED; }
results.add(new StepResult(s.name(), st));
if (st == Status.FAILED) { anyFailed = true; if (s.onFail() == OnFail.STOP) stopped = true; }
}
System.out.println("=== JOB " + jobName + " " + (anyFailed ? "FAILED" : "COMPLETED") + " ===");
results.forEach(r -> System.out.println(" " + r.name() + " → " + r.status()));
return results;
}
public static void main(String[] args) {
Supplier<Status> ok = () -> Status.COMPLETED;
Supplier<Status> boom = () -> { throw new IllegalStateException("집계 오류"); };
// 시나리오 1: 알림 발송(CONTINUE) 실패는 정산 파일 생성을 막지 않는다
run("dailyJob-A", List.of(
new StepDef("aggregateTrades", ok, OnFail.STOP),
new StepDef("sendNotifications", boom, OnFail.CONTINUE),
new StepDef("writeSettlementFile", ok, OnFail.STOP)));
// 시나리오 2: 거래 집계(STOP) 실패는 이후 전부 중단 — 불완전한 데이터로 정산 파일을 만들면 안 된다
run("dailyJob-B", List.of(
new StepDef("aggregateTrades", boom, OnFail.STOP),
new StepDef("sendNotifications", ok, OnFail.CONTINUE),
new StepDef("writeSettlementFile", ok, OnFail.STOP)));
}
}
// 출력:
// === JOB dailyJob-A FAILED ===
// aggregateTrades → COMPLETED
// sendNotifications → FAILED
// writeSettlementFile → COMPLETED
// === JOB dailyJob-B FAILED ===
// aggregateTrades → FAILED
// sendNotifications → SKIPPED
// writeSettlementFile → SKIPPED지금까지 Reader 는 List 였습니다. 실제 배치처럼 10만 행 CSV 를 임시 파일로 만들고 BufferedReader 기반 Reader 로 읽되, Writer 에 커밋 고정 비용(200µs, DB fsync 대체)을 넣어 청크 크기 10/100/1,000/10,000 의 소요 시간을 잽니다.
커밋 비용이 지배적일 때 청크 크기가 100→1,000 에서 효과가 크고 그 이상은 평탄해지는 것이 03 레슨의 곡선 그대로입니다. 마지막으로 "헤더만 있는 빈 파일"과 "6번째 청크에서 Writer 예외"를 넣어 커밋 0 건, 그리고 예외가 나도 close() 가 정확히 한 번 호출되는지 확인합니다.
import java.io.*;
import java.nio.charset.StandardCharsets;
import java.nio.file.*;
import java.util.*;
public class ChunkSizeBench {
interface ItemReader<T> extends AutoCloseable {
T read() throws Exception;
@Override default void close() throws Exception {}
}
interface ItemWriter<T> { void write(List<? extends T> items) throws Exception; }
// 파일 Reader: 한 줄씩. close 횟수를 세어 "예외가 나도 항상 닫히는가" 를 검증한다
static class CsvReader implements ItemReader<String> {
static int closeCount = 0;
private final BufferedReader br;
CsvReader(Path p) throws IOException { br = Files.newBufferedReader(p, StandardCharsets.UTF_8); br.readLine(); } // 헤더
public String read() throws IOException { return br.readLine(); }
@Override public void close() throws IOException { br.close(); closeCount++; }
}
// 쓰기 비용 흉내: 커밋마다 고정 비용 200µs (DB fsync 대체) + 건당 파싱
static long runStep(ItemReader<String> reader, int chunkSize, ItemWriter<String> writer) throws Exception {
long commits = 0;
try (reader) {
List<String> chunk = new ArrayList<>(chunkSize);
String line;
while ((line = reader.read()) != null) {
chunk.add(line.substring(0, line.indexOf(','))); // "가공"
if (chunk.size() == chunkSize) { writer.write(chunk); commits++; chunk.clear(); }
}
if (!chunk.isEmpty()) { writer.write(chunk); commits++; }
}
return commits;
}
static void busySleepMicros(long micros) { long end = System.nanoTime() + micros * 1_000; while (System.nanoTime() < end) Thread.onSpinWait(); }
public static void main(String[] args) throws Exception {
Path dir = Files.createTempDirectory("batch01");
Path csv = dir.resolve("orders.csv");
int rows = 100_000;
try (BufferedWriter w = Files.newBufferedWriter(csv, StandardCharsets.UTF_8)) {
w.write("id,customer,amount"); w.newLine();
for (int i = 1; i <= rows; i++) { w.write(i + ",cust-" + (i % 500) + "," + (i % 90_000)); w.newLine(); }
}
ItemWriter<String> writer = chunk -> busySleepMicros(200); // 커밋 고정 비용
runStep(new CsvReader(csv), 1000, writer); // 워밍업
System.out.println("chunkSize commits elapsed");
for (int size : new int[]{10, 100, 1_000, 10_000}) {
long t0 = System.nanoTime();
long commits = runStep(new CsvReader(csv), size, writer);
System.out.printf("%,9d %,7d %,6d ms%n", size, commits, (System.nanoTime() - t0) / 1_000_000);
}
// 엣지 1: 빈 파일(헤더만) → 커밋 0
Path empty = dir.resolve("empty.csv");
Files.writeString(empty, "id,customer,amount\n", StandardCharsets.UTF_8);
System.out.println("빈 입력 commits = " + runStep(new CsvReader(empty), 100, writer));
// 엣지 2: 특정 청크에서 writer 예외 → 그래도 reader 는 닫힌다
int before = CsvReader.closeCount;
try {
runStep(new CsvReader(csv), 1000, chunk -> { if (chunk.get(0).equals("5001")) throw new IllegalStateException("chunk 6 쓰기 실패"); });
} catch (IllegalStateException e) {
System.out.println("예외: " + e.getMessage() + ", reader close 호출 = " + (CsvReader.closeCount - before));
}
}
}
// 출력 (시간 값은 환경에 따라 다름):
// chunkSize commits elapsed
// 10 10,000 4,306 ms
// 100 1,000 848 ms
// 1,000 100 180 ms
// 10,000 10 169 ms
// 빈 입력 commits = 0
// 예외: chunk 6 쓰기 실패, reader close 호출 = 1