아래 코드는 java-src/batch/01_architecture/ 에 있습니다. javac *.java && java Main 으로 실행됩니다.
Reader 는 한 건씩, Writer 는 묶음으로. ItemReader.of(list) 는 테스트용 편의 메서드입니다.
import java.util.Iterator;
import java.util.List;
public interface ItemReader<T> extends AutoCloseable {
T read() throws Exception; // 다음 항목, 끝이면 null
default void open() throws Exception {} // 파일/커서 열기
@Override default void close() throws Exception {}
static <T> ItemReader<T> of(List<T> items) {
Iterator<T> it = items.iterator();
return () -> it.hasNext() ? it.next() : null;
}
}
@FunctionalInterface
public interface ItemProcessor<I, O> {
O process(I item) throws Exception; // null 반환 = 필터링
static <T> ItemProcessor<T, T> identity() { return item -> item; }
}
@FunctionalInterface
public interface ItemWriter<T> {
void write(List<? extends T> items) throws Exception; // 청크 단위
}
// 사용:
// ItemReader<String> r = ItemReader.of(List.of("a", "b"));
// System.out.println(r.read() + " " + r.read() + " " + r.read());
// 출력: a b null실행이 끝나면 값이 바뀌지 않으므로 record 가 적합합니다. summary() 하나로 로그 한 줄이 나옵니다.
import java.time.Duration;
import java.time.LocalDateTime;
public record StepExecution(
String stepName, Status status,
LocalDateTime startTime, LocalDateTime endTime,
long readCount, long filterCount, long writeCount, long commitCount,
String exitMessage) {
public enum Status { STARTED, COMPLETED, FAILED }
public Duration duration() { return Duration.between(startTime, endTime); }
public String summary() {
return String.format("[%s] %s read=%d filtered=%d written=%d commits=%d %dms%s",
stepName, status, readCount, filterCount, writeCount, commitCount,
duration().toMillis(), exitMessage == null ? "" : " msg=" + exitMessage);
}
}
// 출력 예: [regradeMembers] COMPLETED read=1000 filtered=0 written=1000 commits=10 8ms읽기 → 가공 → 청크 버퍼 → N건 모이면 쓰기. 예외가 나면 FAILED 상태와 메시지를 기록하고, 어떤 경우에도 reader 를 닫습니다.
public class Step<I, O> {
private final String name;
private final ItemReader<I> reader;
private final ItemProcessor<I, O> processor;
private final ItemWriter<O> writer;
private final int chunkSize;
public Step(String name, ItemReader<I> reader, ItemProcessor<I, O> processor,
ItemWriter<O> writer, int chunkSize) {
if (chunkSize <= 0) throw new IllegalArgumentException("chunkSize must be > 0");
this.name = name; this.reader = reader; this.processor = processor;
this.writer = writer; this.chunkSize = chunkSize;
}
public StepExecution execute() {
LocalDateTime start = LocalDateTime.now();
long read = 0, filtered = 0, written = 0, commits = 0;
try {
reader.open();
List<O> chunk = new ArrayList<>(chunkSize);
boolean done = false;
while (!done) {
chunk.clear(); // 이전 청크 참조 해제
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) { filtered++; continue; } // 필터
chunk.add(out);
}
if (!chunk.isEmpty()) {
writer.write(chunk); // 커밋 포인트
written += chunk.size();
commits++;
}
}
return new StepExecution(name, StepExecution.Status.COMPLETED, start, LocalDateTime.now(),
read, filtered, written, commits, null);
} catch (Exception e) {
return new StepExecution(name, StepExecution.Status.FAILED, start, LocalDateTime.now(),
read, filtered, written, commits, e.getClass().getSimpleName() + ": " + e.getMessage());
} finally {
try { reader.close(); } catch (Exception ignored) {}
}
}
}
// 1000건, chunk=100 이면: commits=10, 마지막 청크가 꽉 차지 않아도 write 됨public class JobRunner {
public record JobExecution(String jobName, StepExecution.Status status,
LocalDateTime startTime, LocalDateTime endTime,
List<StepExecution> steps) {
public Duration duration() { return Duration.between(startTime, endTime); }
}
private final String jobName;
private final List<Step<?, ?>> steps = new ArrayList<>();
public JobRunner(String jobName) { this.jobName = jobName; }
public JobRunner addStep(Step<?, ?> step) { steps.add(step); return this; }
public JobExecution run() {
LocalDateTime start = LocalDateTime.now();
System.out.printf("=== JOB %s START %s ===%n", jobName, start);
List<StepExecution> results = new ArrayList<>();
StepExecution.Status jobStatus = StepExecution.Status.COMPLETED;
for (Step<?, ?> step : steps) {
StepExecution exec = step.execute();
results.add(exec);
System.out.println(" " + exec.summary());
if (exec.status() == StepExecution.Status.FAILED) {
jobStatus = StepExecution.Status.FAILED;
break; // 뒤 Step 은 실행하지 않음
}
}
JobExecution job = new JobExecution(jobName, jobStatus, start, LocalDateTime.now(), results);
System.out.printf("=== JOB %s %s (%dms) ===%n", jobName, jobStatus, job.duration().toMillis());
return job;
}
}
// 출력:
// === JOB failingJob START 2026-09-08T02:00:00.010 ===
// [failingStep] FAILED read=250 filtered=0 written=200 commits=2 1ms msg=IllegalStateException: id=250 데이터 오류
// === JOB failingJob FAILED (2ms) ===written=200 commits=2 에 주목하세요. 250번째에서 죽었지만 앞선 두 청크(200건)는 이미 Writer 에 전달되었습니다. 이것이 04 레슨에서 다룰 "청크 단위 커밋"의 출발점입니다.
두 Step 을 하나의 Job 으로 묶습니다. 두 번째 Step 은 Processor 가 대부분 null 을 반환해 입력보다 출력이 훨씬 적습니다.
enum Grade { BRONZE, SILVER, GOLD, VIP }
record Member(long id, String name, long yearlyPurchase, Grade grade) {
Member withGrade(Grade g) { return new Member(id, name, yearlyPurchase, g); }
}
record Invoice(long memberId, long amount, LocalDate dueDate, boolean paid) {}
record Notification(long memberId, String message) {}
static Grade calcGrade(long p) {
if (p >= 5_000_000) return Grade.VIP;
if (p >= 2_000_000) return Grade.GOLD;
if (p >= 500_000) return Grade.SILVER;
return Grade.BRONZE;
}
// Step 1: 등급 재산정
List<Member> updated = new ArrayList<>();
Step<Member, Member> regrade = new Step<>("regradeMembers",
ItemReader.of(members),
m -> m.withGrade(calcGrade(m.yearlyPurchase())),
chunk -> updated.addAll(chunk), // 실제로는 UPDATE ... 배치 실행
100);
// Step 2: 미납 알림 대상 추출 (필터링)
LocalDate today = LocalDate.of(2026, 9, 8);
List<Notification> notifications = new ArrayList<>();
Step<Invoice, Notification> unpaid = new Step<>("extractUnpaid",
ItemReader.of(invoices),
inv -> {
if (inv.paid() || !inv.dueDate().isBefore(today)) return null; // 대상 아님
long overdue = today.toEpochDay() - inv.dueDate().toEpochDay();
return new Notification(inv.memberId(),
String.format("%,d원이 %d일 연체되었습니다", inv.amount(), overdue));
},
notifications::addAll,
200);
new JobRunner("dailyMemberJob").addStep(regrade).addStep(unpaid).run();
// 출력:
// === JOB dailyMemberJob START 2026-09-08T02:00:00.123 ===
// --- step regradeMembers ...
// [regradeMembers] COMPLETED read=1000 filtered=0 written=1000 commits=10 8ms
// --- step extractUnpaid ...
// [extractUnpaid] COMPLETED read=1000 filtered=752 written=248 commits=5 3ms
// === JOB dailyMemberJob COMPLETED (15ms) ===filtered=752 written=248 — 메타데이터가 있으니 "1000건 읽었는데 248건만 나온 것이 정상(필터)"임을 바로 알 수 있습니다.