소스는 java-src/project/stage4/(13개 파일). 실행하면 data/ 폴더에 orders.csv(약 5MB)와 정산 파일이 생긴다. 데이터 흐름 순서 — 생성기 → 레코드 → 리더 → 파서 → 게이트웨이 → 재시도 → 상태 → 체크포인트 → 출력 → 잡 → 리포트 → Main — 로 읽는다.
import java.io.BufferedWriter;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.LocalDate;
import java.util.Random;
/**
* 주문 CSV 생성기. 고정 시드(Random(2024))라 실행할 때마다 같은 파일이 나온다.
* 약 1.2%는 일부러 잘못된 행(수량 0/음수, 숫자 아닌 단가, 컬럼 누락, 모르는 결제수단)을 섞는다.
*
* 형식: order_id,member_id,product_id,category,quantity,unit_price,ordered_at,pay_type
*/
public class SampleDataGenerator {
static final String[][] PRODUCTS = { // id, category, unit_price
{"101", "전자", "25000"}, {"102", "전자", "8000"}, {"103", "도서", "32000"},
{"104", "생활", "15000"}, {"105", "문구", "2500"}, {"106", "전자", "89000"},
{"107", "도서", "38000"}, {"108", "생활", "27000"}, {"109", "생활", "12000"}};
static final String[] PAY_TYPES = {"CARD", "CARD", "CARD", "BANK", "BANK", "POINT"}; // 카드 50%, 계좌 33%, 포인트 17%
public static Path generate(Path file, int rows, LocalDate date) throws IOException {
Files.createDirectories(file.getParent());
Random rnd = new Random(2024);
try (BufferedWriter w = Files.newBufferedWriter(file, StandardCharsets.UTF_8)) {
w.write("order_id,member_id,product_id,category,quantity,unit_price,ordered_at,pay_type");
w.newLine();
for (int i = 1; i <= rows; i++) {
String[] p = PRODUCTS[rnd.nextInt(PRODUCTS.length)];
String orderId = String.format("ORD-%07d", i);
int memberId = 1 + rnd.nextInt(500);
int qty = 1 + rnd.nextInt(5);
String payType = PAY_TYPES[rnd.nextInt(PAY_TYPES.length)];
String price = p[2];
double r = rnd.nextDouble(); // 오염 주입
if (r < 0.004) qty = 0;
else if (r < 0.007) price = "N/A";
else if (r < 0.009) payType = null; // 컬럼 누락
else if (r < 0.011) qty = -qty;
else if (r < 0.012) payType = "BITCOIN";
StringBuilder sb = new StringBuilder()
.append(orderId).append(',').append(memberId).append(',').append(p[0]).append(',')
.append(p[1]).append(',').append(qty).append(',').append(price).append(',').append(date);
if (payType != null) sb.append(',').append(payType);
w.write(sb.toString());
w.newLine();
}
}
return file;
}
}상품 9개는 3단계와 같다. 오염 비율은 수량 0(0.4%), 단가 문자열(0.3%), 컬럼 누락(0.2%), 음수 수량(0.2%), 결제수단 오류(0.1%)로 합계 1.2%. 리포트의 스킵 사유 건수가 이 비율과 맞는지 4절에서 확인한다.
import java.time.LocalDate;
/** CSV 한 행을 파싱·검증한 결과. 여기까지 왔으면 "정상 주문"이다. */
public record OrderRecord(long lineNo, String orderId, int memberId, int productId, String category,
int quantity, long unitPrice, LocalDate orderedAt, PayType payType) {
public enum PayType {
CARD(250), BANK(100), POINT(0); // 수수료율 (bp: 1/10000)
public final int feeBasisPoints;
PayType(int feeBasisPoints) { this.feeBasisPoints = feeBasisPoints; }
}
public long amount() { return quantity * unitPrice; }
public long fee() { return amount() * payType.feeBasisPoints / 10_000; }
}수수료율을 enum 상수의 필드로 둔 것은 "결제수단별 수수료"라는 규칙이 결제수단 정의 옆에 있어야 하기 때문이다. 카드 2.5%, 계좌 1.0%, 포인트 0%. 정산액 = 매출 − 수수료.
/** 이 행은 처리할 수 없으니 건너뛰라는 신호. reason은 집계 키로 쓰인다. */
public class SkipException extends Exception {
private final String reason;
public SkipException(String reason, String detail) {
super(reason + ": " + detail);
this.reason = reason;
}
public String reason() { return reason; }
}2·3단계의 ShopException은 unchecked였다. SkipException은 checked다. "이 행을 건너뛴다"는 것은 잡의 정상 흐름이지 버그가 아니고, 호출자(SettlementJob)가 반드시 처리해야 하므로 컴파일러가 강제하게 했다.
import java.io.BufferedReader;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
/**
* CSV를 청크 단위로 읽는다. 파일 전체를 메모리에 올리지 않고 BufferedReader로 한 줄씩 읽어
* chunkSize개씩 묶어 돌려준다. startAfterLine까지는 건너뛴다(체크포인트 재시작).
*/
public class ChunkReader implements AutoCloseable {
public record RawLine(long lineNo, String text) {}
private final BufferedReader reader;
private long lineNo = 0; // 헤더 = 1행
public ChunkReader(Path file, long startAfterLine) throws IOException {
this.reader = Files.newBufferedReader(file, StandardCharsets.UTF_8);
String header = reader.readLine(); // 헤더 소비
lineNo = 1;
if (header == null) throw new IOException("빈 파일: " + file);
while (lineNo < startAfterLine) { // 이미 커밋된 행은 읽고 버린다
if (reader.readLine() == null) break;
lineNo++;
}
}
/** 다음 청크. 파일 끝이면 빈 리스트. */
public List<RawLine> nextChunk(int chunkSize) throws IOException {
List<RawLine> chunk = new ArrayList<>(chunkSize);
String line;
while (chunk.size() < chunkSize && (line = reader.readLine()) != null) {
lineNo++;
if (line.isBlank()) continue;
chunk.add(new RawLine(lineNo, line));
}
return chunk;
}
public long currentLine() { return lineNo; }
@Override
public void close() throws IOException { reader.close(); }
}재시작 시 startAfterLine까지 읽고 버린다. 파일 오프셋(바이트)을 기억하면 seek로 건너뛸 수 있지만 행 번호가 더 직관적이고, 4만 행을 읽고 버리는 비용은 밀리초 단위다. 행 번호를 RawLine에 함께 담아 스킵 메시지에 "몇 행"인지 남긴다.
import java.time.LocalDate;
import java.time.format.DateTimeParseException;
/** 문자열 행 → OrderRecord. 규칙 위반은 SkipException(사유 코드)으로 보고한다. */
public class OrderParser {
private static final int FIELD_COUNT = 8;
public OrderRecord parse(ChunkReader.RawLine raw) throws SkipException {
String[] f = raw.text().split(",", -1);
if (f.length != FIELD_COUNT)
throw new SkipException("FIELD_COUNT", raw.lineNo() + "행 컬럼 " + f.length + "개");
int quantity;
long unitPrice;
int memberId;
int productId;
try {
memberId = Integer.parseInt(f[1]);
productId = Integer.parseInt(f[2]);
quantity = Integer.parseInt(f[4]);
unitPrice = Long.parseLong(f[5]);
} catch (NumberFormatException e) {
throw new SkipException("NUMBER_FORMAT", raw.lineNo() + "행 " + e.getMessage());
}
if (quantity <= 0)
throw new SkipException("QUANTITY", raw.lineNo() + "행 수량 " + quantity);
if (unitPrice <= 0)
throw new SkipException("UNIT_PRICE", raw.lineNo() + "행 단가 " + unitPrice);
LocalDate date;
try {
date = LocalDate.parse(f[6]);
} catch (DateTimeParseException e) {
throw new SkipException("DATE", raw.lineNo() + "행 " + f[6]);
}
OrderRecord.PayType payType;
try {
payType = OrderRecord.PayType.valueOf(f[7]);
} catch (IllegalArgumentException e) {
throw new SkipException("PAY_TYPE", raw.lineNo() + "행 " + f[7]);
}
return new OrderRecord(raw.lineNo(), f[0], memberId, productId, f[3], quantity, unitPrice, date, payType);
}
}split(",", -1)의 -1은 끝의 빈 컬럼을 버리지 않게 한다. 사유 코드(FIELD_COUNT, NUMBER_FORMAT, ...)는 리포트의 집계 키가 되므로 자유 문장이 아니라 고정 문자열이다. 검사 순서가 값 → 형식 → 의미 순인 것도 의도적이다. 컬럼 수가 틀리면 숫자 파싱을 시도할 이유가 없다.
/**
* 결제 게이트웨이 시뮬레이션. 실제라면 네트워크 호출이다.
* - 일시 장애(TransientFailure): 약 5%. 같은 주문을 다시 보내면 성공할 수 있다 → 재시도 대상.
* - 영구 거절(Declined): 약 0.1%. 몇 번을 보내도 같다 → 재시도 무의미, 스킵 대상.
* 실패 여부를 Random이 아니라 (주문ID, 시도 횟수)의 해시로 정하므로 실행마다 결과가 같다.
*/
public class PaymentGateway {
public static class TransientFailure extends Exception {
public TransientFailure(String msg) { super(msg); }
}
public static class Declined extends Exception {
public Declined(String msg) { super(msg); }
}
private int calls = 0;
/** 승인 코드를 돌려준다. */
public String charge(OrderRecord order, int attempt) throws TransientFailure, Declined {
calls++;
long h = mix(order.orderId().hashCode(), 0);
if (Math.floorMod(h, 1000) == 0)
throw new Declined(order.orderId() + " 한도 초과");
long a = mix(order.orderId().hashCode(), attempt);
if (Math.floorMod(a, 100) < 5)
throw new TransientFailure(order.orderId() + " 타임아웃 (시도 " + attempt + ")");
return "APV-" + Long.toHexString(a & 0xFFFFFF);
}
public int calls() { return calls; }
private static long mix(long id, long attempt) {
long x = id * 0x9E3779B97F4A7C15L + attempt * 0xBF58476D1CE4E5B9L;
x ^= (x >>> 31);
x *= 0x94D049BB133111EBL;
return x ^ (x >>> 29);
}
}영구 거절은 attempt를 0으로 고정한 해시라 시도 횟수와 무관하게 같은 주문은 항상 거절된다. 일시 장애는 attempt를 포함한 해시라 1차에 실패한 주문도 2차에는 95% 확률로 성공한다. 두 예외는 checked이고 서로 다른 타입이다 — 호출자가 "재시도할 것"과 "포기할 것"을 타입으로 구분하게 하기 위해서다.
/**
* 일시 장애를 최대 maxAttempts번까지 지수 백오프로 재시도한다.
* 결과: 성공(승인 코드) / 재시도 소진(SkipException RETRY_EXHAUSTED) / 영구 거절(SkipException DECLINED).
*/
public class RetryTemplate {
private final PaymentGateway gateway;
private final int maxAttempts;
private final long baseDelayMs;
private int retries = 0; // 실제로 재시도한 횟수 (통계용)
public RetryTemplate(PaymentGateway gateway, int maxAttempts, long baseDelayMs) {
this.gateway = gateway;
this.maxAttempts = maxAttempts;
this.baseDelayMs = baseDelayMs;
}
public String chargeWithRetry(OrderRecord order) throws SkipException {
for (int attempt = 1; ; attempt++) {
try {
return gateway.charge(order, attempt);
} catch (PaymentGateway.Declined e) {
throw new SkipException("DECLINED", e.getMessage()); // 즉시 포기
} catch (PaymentGateway.TransientFailure e) {
if (attempt >= maxAttempts)
throw new SkipException("RETRY_EXHAUSTED", e.getMessage());
retries++;
backoff(attempt); // 1x, 2x, 4x ...
}
}
}
private void backoff(int attempt) {
long delay = baseDelayMs * (1L << (attempt - 1));
if (delay <= 0) return;
try {
Thread.sleep(delay);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
public int retries() { return retries; }
}게이트웨이의 두 예외를 하나의 SkipException으로 정규화해서 돌려준다. 잡은 "왜 실패했는지"를 사유 코드로만 알면 되고, 게이트웨이 예외 타입을 알 필요가 없다. 지연은 base × 2^(attempt−1). 시연에서는 baseDelayMs = 0이라 대기가 없다(10만 건 × 5% × 대기 시간은 학습용으로 너무 길다).
import java.util.Map;
import java.util.Properties;
import java.util.TreeMap;
/**
* 배치의 누적 상태 = 마지막 커밋 행 + 카운터 + 정산 집계.
* 청크 하나를 처리할 때는 "빈 JobState"에 작업하고, 성공하면 merge(커밋), 실패하면 버린다(롤백).
* Properties로 직렬화해 체크포인트 파일에 저장/복원한다.
*/
public class JobState {
/** 정산 집계 한 줄: 건수·매출·수수료. 정산액(net) = 매출 - 수수료. */
public static class Entry {
public long count, gross, fee;
void add(long amount, long feeAmount) { count++; gross += amount; fee += feeAmount; }
void merge(Entry o) { count += o.count; gross += o.gross; fee += o.fee; }
public long net() { return gross - fee; }
String serialize() { return count + "," + gross + "," + fee; }
static Entry parse(String s) {
String[] f = s.split(",");
Entry e = new Entry();
e.count = Long.parseLong(f[0]); e.gross = Long.parseLong(f[1]); e.fee = Long.parseLong(f[2]);
return e;
}
}
public long lastCommittedLine = 1; // 헤더까지 읽은 상태
public int chunksCommitted = 0;
public long read = 0, processed = 0, skipped = 0, retried = 0;
public final Map<String, Long> skipReasons = new TreeMap<>();
public final Map<String, Entry> byCategory = new TreeMap<>();
public final Map<String, Entry> byPayType = new TreeMap<>();
public void recordSuccess(OrderRecord o) {
processed++;
byCategory.computeIfAbsent(o.category(), k -> new Entry()).add(o.amount(), o.fee());
byPayType.computeIfAbsent(o.payType().name(), k -> new Entry()).add(o.amount(), o.fee());
}
public void recordSkip(String reason) {
skipped++;
skipReasons.merge(reason, 1L, Long::sum);
}
/** 청크 작업 결과를 누적 상태에 반영 = 커밋. */
public void merge(JobState work) {
read += work.read; processed += work.processed; skipped += work.skipped; retried += work.retried;
work.skipReasons.forEach((k, v) -> skipReasons.merge(k, v, Long::sum));
work.byCategory.forEach((k, v) -> byCategory.computeIfAbsent(k, x -> new Entry()).merge(v));
work.byPayType.forEach((k, v) -> byPayType.computeIfAbsent(k, x -> new Entry()).merge(v));
lastCommittedLine = work.lastCommittedLine;
chunksCommitted++;
}
public long totalGross() { return byCategory.values().stream().mapToLong(e -> e.gross).sum(); }
public long totalFee() { return byCategory.values().stream().mapToLong(e -> e.fee).sum(); }
// ---- 직렬화 ----
public Properties toProperties() {
Properties p = new Properties();
p.setProperty("lastCommittedLine", String.valueOf(lastCommittedLine));
p.setProperty("chunksCommitted", String.valueOf(chunksCommitted));
p.setProperty("read", String.valueOf(read));
p.setProperty("processed", String.valueOf(processed));
p.setProperty("skipped", String.valueOf(skipped));
p.setProperty("retried", String.valueOf(retried));
skipReasons.forEach((k, v) -> p.setProperty("skip." + k, String.valueOf(v)));
byCategory.forEach((k, v) -> p.setProperty("cat." + k, v.serialize()));
byPayType.forEach((k, v) -> p.setProperty("pay." + k, v.serialize()));
return p;
}
public static JobState fromProperties(Properties p) {
JobState s = new JobState();
s.lastCommittedLine = Long.parseLong(p.getProperty("lastCommittedLine"));
s.chunksCommitted = Integer.parseInt(p.getProperty("chunksCommitted"));
s.read = Long.parseLong(p.getProperty("read"));
s.processed = Long.parseLong(p.getProperty("processed"));
s.skipped = Long.parseLong(p.getProperty("skipped"));
s.retried = Long.parseLong(p.getProperty("retried"));
for (String key : p.stringPropertyNames()) {
if (key.startsWith("skip.")) s.skipReasons.put(key.substring(5), Long.parseLong(p.getProperty(key)));
else if (key.startsWith("cat.")) s.byCategory.put(key.substring(4), Entry.parse(p.getProperty(key)));
else if (key.startsWith("pay.")) s.byPayType.put(key.substring(4), Entry.parse(p.getProperty(key)));
}
return s;
}
}merge가 이 클래스의 핵심이다. 카운터는 더하고, 맵은 키별로 합치고, lastCommittedLine은 덮어쓰고, chunksCommitted는 1 증가. 집계 맵이 TreeMap인 것은 리포트와 체크포인트 파일의 순서를 고정하기 위해서다.
필드가 public인 것은 이 클래스가 캡슐화할 규칙이 없는 순수한 상태 덩어리이기 때문이다 — 2단계에서 getter를 강조했던 것과 대비된다. 규칙이 없으면 getter도 필요 없다.
import java.io.IOException;
import java.io.Reader;
import java.io.Writer;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.util.Optional;
import java.util.Properties;
/**
* 체크포인트 파일 저장/복원. 저장은 "임시 파일에 쓰고 원자적으로 교체"라서
* 쓰는 도중 프로세스가 죽어도 반쯤 쓰인 체크포인트가 남지 않는다.
*/
public class CheckpointStore {
private final Path file;
public CheckpointStore(Path file) { this.file = file; }
public void save(JobState state) throws IOException {
Path tmp = file.resolveSibling(file.getFileName() + ".tmp");
try (Writer w = Files.newBufferedWriter(tmp, StandardCharsets.UTF_8)) {
state.toProperties().store(w, "settlement checkpoint");
}
AtomicMove.replace(tmp, file); // Windows 잠금 재시도 포함
}
public Optional<JobState> load() throws IOException {
if (!Files.exists(file)) return Optional.empty();
Properties p = new Properties();
try (Reader r = Files.newBufferedReader(file, StandardCharsets.UTF_8)) {
p.load(r);
}
return Optional.of(JobState.fromProperties(p));
}
public void delete() throws IOException { Files.deleteIfExists(file); }
}import java.io.BufferedWriter;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.util.Map;
/**
* 정산 결과 파일 출력. 임시 파일(.tmp)에 전부 쓴 뒤 최종 이름으로 원자적 이동.
* 다운스트림(회계 시스템 등)이 "파일이 존재한다 = 완성됐다"로 믿을 수 있게 하기 위함이다.
*/
public class SettlementWriter {
public Path write(Path target, JobState state) throws IOException {
Files.createDirectories(target.getParent());
Path tmp = target.resolveSibling(target.getFileName() + ".tmp");
try (BufferedWriter w = Files.newBufferedWriter(tmp, StandardCharsets.UTF_8)) {
w.write("group,key,count,gross,fee,net");
w.newLine();
writeGroup(w, "CATEGORY", state.byCategory);
writeGroup(w, "PAY_TYPE", state.byPayType);
w.write(String.format("TOTAL,ALL,%d,%d,%d,%d",
state.processed, state.totalGross(), state.totalFee(), state.totalGross() - state.totalFee()));
w.newLine();
}
AtomicMove.replace(tmp, target);
return target;
}
private void writeGroup(BufferedWriter w, String group, Map<String, JobState.Entry> entries) throws IOException {
for (Map.Entry<String, JobState.Entry> e : entries.entrySet()) {
JobState.Entry v = e.getValue();
w.write(String.format("%s,%s,%d,%d,%d,%d", group, e.getKey(), v.count, v.gross, v.fee, v.net()));
w.newLine();
}
}
}두 클래스의 save/write가 같은 패턴이다: .tmp에 쓰고 → try-with-resources로 닫아 flush 보장 → ATOMIC_MOVE. load는 Optional<JobState>를 돌려준다 — "체크포인트가 없음"은 첫 실행의 정상 상황이지 오류가 아니다. 3단계 Repository.findById와 같은 이유다.
Files.move(..., ATOMIC_MOVE) 를 직접 부르지 않고 AtomicMove.replace 로 감싼 이유가 있다. Windows 에서는 대상 파일을 백신이나 OneDrive 같은 동기화 프로그램이 검사하느라 잠시 잡고 있으면 rename 이 AccessDeniedException 으로 거부된다.
리눅스 서버에서는 거의 없는 일이지만 개발 PC 에서 체크포인트를 20번째 저장할 때 갑자기 터지는 식으로 나타난다. 잠금은 보통 수십 ms 안에 풀리므로 짧게 재시도하면 된다.
public final class AtomicMove {
public static void replace(Path tmp, Path target) throws IOException {
for (int attempt = 1; ; attempt++) {
try {
Files.move(tmp, target, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE);
return;
} catch (AccessDeniedException e) {
if (attempt >= 5) throw e; // 0.1+0.2+0.3+0.4 = 1초 후 포기
try { Thread.sleep(100L * attempt); }
catch (InterruptedException ie) { Thread.currentThread().interrupt(); throw e; }
}
}
}
}배치 05 의 재시도 패턴(지수 백오프, 최대 횟수, 포기 시 원래 예외 전파)이 파일 시스템 호출에도 그대로 적용된다. 체크포인트와 정산 파일 두 곳이 같은 메서드를 거치므로 고칠 곳도 한 곳이다.
import java.io.IOException;
import java.nio.file.Path;
import java.util.List;
/**
* 일마감 정산 잡. 청크 = 트랜잭션 단위.
*
* for each chunk:
* work = new JobState() ← 트랜잭션 시작 (작업용 상태)
* for each line: parse → charge(retry) → aggregate (실패 행은 skip)
* committed.merge(work); checkpoint.save(committed) ← 커밋
* 예외 발생 시 work를 버린다 ← 롤백 (committed는 그대로)
*/
public class SettlementJob {
/** 장애 주입용: 이 청크를 처리하던 중 프로세스가 죽는 상황을 흉내 낸다. */
public static class SimulatedCrash extends RuntimeException {
public SimulatedCrash(String msg) { super(msg); }
}
private final Path input;
private final int chunkSize;
private final CheckpointStore checkpoint;
private final OrderParser parser = new OrderParser();
private final PaymentGateway gateway = new PaymentGateway();
private final RetryTemplate retry;
private final int crashAtChunk; // 0이면 장애 없음
public SettlementJob(Path input, int chunkSize, CheckpointStore checkpoint, int crashAtChunk) {
this.input = input;
this.chunkSize = chunkSize;
this.checkpoint = checkpoint;
this.crashAtChunk = crashAtChunk;
this.retry = new RetryTemplate(gateway, 3, 0); // 시연이므로 백오프 대기 0ms
}
public JobState run() throws IOException {
JobState committed = checkpoint.load().orElseGet(JobState::new);
if (committed.chunksCommitted > 0) {
System.out.printf("체크포인트 발견: 청크 %d개 커밋됨, %d행부터 재개%n",
committed.chunksCommitted, committed.lastCommittedLine + 1);
}
try (ChunkReader reader = new ChunkReader(input, committed.lastCommittedLine)) {
List<ChunkReader.RawLine> chunk;
while (!(chunk = reader.nextChunk(chunkSize)).isEmpty()) {
int chunkNo = committed.chunksCommitted + 1;
JobState work = new JobState(); // ---- 트랜잭션 시작 ----
try {
processChunk(chunkNo, chunk, work);
work.lastCommittedLine = reader.currentLine();
committed.merge(work); // ---- 커밋 (메모리) ----
checkpoint.save(committed); // ---- 커밋 (디스크) ----
} catch (RuntimeException e) {
// work는 버려진다 = 롤백. committed와 체크포인트는 마지막 성공 청크 그대로.
System.out.printf("청크 %d 롤백: %s (작업 중이던 %d행 폐기, 체크포인트는 청크 %d 유지)%n",
chunkNo, e.getMessage(), work.read, committed.chunksCommitted);
throw e;
}
if (chunkNo % 20 == 0) {
System.out.printf(" 청크 %3d 커밋 | 누적 읽음 %,7d 처리 %,7d 스킵 %,5d 재시도 %,5d%n",
chunkNo, committed.read, committed.processed, committed.skipped, committed.retried);
}
}
}
return committed;
}
private void processChunk(int chunkNo, List<ChunkReader.RawLine> chunk, JobState work) {
int i = 0;
for (ChunkReader.RawLine raw : chunk) {
if (chunkNo == crashAtChunk && i == chunk.size() / 2)
throw new SimulatedCrash("전원 차단 시뮬레이션");
i++;
work.read++;
try {
OrderRecord order = parser.parse(raw); // 검증 실패 → SkipException
int before = retry.retries();
retry.chargeWithRetry(order); // 결제 실패 → SkipException
work.retried += retry.retries() - before;
work.recordSuccess(order);
} catch (SkipException e) {
work.recordSkip(e.reason());
if (work.skipped <= 2 && chunkNo <= 1) { // 첫 청크의 처음 두 건만 예시로 출력
System.out.println(" 스킵 예시: " + e.getMessage());
}
}
}
}
public int gatewayCalls() { return gateway.calls(); }
}run의 try 블록 네 줄이 트랜잭션이다. processChunk가 정상 종료하면 merge + save. RuntimeException이 나오면 work는 참조를 잃고 GC 대상이 되며, committed는 손대지 않았다. 롤백에 코드가 필요 없다 — 작업 중 상태를 처음부터 별도 객체에 두었기 때문이다. DB라면 ROLLBACK 명령이 필요하지만, 메모리 상태는 "합치지 않으면" 그것이 롤백이다.
processChunk 안에서 SkipException은 잡고 RuntimeException은 잡지 않는다. 스킵은 행 단위 실패(계속 진행), 런타임 예외는 청크 단위 실패(롤백)다. 두 층위를 구분하는 것이 배치 오류 처리의 기본이다.
/** 최종 리포트 출력. */
public class BatchReport {
public static void print(JobState s, long elapsedMs) {
System.out.println("----- 정산 리포트 -----");
System.out.printf("커밋 청크: %d, 마지막 커밋 행: %,d%n", s.chunksCommitted, s.lastCommittedLine);
System.out.printf("읽음 %,d / 처리 %,d / 스킵 %,d (%.2f%%) / 재시도 %,d%n",
s.read, s.processed, s.skipped, 100.0 * s.skipped / s.read, s.retried);
System.out.println("스킵 사유:");
s.skipReasons.forEach((k, v) -> System.out.printf(" %-16s %,6d%n", k, v));
System.out.println("카테고리별 정산:");
s.byCategory.forEach((k, v) -> System.out.printf(" %-6s %,7d건 매출 %,15d 수수료 %,12d 정산 %,15d%n",
k, v.count, v.gross, v.fee, v.net()));
System.out.println("결제수단별 정산:");
s.byPayType.forEach((k, v) -> System.out.printf(" %-6s %,7d건 매출 %,15d 수수료 %,12d 정산 %,15d%n",
k, v.count, v.gross, v.fee, v.net()));
System.out.printf("총 매출 %,d / 총 수수료 %,d / 총 정산액 %,d%n",
s.totalGross(), s.totalFee(), s.totalGross() - s.totalFee());
System.out.printf("소요 시간: %dms (환경에 따라 다름)%n", elapsedMs);
}
}import java.nio.file.Files;
import java.nio.file.Path;
import java.time.LocalDate;
import java.util.List;
/**
* 4단계: 일마감 정산 배치.
* javac -encoding UTF-8 *.java && java -Dstdout.encoding=UTF-8 Main
*
* 1) 주문 CSV 10만 건 생성 (data/orders.csv)
* 2) 1차 실행: 청크 40에서 장애 → 롤백 후 종료
* 3) 2차 실행: 체크포인트에서 재개 → 완료 → 정산 파일 출력
* 4) 검증: 처음부터 한 번에 돌린 결과와 동일한지 비교
*/
public class Main {
static final int ROWS = 100_000;
static final int CHUNK_SIZE = 1_000;
static final LocalDate BUSINESS_DATE = LocalDate.of(2024, 3, 31);
public static void main(String[] args) throws Exception {
Path dataDir = Path.of("data");
Path input = dataDir.resolve("orders.csv");
Path checkpointFile = dataDir.resolve("settlement.checkpoint");
Path output = dataDir.resolve("settlement_" + BUSINESS_DATE + ".csv");
System.out.println("=== 0. 샘플 데이터 생성 ===");
SampleDataGenerator.generate(input, ROWS, BUSINESS_DATE);
System.out.printf("%s (%,d 행, %,d KB)%n", input, ROWS, Files.size(input) / 1024);
CheckpointStore checkpoint = new CheckpointStore(checkpointFile);
checkpoint.delete(); // 이전 실행 흔적 제거
System.out.println();
System.out.println("=== 1. 1차 실행: 청크 40 처리 중 장애 ===");
long t0 = System.nanoTime();
try {
new SettlementJob(input, CHUNK_SIZE, checkpoint, 40).run();
} catch (SettlementJob.SimulatedCrash e) {
System.out.println("잡 비정상 종료: " + e.getMessage());
}
JobState afterCrash = checkpoint.load().orElseThrow();
System.out.printf("체크포인트 상태: 청크 %d, 행 %,d, 처리 %,d, 스킵 %,d%n",
afterCrash.chunksCommitted, afterCrash.lastCommittedLine, afterCrash.processed, afterCrash.skipped);
System.out.println();
System.out.println("=== 2. 2차 실행: 체크포인트에서 재개 ===");
SettlementJob resumed = new SettlementJob(input, CHUNK_SIZE, checkpoint, 0);
JobState result = resumed.run();
long elapsed = (System.nanoTime() - t0) / 1_000_000;
Path written = new SettlementWriter().write(output, result);
System.out.println("정산 파일: " + written + " (" + Files.size(written) + " bytes)");
checkpoint.delete(); // 완료된 잡의 체크포인트는 지운다
System.out.println();
System.out.println("=== 3. 최종 리포트 ===");
BatchReport.print(result, elapsed);
System.out.println();
System.out.println("=== 4. 검증: 장애 없이 한 번에 돌린 결과와 비교 ===");
JobState straight = new SettlementJob(input, CHUNK_SIZE, new CheckpointStore(dataDir.resolve("verify.checkpoint")), 0).run();
Files.deleteIfExists(dataDir.resolve("verify.checkpoint"));
boolean same = straight.processed == result.processed
&& straight.skipped == result.skipped
&& straight.totalGross() == result.totalGross()
&& straight.totalFee() == result.totalFee()
&& straight.skipReasons.equals(result.skipReasons);
System.out.println("재시작 결과 == 단일 실행 결과: " + same);
System.out.println("게이트웨이 호출 수 (재개 실행): " + resumed.gatewayCalls());
System.out.println();
System.out.println("=== 5. 정산 파일 내용 ===");
List<String> lines = Files.readAllLines(written);
lines.forEach(System.out::println);
}
}// 출력:
=== 0. 샘플 데이터 생성 ===
data\orders.csv (100,000 행, 5,050 KB)
=== 1. 1차 실행: 청크 40 처리 중 장애 ===
스킵 예시: QUANTITY: 28행 수량 -5
스킵 예시: QUANTITY: 58행 수량 0
청크 20 커밋 | 누적 읽음 20,000 처리 19,735 스킵 265 재시도 1,018
청크 40 롤백: 전원 차단 시뮬레이션 (작업 중이던 500행 폐기, 체크포인트는 청크 39 유지)
잡 비정상 종료: 전원 차단 시뮬레이션
체크포인트 상태: 청크 39, 행 39,001, 처리 38,495, 스킵 505
=== 2. 2차 실행: 체크포인트에서 재개 ===
체크포인트 발견: 청크 39개 커밋됨, 39002행부터 재개
청크 40 커밋 | 누적 읽음 40,000 처리 39,486 스킵 514 재시도 2,025
청크 60 커밋 | 누적 읽음 60,000 처리 59,209 스킵 791 재시도 2,999
청크 80 커밋 | 누적 읽음 80,000 처리 78,959 스킵 1,041 재시도 4,005
청크 100 커밋 | 누적 읽음 100,000 처리 98,723 스킵 1,277 재시도 5,053
정산 파일: data\settlement_2024-03-31.csv (442 bytes)
=== 3. 최종 리포트 ===
----- 정산 리포트 -----
커밋 청크: 100, 마지막 커밋 행: 100,001
읽음 100,000 / 처리 98,723 / 스킵 1,277 (1.28%) / 재시도 5,053
스킵 사유:
DECLINED 95
FIELD_COUNT 194
NUMBER_FORMAT 305
PAY_TYPE 102
QUANTITY 571
RETRY_EXHAUSTED 10
카테고리별 정산:
도서 21,849건 매출 2,291,974,000 수수료 36,195,040 정산 2,255,778,960
문구 11,030건 매출 83,045,000 수수료 1,308,662 정산 81,736,338
생활 33,167건 매출 1,793,259,000 수수료 28,348,665 정산 1,764,910,335
전자 32,677건 매출 3,960,146,000 수수료 62,558,665 정산 3,897,587,335
결제수단별 정산:
BANK 32,947건 매출 2,697,682,000 수수료 26,976,820 정산 2,670,705,180
CARD 49,207건 매출 4,057,434,500 수수료 101,434,212 정산 3,956,000,288
POINT 16,569건 매출 1,373,307,500 수수료 0 정산 1,373,307,500
총 매출 8,128,424,000 / 총 수수료 128,411,032 / 총 정산액 8,000,012,968
소요 시간: 29136ms (환경에 따라 다름)
=== 4. 검증: 장애 없이 한 번에 돌린 결과와 비교 ===
스킵 예시: QUANTITY: 28행 수량 -5
스킵 예시: QUANTITY: 58행 수량 0
청크 20 커밋 | 누적 읽음 20,000 처리 19,735 스킵 265 재시도 1,018
청크 40 커밋 | 누적 읽음 40,000 처리 39,486 스킵 514 재시도 2,025
청크 60 커밋 | 누적 읽음 60,000 처리 59,209 스킵 791 재시도 2,999
청크 80 커밋 | 누적 읽음 80,000 처리 78,959 스킵 1,041 재시도 4,005
청크 100 커밋 | 누적 읽음 100,000 처리 98,723 스킵 1,277 재시도 5,053
재시작 결과 == 단일 실행 결과: true
게이트웨이 호출 수 (재개 실행): 63387
=== 5. 정산 파일 내용 ===
group,key,count,gross,fee,net
CATEGORY,도서,21849,2291974000,36195040,2255778960
CATEGORY,문구,11030,83045000,1308662,81736338
CATEGORY,생활,33167,1793259000,28348665,1764910335
CATEGORY,전자,32677,3960146000,62558665,3897587335
PAY_TYPE,BANK,32947,2697682000,26976820,2670705180
PAY_TYPE,CARD,49207,4057434500,101434212,3956000288
PAY_TYPE,POINT,16569,1373307500,0,1373307500
TOTAL,ALL,98723,8128424000,128411032,8000012968소요 시간과 파일 경로의 구분자(\ / /)만 환경에 따라 다르고, 나머지 모든 숫자는 실행마다 같다.