3절에서는 RetryTemplate + SkipPolicy + BatchReport 를 하나의 통합 루프에 넣었습니다. 여기서는 같은 부품을 다른 각도에서 다시 씁니다.
백오프 정책의 실제 대기 시간을 재고, "재시도할 것인가"를 예외 클래스가 아니라 Predicate 로 판단하고, Writer 수준 실패를 Spring Batch 식 scan 모드로 격리하고, 백오프 대기 중 인터럽트로 배치를 정상 중단시키고, 스킵 한도가 마지막 건에서 걸릴 때 데드 레터에 무엇이 남는지 확인합니다.
2.2 절의 표는 "계획된" 대기 시간입니다. 실제로 Thread.sleep 이 얼마나 걸리는지, 정책마다 총 대기가 어떻게 달라지는지를 같은 결정적 실패 시퀀스(3번 실패 후 4번째 성공)로 측정합니다. 지터는 Random(42) 시드를 고정해 재현 가능하게 했습니다. 재시도 총 대기 시간을 예측할 수 있어야 "배치가 왜 30분 더 걸렸는지"를 설명할 수 있습니다.
import java.util.*;
import java.util.concurrent.Callable;
public class BackoffCompare {
@FunctionalInterface interface BackoffPolicy { long delayMillis(int attempt); }
static BackoffPolicy none() { return a -> 0; }
static BackoffPolicy fixed(long ms) { return a -> ms; }
static BackoffPolicy exponential(long init, double mul, long max) {
return a -> (long) Math.min(init * Math.pow(mul, a - 1), max);
}
static BackoffPolicy exponentialWithJitter(long init, double mul, long max, Random rnd) { // 시드 고정 → 재현 가능
BackoffPolicy base = exponential(init, mul, max);
return a -> { long d = base.delayMillis(a); return d == 0 ? 0 : d / 2 + rnd.nextLong(d / 2 + 1); };
}
static class Transient extends RuntimeException { Transient(int n) { super("timeout #" + n); } }
/** 처음 failures 번은 실패, 그 다음 성공하는 결정적 액션 */
static Callable<String> failingThenOk(int failures) {
int[] calls = {0};
return () -> { if (++calls[0] <= failures) throw new Transient(calls[0]); return "ok@" + calls[0]; };
}
record Run(String result, List<Long> delays, long elapsedMs) {}
static Run retry(Callable<String> action, int maxAttempts, BackoffPolicy backoff) throws Exception {
List<Long> delays = new ArrayList<>();
long t0 = System.nanoTime();
for (int attempt = 1; ; attempt++) {
try { return new Run(action.call(), delays, (System.nanoTime() - t0) / 1_000_000); }
catch (Transient e) {
if (attempt >= maxAttempts) throw e;
long d = backoff.delayMillis(attempt);
delays.add(d);
Thread.sleep(d);
}
}
}
public static void main(String[] args) throws Exception {
Random rnd = new Random(42);
Map<String, BackoffPolicy> policies = new LinkedHashMap<>();
policies.put("none", none());
policies.put("fixed(30)", fixed(30));
policies.put("exponential(10,2,1000)", exponential(10, 2, 1000));
policies.put("jitter(10,2,1000)", exponentialWithJitter(10, 2, 1000, rnd));
System.out.printf("%-24s %-8s %-22s %s%n", "policy", "result", "delays(ms)", "measured");
for (var e : policies.entrySet()) {
Run r = retry(failingThenOk(3), 5, e.getValue()); // 3번 실패 후 4번째 성공
long best = r.elapsedMs();
for (int rep = 0; rep < 2; rep++) best = Math.min(best, retry(failingThenOk(3), 5, e.getValue()).elapsedMs()); // 최소값
long planned = r.delays().stream().mapToLong(Long::longValue).sum();
System.out.printf("%-24s %-8s %-22s plan %3dms / real %3dms%n", e.getKey(), r.result(), r.delays(), planned, best);
}
}
}
// 출력 (real 은 환경·부하에 따라 다름. delays 는 시드 고정으로 항상 같음):
// policy result delays(ms) measured
// none ok@4 [0, 0, 0] plan 0ms / real 467ms
// fixed(30) ok@4 [30, 30, 30] plan 90ms / real 106ms
// exponential(10,2,1000) ok@4 [10, 20, 40] plan 70ms / real 371ms
// jitter(10,2,1000) ok@4 [6, 16, 28] plan 50ms / real 144msnone 이 계획 0ms 인데도 수백 ms 가 나온 것은 첫 실행의 클래스 로딩·JIT 과 측정 시점의 CPU 경합 때문입니다 — 실측은 항상 이런 노이즈를 포함하므로 여러 번 중 최소값을 보는 습관이 필요합니다. 지터는 같은 지수 정책보다 대기가 짧게(50%~100%) 흩어지고, 시드를 고정하면 테스트에서도 재현됩니다.
Predicate<Exception> 으로예제 2 의 RetryTemplate 은 예외 클래스로 재시도 여부를 정합니다. 그런데 실무에서는 같은 ApiException 이라도 HTTP 429/503 은 일시적, 400/404 는 영구적이고, DB 데드락은 벤더 메시지 문자열로만 구분되는 경우가 많습니다. 판별 로직을 Predicate<Exception> 으로 받으면 상태 코드·메시지·원인 체인 어떤 기준이든 or() 로 합성해 넘길 수 있습니다.
import java.util.*;
import java.util.concurrent.Callable;
import java.util.function.Predicate;
public class PredicateRetry {
static class ApiException extends RuntimeException {
final int status;
ApiException(int status, String msg) { super(status + " " + msg); this.status = status; }
}
static class RetryExhausted extends Exception { RetryExhausted(String m, Throwable c) { super(m, c); } }
/** 클래스 집합 대신 Predicate 로 "재시도할 것인가" 를 판단하는 제네릭 재시도 함수 */
static <T> T retry(Callable<T> action, int maxAttempts, Predicate<Exception> retryOn, List<String> log) throws Exception {
for (int attempt = 1; ; attempt++) {
try { return action.call(); }
catch (Exception e) {
if (!retryOn.test(e)) throw e; // 재시도 대상 아님 → 즉시 전파
if (attempt >= maxAttempts) throw new RetryExhausted("재시도 " + maxAttempts + "회 소진", e);
log.add("retry#" + attempt + "(" + e.getMessage() + ")");
}
}
}
// HTTP 상태 코드 기준: 429/503/504 만 일시적. 같은 예외 클래스라도 상태로 갈린다
static final Predicate<Exception> TRANSIENT_STATUS =
e -> e instanceof ApiException a && Set.of(429, 503, 504).contains(a.status);
// DB 데드락은 메시지로만 구분되는 경우가 많다 (SQLException 의 vendor 메시지)
static final Predicate<Exception> DEADLOCK_MESSAGE =
e -> e.getMessage() != null && e.getMessage().toLowerCase().contains("deadlock");
/** 스크립트대로 응답하는 가짜 API: id 별 상태 코드 순서 */
static Callable<String> scripted(String id, int... statuses) {
int[] i = {0};
return () -> {
int s = statuses[Math.min(i[0]++, statuses.length - 1)];
if (s == 200) return id + " sent";
throw new ApiException(s, switch (s) { case 429 -> "Too Many Requests"; case 400 -> "Bad Request";
case 503 -> "Service Unavailable"; default -> "Error"; });
};
}
public static void main(String[] args) {
Map<String, Callable<String>> calls = new LinkedHashMap<>();
calls.put("ORD-1", scripted("ORD-1", 503, 200)); // 일시 오류 1회 후 성공
calls.put("ORD-2", scripted("ORD-2", 400)); // 영구 오류 → 재시도 없이 즉시
calls.put("ORD-3", scripted("ORD-3", 429, 429, 429, 429)); // 계속 429 → 소진
calls.put("ORD-4", scripted("ORD-4", 200));
calls.put("ORD-5", () -> { throw new IllegalStateException("Deadlock found when trying to get lock"); });
Predicate<Exception> retryOn = TRANSIENT_STATUS.or(DEADLOCK_MESSAGE);
for (var e : calls.entrySet()) {
List<String> log = new ArrayList<>();
String outcome;
try { outcome = retry(e.getValue(), 3, retryOn, log); }
catch (RetryExhausted ex) { outcome = "EXHAUSTED (cause: " + ex.getCause().getMessage() + ")"; }
catch (Exception ex) { outcome = "SKIP (" + ex.getMessage() + ")"; }
System.out.printf("%-6s %-45s %s%n", e.getKey(), outcome, log);
}
}
}
// 출력:
// ORD-1 ORD-1 sent [retry#1(503 Service Unavailable)]
// ORD-2 SKIP (400 Bad Request) []
// ORD-3 EXHAUSTED (cause: 429 Too Many Requests) [retry#1(429 Too Many Requests), retry#2(429 Too Many Requests)]
// ORD-4 ORD-4 sent []
// ORD-5 EXHAUSTED (cause: Deadlock found when trying to get lock) [retry#1(Deadlock found when trying to get lock), retry#2(Deadlock found when trying to get lock)]ORD-2 는 같은 ApiException 인데도 400 이라 재시도 로그가 비어 있습니다 — 클래스 기준이었다면 3번 헛되이 재시도했을 것입니다. ORD-5 는 IllegalStateException 이지만 메시지의 "deadlock" 으로 재시도 대상이 되었습니다.
Predicate 는 클래스 화이트리스트를 포함하는 상위 개념이므로(c::isInstance 를 넘기면 동일), 기준이 하나뿐이라면 클래스 집합이, 둘 이상이면 Predicate 가 낫습니다.
2.5 절 끝에서 언급한 Spring Batch 의 scan 모드입니다. DB 배치 INSERT 는 청크 안에 제약 위반 행이 하나만 있어도 통째로 거부하고 어느 행인지는 알려주지 않습니다. 그래서 청크를 롤백한 뒤 한 건씩 다시 써서 범인을 찾아 스킵하고 나머지는 건별 커밋합니다. 정상 청크는 writer 1회·커밋 1회로 빠르게, 문제 청크만 건별로 느리게 처리하는 것이 핵심입니다.
import java.util.*;
public class WriterScanMode {
record Order(long id, long amount) {}
static class Db {
final Map<Long, Long> committed = new LinkedHashMap<>();
Map<Long, Long> pending;
int commits, rollbacks, writerCalls;
void begin() { pending = new LinkedHashMap<>(); }
void commit() { committed.putAll(pending); pending = null; commits++; }
void rollback() { pending = null; rollbacks++; }
/** 배치 INSERT: 청크 안에 잘못된 행이 하나라도 있으면 DB 가 통째로 거부. 어느 행인지는 알려주지 않는다 */
void batchInsert(List<Order> orders) {
writerCalls++;
if (orders.stream().anyMatch(o -> o.amount() <= 0))
throw new IllegalArgumentException("CHECK 제약 위반 (amount > 0), 배치 크기 " + orders.size());
for (Order o : orders) pending.put(o.id(), o.amount());
}
}
static void processChunk(Db db, List<Order> chunk, List<String> deadLetter) {
db.begin();
try {
db.batchInsert(chunk); // 1) 청크 통째로 시도
db.commit();
System.out.println(" 청크 커밋 (" + chunk.size() + "건, writer 1회)");
return;
} catch (IllegalArgumentException e) {
db.rollback(); // 2) 실패 → 청크 롤백
System.out.println(" 청크 쓰기 실패: " + e.getMessage() + " -> scan 모드 진입");
}
int ok = 0;
for (Order o : chunk) { // 3) 한 건씩 다시 써서 범인을 찾는다
db.begin();
try { db.batchInsert(List.of(o)); db.commit(); ok++; }
catch (IllegalArgumentException e) {
db.rollback();
deadLetter.add(o.id() + "," + e.getMessage().split(",")[0] + "," + o);
System.out.println(" id=" + o.id() + " 스킵 -> 데드 레터");
}
}
System.out.println(" scan 완료: " + ok + "/" + chunk.size() + "건 커밋 (건별 커밋 " + ok + "회)");
}
public static void main(String[] args) {
Db db = new Db();
List<String> deadLetter = new ArrayList<>();
List<Order> chunk1 = new ArrayList<>(), chunk2 = new ArrayList<>();
for (long i = 1; i <= 10; i++) chunk1.add(new Order(i, i * 100));
for (long i = 11; i <= 20; i++) chunk2.add(new Order(i, i == 17 ? 0 : i * 100)); // 17번이 독약
processChunk(db, chunk1, deadLetter);
processChunk(db, chunk2, deadLetter);
System.out.println("committed=" + db.committed.size() + " commits=" + db.commits + " rollbacks=" + db.rollbacks
+ " writerCalls=" + db.writerCalls);
System.out.println("deadLetter=" + deadLetter);
}
}
// 출력:
// 청크 커밋 (10건, writer 1회)
// 청크 쓰기 실패: CHECK 제약 위반 (amount > 0), 배치 크기 10 -> scan 모드 진입
// id=17 스킵 -> 데드 레터
// scan 완료: 9/10건 커밋 (건별 커밋 9회)
// committed=19 commits=10 rollbacks=2 writerCalls=12
// deadLetter=[17,CHECK 제약 위반 (amount > 0),Order[id=17, amount=0]]두 번째 청크는 writer 호출이 1(실패) + 10(scan) = 11회, 커밋 9회, 롤백 2회(청크 1회 + 범인 1회)입니다. scan 모드는 비싸므로 Writer 실패가 잦다면 Processor 단계에서 미리 검증해 스킵하는 것이 낫습니다 — scan 은 "Processor 가 잡지 못한 실패"를 위한 안전망입니다. 이 결과에도 SkipPolicy 의 한도 검사를 붙여야 함은 물론입니다.
재시도의 Thread.sleep 은 배치가 가장 오래 멈춰 있는 지점이고, 운영자가 "지금 멈춰"라고 하면 바로 그 대기 중일 확률이 높습니다. InterruptedException 을 삼키지 않고 전파하면 진행 중 청크를 롤백하고 마지막 커밋 위치를 남긴 채 깨끗하게 끝낼 수 있습니다. 이 변형은 워커 스레드를 띄우고 백오프 대기 중에 interrupt() 를 보내 그 흐름을 확인합니다.
import java.util.*;
import java.util.concurrent.Callable;
import java.util.concurrent.atomic.AtomicLong;
public class InterruptDuringBackoff {
static class Transient extends RuntimeException { Transient(String m) { super(m); } }
/** id 5 부터는 항상 타임아웃 → 재시도가 백오프 대기에 들어간다 */
static String send(long id) {
if (id >= 5) throw new Transient("timeout id=" + id);
return "sent " + id;
}
static <T> T retry(Callable<T> action, int maxAttempts, long backoffMs) throws Exception {
for (int attempt = 1; ; attempt++) {
try { return action.call(); }
catch (Transient e) {
if (attempt >= maxAttempts) throw e;
Thread.sleep(backoffMs); // 인터럽트되면 InterruptedException → 그대로 전파 (삼키지 않는다)
}
}
}
public static void main(String[] args) throws Exception {
int chunkSize = 3;
AtomicLong checkpoint = new AtomicLong();
String[] exit = {"RUNNING"};
Thread worker = new Thread(() -> {
List<String> chunk = new ArrayList<>();
long id = 1;
try {
while (id <= 100) {
final long cur = id;
chunk.add(retry(() -> send(cur), 5, 200));
if (chunk.size() == chunkSize) { checkpoint.set(id); chunk.clear(); } // 청크 커밋
id++;
}
exit[0] = "COMPLETED";
} catch (InterruptedException e) {
chunk.clear(); // 진행 중 청크 롤백
Thread.currentThread().interrupt(); // 플래그 복원 (관례)
exit[0] = "STOPPED at id=" + id + ", rolled back " + (id - 1 - checkpoint.get()) + " uncommitted item(s)";
} catch (Exception e) {
exit[0] = "FAILED " + e;
}
}, "batch-worker");
worker.start();
Thread.sleep(300); // id 1~3 커밋, id 4 처리 후 id 5 의 백오프 대기 중
System.out.println("운영자: 중단 요청 (interrupt)");
worker.interrupt();
worker.join(2_000);
System.out.println("worker alive? " + worker.isAlive());
System.out.println("exit = " + exit[0]);
System.out.println("checkpoint = " + checkpoint.get() + " -> 재시작 시 id " + (checkpoint.get() + 1) + " 부터");
}
}
// 출력:
// 운영자: 중단 요청 (interrupt)
// worker alive? false
// exit = STOPPED at id=5, rolled back 1 uncommitted item(s)
// checkpoint = 3 -> 재시작 시 id 4 부터id 5 는 최대 5회 × 200ms = 1초를 기다릴 참이었지만 300ms 시점의 인터럽트로 즉시 깨어났고, 청크에 들어 있던 id 4 는 롤백되어 재시작이 4 부터 이어집니다.
예제 2 의 RetryTemplate.execute 가 throws Exception 인 이유가 이것입니다 — catch (InterruptedException e) {} 로 삼키면 배치는 인터럽트를 무시하고 남은 800ms 를 계속 기다립니다(02 레슨 멀티스레드의 관례).
shutdownNow(), Ctrl+C 훅, ExecutorService.close() 가 모두 이 경로로 들어옵니다.
skipLimit 초과는 "N+1 번째 스킵 시도" 순간에 발생합니다. 그 건이 입력의 마지막 행이면 배치는 99% 를 끝내고 FAILED 가 되고, 한도를 넘긴 그 행은 데드 레터에 기록되지 않습니다(shouldSkip 이 기록 전에 던지므로). 같은 입력(잘못된 행 3, 7, 10)을 한도 3 과 2 로 돌려 리포트와 데드 레터 파일 내용을 대조합니다.
import java.io.*;
import java.nio.charset.StandardCharsets;
import java.nio.file.*;
import java.util.*;
public class SkipLimitEdge {
static class SkipLimitExceeded extends RuntimeException {
SkipLimitExceeded(String m, Throwable c) { super(m, c); }
}
static class SkipPolicy {
final int limit; int count;
SkipPolicy(int limit) { this.limit = limit; }
boolean shouldSkip(Throwable t) {
if (!(t instanceof IllegalArgumentException)) return false;
if (++count > limit) throw new SkipLimitExceeded("스킵 한도 " + limit + " 초과 (" + count + "번째 스킵)", t);
return true;
}
}
record Result(String status, long total, long ok, long skipped, List<String> deadLetter) {}
static Result run(Path input, Path deadLetterFile, int skipLimit) throws IOException {
SkipPolicy policy = new SkipPolicy(skipLimit);
long total = 0, ok = 0;
String status = "COMPLETED";
try (BufferedReader r = Files.newBufferedReader(input, StandardCharsets.UTF_8);
BufferedWriter dl = Files.newBufferedWriter(deadLetterFile, StandardCharsets.UTF_8)) {
dl.write("line,reason,raw"); dl.newLine();
String line; long lineNo = 0;
while ((line = r.readLine()) != null) {
lineNo++; total++;
try { Long.parseLong(line.split(",")[1]); ok++; }
catch (RuntimeException e) {
if (!policy.shouldSkip(e)) throw e; // 한도 초과 → SkipLimitExceeded 전파
dl.write(lineNo + "," + e.getMessage().replace(',', ';') + "," + line); dl.newLine();
}
}
} catch (SkipLimitExceeded e) {
status = "FAILED: " + e.getMessage();
}
List<String> dead = Files.readAllLines(deadLetterFile, StandardCharsets.UTF_8);
return new Result(status, total, ok, policy.count > skipLimit ? skipLimit : policy.count, dead.subList(1, dead.size()));
}
public static void main(String[] args) throws IOException {
Path dir = Files.createTempDirectory("skipedge");
Path input = dir.resolve("orders.csv");
List<String> lines = new ArrayList<>();
for (int i = 1; i <= 10; i++) lines.add("ORD-" + i + "," + (i == 3 || i == 7 || i == 10 ? i + "abc" : i * 100));
Files.write(input, lines, StandardCharsets.UTF_8); // 잘못된 행: 3, 7, 10(마지막)
for (int limit : new int[]{3, 2}) {
Result r = run(input, dir.resolve("dead-" + limit + ".csv"), limit);
System.out.println("skipLimit=" + limit + " -> " + r.status());
System.out.println(" 총 " + r.total() + " = 성공 " + r.ok() + " + 스킵 " + r.skipped()
+ (r.total() == r.ok() + r.skipped() ? " (검산 OK)" : " (불일치: 마지막 건은 기록 전에 중단됨)"));
r.deadLetter().forEach(l -> System.out.println(" deadletter: " + l));
}
}
}
// 출력:
// skipLimit=3 -> COMPLETED
// 총 10 = 성공 7 + 스킵 3 (검산 OK)
// deadletter: 3,For input string: "3abc",ORD-3,3abc
// deadletter: 7,For input string: "7abc",ORD-7,7abc
// deadletter: 10,For input string: "10abc",ORD-10,10abc
// skipLimit=2 -> FAILED: 스킵 한도 2 초과 (3번째 스킵)
// 총 10 = 성공 7 + 스킵 2 (불일치: 마지막 건은 기록 전에 중단됨)
// deadletter: 3,For input string: "3abc",ORD-3,3abc
// deadletter: 7,For input string: "7abc",ORD-7,7abc한도 2 에서는 총 = 성공 + 스킵 검산이 깨집니다. 한도를 넘긴 10번 행이 데드 레터에 없기 때문인데, 그 행의 정보는 SkipLimitExceeded 의 cause 와 메시지에만 있습니다.
실무에서는 FAILED 로그에 이 원인을 반드시 남겨야 하고(실수 6), 리포트의 검산이 안 맞는 것 자체가 "한도로 중단됐다"는 신호로 읽혀야 합니다. 그리고 한도 초과가 마지막 행에서 났다면 성공한 7건은 이미 커밋된 상태이므로, 재실행은 04 레슨의 체크포인트·멱등성 위에서 이루어져야 합니다.