java-src/batch/04_transaction/ 에서 javac *.java && java Main. 기본 1,000건 이체, 청크 100 입니다.
pending 버퍼 하나로 트랜잭션의 계약(격리, 취소 가능)을 구현합니다.
public class InMemoryDatabase {
private final Map<String, Long> committed = new HashMap<>();
private Map<String, Long> pending; // null 이면 트랜잭션 밖
private int commitCount, rollbackCount;
public void begin() {
if (pending != null) throw new IllegalStateException("이미 트랜잭션 진행 중");
pending = new HashMap<>();
}
public void commit() { requireTx(); committed.putAll(pending); pending = null; commitCount++; }
public void rollback() { requireTx(); pending = null; rollbackCount++; } // 변경분 폐기
public Long get(String key) { // 자기 변경 우선 (read-your-writes)
if (pending != null && pending.containsKey(key)) return pending.get(key);
return committed.get(key);
}
public long getOrDefault(String key, long def) { Long v = get(key); return v == null ? def : v; }
public boolean contains(String key) { return get(key) != null; }
public void put(String key, long value) { requireTx(); pending.put(key, value); }
public boolean inTransaction() { return pending != null; }
private void requireTx() {
if (pending == null) throw new IllegalStateException("begin() 없이 쓰기/커밋 시도");
}
}
// 사용
db.putCommitted("balance:A", 100);
db.begin();
db.put("balance:A", 50);
System.out.println(db.get("balance:A")); // 50 (자기 변경은 보인다)
db.rollback();
System.out.println(db.get("balance:A")); // 100 (변경 폐기)
db.begin(); db.put("balance:A", 50); db.commit();
System.out.println(db.get("balance:A")); // 50 (확정)
db.put("balance:A", 0); // IllegalStateException: begin() 없이 쓰기/커밋 시도JDBC 로 옮기면 db.begin() → conn.setAutoCommit(false), db.put() → PreparedStatement.executeUpdate(), db.commit() → conn.commit() 입니다. 실제 배치에서는 청크 하나에 커넥션 하나를 잡고 청크 끝에 커밋합니다.
public record Checkpoint(String jobName, long lastCommittedLine, int lastCommittedChunk) {
public static Path fileFor(Path dir, String jobName) { return dir.resolve(jobName + ".checkpoint"); }
public static Optional<Checkpoint> load(Path file) throws IOException {
if (!Files.exists(file)) return Optional.empty();
List<String> lines = Files.readAllLines(file, StandardCharsets.UTF_8);
if (lines.size() < 3) return Optional.empty();
return Optional.of(new Checkpoint(lines.get(0), Long.parseLong(lines.get(1)), Integer.parseInt(lines.get(2))));
}
public void save(Path file) throws IOException {
Path tmp = file.resolveSibling(file.getFileName() + ".tmp");
Files.writeString(tmp, jobName + "\n" + lastCommittedLine + "\n" + lastCommittedChunk + "\n", StandardCharsets.UTF_8);
Files.move(tmp, file, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); // 반쯤 쓰인 체크포인트 방지
}
public static void clear(Path file) throws IOException { Files.deleteIfExists(file); }
}
// 파일 내용 예 (data/restartJob.checkpoint):
// restartJob
// 400
// 4체크포인트 파일 자체도 "임시 파일 + 원자적 이동"으로 씁니다. 체크포인트를 쓰다 죽어서 40 만 남으면 재시작이 40번째 줄부터 시작되는 사고가 나기 때문입니다.
public Result run(int crashAtChunk) throws IOException {
Optional<Checkpoint> cp = Checkpoint.load(checkpointFile);
long skipLines = cp.map(Checkpoint::lastCommittedLine).orElse(0L);
int chunkNo = cp.map(Checkpoint::lastCommittedChunk).orElse(0);
if (cp.isPresent())
System.out.printf(" [restart] 체크포인트 발견: 청크 %d / %d줄까지 완료 -> 이어서 시작%n", chunkNo, skipLines);
int committed = 0, rolledBack = 0;
long processed = 0, alreadyDone = 0, lineNo = 0;
try (BufferedReader r = Files.newBufferedReader(input, StandardCharsets.UTF_8)) {
r.readLine(); // 헤더
for (long i = 0; i < skipLines; i++) { r.readLine(); lineNo++; } // 재시작: 완료분 건너뛰기
while (true) {
List<Transfer> chunk = new ArrayList<>(chunkSize);
String line;
while (chunk.size() < chunkSize && (line = r.readLine()) != null) {
lineNo++;
chunk.add(Transfer.parse(line));
}
if (chunk.isEmpty()) break;
chunkNo++;
db.begin(); // ── 청크 트랜잭션 시작
long processedInChunk = 0, doneInChunk = 0;
try {
for (Transfer t : chunk) {
if (db.contains("done:" + t.txId())) { doneInChunk++; continue; } // 멱등: 이미 처리됨
apply(t); // 잔액 변경 + done 마킹 (같은 tx)
processedInChunk++;
}
if (chunkNo == crashAtChunk)
throw new RuntimeException("프로세스 장애 시뮬레이션 (청크 " + chunkNo + " 커밋 직전)");
db.commit(); // ── 커밋 포인트
committed++;
processed += processedInChunk;
alreadyDone += doneInChunk;
new Checkpoint(jobName, lineNo, chunkNo).save(checkpointFile); // 커밋 직후 재시작 지점 기록
System.out.printf(" 청크 %d 커밋 (줄 %d까지, 처리 %d건, 이미완료 %d건)%n", chunkNo, lineNo, processedInChunk, doneInChunk);
} catch (InsufficientBalanceException e) {
db.rollback(); // ── 이 청크의 변경 전부 취소
rolledBack++;
System.out.printf(" 청크 %d 롤백 (%s) -> 이 청크의 %d건 변경 전부 취소, 이전 청크는 유지%n",
chunkNo, e.getMessage(), processedInChunk);
new Checkpoint(jobName, lineNo, chunkNo).save(checkpointFile); // 다음 청크로 전진 (정책)
}
}
} catch (RuntimeException crash) {
if (db.inTransaction()) db.rollback(); // 실제 DB 는 커넥션 끊김 시 자동
System.out.println(" !!! " + crash.getMessage());
System.out.printf(" !!! 체크포인트는 청크 %d 에 머물러 있음 -> 다음 실행 시 청크 %d 부터 재개%n", chunkNo - 1, chunkNo);
throw crash;
}
Checkpoint.clear(checkpointFile); // 정상 종료: 다음 실행은 처음부터
return new Result(committed, rolledBack, processed, alreadyDone);
}
private void apply(Transfer t) throws InsufficientBalanceException {
String fromKey = "balance:" + t.from(), toKey = "balance:" + t.to();
long fromBal = db.getOrDefault(fromKey, 0);
if (fromBal < t.amount())
throw new InsufficientBalanceException(t.txId() + ": " + t.from() + " 잔액 " + fromBal + " < 이체 " + t.amount());
db.put(fromKey, fromBal - t.amount());
db.put(toKey, db.getOrDefault(toKey, 0) + t.amount());
db.put("done:" + t.txId(), 1L); // 완료 마킹도 같은 트랜잭션
}
// 출력 (rows=1000, chunk=100, 350번째 이체가 잔액 부족):
// 청크 1 커밋 (줄 100까지, 처리 100건, 이미완료 0건)
// 청크 2 커밋 (줄 200까지, 처리 100건, 이미완료 0건)
// 청크 3 커밋 (줄 300까지, 처리 100건, 이미완료 0건)
// 청크 4 롤백 (TX-000350: ACC-04 잔액 1009000 < 이체 999999999) -> 이 청크의 49건 변경 전부 취소, 이전 청크는 유지
// 청크 5 커밋 (줄 500까지, 처리 100건, 이미완료 0건)
// ...
// 청크 10 커밋 (줄 1000까지, 처리 100건, 이미완료 0건)
// 결과: 커밋 9청크, 롤백 1청크, 처리 900건
// 잔액 총합 20,000,000 -> 20,000,000 (이체는 총합을 바꾸지 않아야 함: OK)
// DB 통계: commit=9 rollback=1"잔액 총합이 변하지 않는다"는 것이 이체 배치의 불변식(invariant)입니다. 청크 4 가 반쯤 반영되었다면 총합이 깨졌을 것입니다. 롤백 덕분에 유지됩니다.
같은 TransferBatch 를 세 번 실행합니다. 첫 실행은 청크 5 커밋 직전에 죽고, 두 번째는 이어서 끝내고, 세 번째는 아무것도 바꾸지 않습니다.
InMemoryDatabase db3 = newBank(); // 20계좌 × 1,000,000
TransferBatch batch = new TransferBatch("restartJob", db3, transfers, dataDir, chunkSize);
try { batch.run(5); } // 청크 5 에서 장애
catch (RuntimeException e) {
System.out.println(" --- 프로세스 종료됨. 완료 마킹된 건수: " + db3.countOfPrefix("done:"));
}
// 출력:
// 청크 1 커밋 ... 청크 3 커밋, 청크 4 롤백
// !!! 프로세스 장애 시뮬레이션 (청크 5 커밋 직전)
// !!! 체크포인트는 청크 4 에 머물러 있음 -> 다음 실행 시 청크 5 부터 재개
// --- 프로세스 종료됨. 완료 마킹된 건수: 300 ← 청크 5 의 미커밋 변경은 사라짐
TransferBatch.Result r3 = batch.run(0); // 재실행
// 출력:
// [restart] 체크포인트 발견: 청크 4 / 400줄까지 완료 -> 이어서 시작
// 청크 5 커밋 (줄 500까지, 처리 100건, 이미완료 0건) ... 청크 10 커밋
// 결과: 커밋 6청크, 롤백 0청크, 처리 600건, 이미완료 0건
// 완료 마킹 900건, 잔액 총합 20,000,000 (OK)
Map<String, Long> before = db3.snapshot();
TransferBatch.Result r3b = batch.run(0); // 한 번 더 (체크포인트 없음 → 처음부터)
System.out.println(db3.snapshot().equals(before) ? "OK" : "FAIL");
// 출력:
// 청크 1 커밋 (줄 100까지, 처리 0건, 이미완료 100건) ← done 마킹 덕분에 건너뜀
// ... 청크 4 롤백 (여전히 잔액 부족) ...
// 처리 0건, 이미완료 900건 -> 잔액 변화 없음: OK세 번째 실행이 "OK"인 이유가 멱등성입니다. 체크포인트가 없어도(정상 종료 후 삭제됨) done: 마킹이 이중 처리를 막습니다. 체크포인트는 성능(완료분을 읽지 않음)을, 멱등성은 정확성(읽어도 두 번 반영 안 됨)을 담당합니다. 둘 다 필요합니다.
static void writeSettlement(InMemoryDatabase db, Path dest, boolean crash) throws IOException {
Path tmp = dest.resolveSibling(dest.getFileName() + ".tmp");
try (BufferedWriter w = Files.newBufferedWriter(tmp, StandardCharsets.UTF_8)) {
w.write("account,balance"); w.newLine();
int n = 0;
for (Map.Entry<String, Long> e : db.snapshot().entrySet()) {
if (!e.getKey().startsWith("balance:")) continue;
w.write(e.getKey().substring("balance:".length()) + "," + e.getValue()); w.newLine();
if (crash && ++n == 10) throw new RuntimeException("쓰기 도중 장애 시뮬레이션");
}
}
Files.move(tmp, dest, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
}
try { writeSettlement(db3, settlement, true); }
catch (RuntimeException e) {
System.out.println("정식 파일 존재? " + Files.exists(settlement) + ", 임시 파일 존재? " + Files.exists(tmp));
}
writeSettlement(db3, settlement, false);
// 출력:
// 1차 실행: 쓰기 도중 장애 시뮬레이션
// 정식 파일 존재? false, 임시 파일 존재? true
// 2차 실행: 정식 파일 존재? true, 줄 수 21, 임시 파일 존재? false재실행이 .tmp 를 덮어쓰고 완성 후 이동하므로 "이어 쓰기"가 필요 없습니다. 재생성이 싼 출력은 이 방식이 가장 단순하고 안전합니다.