java-src/batch/05_skip_retry/ 에서 javac *.java && java Main 100000.
@FunctionalInterface
public interface BackoffPolicy {
long delayMillis(int attempt); // attempt 는 1부터
static BackoffPolicy none() { return attempt -> 0; }
static BackoffPolicy fixed(long millis) { return attempt -> millis; }
static BackoffPolicy exponential(long initialMillis, double multiplier, long maxMillis) {
return attempt -> (long) Math.min(initialMillis * Math.pow(multiplier, attempt - 1), maxMillis);
}
static BackoffPolicy exponentialWithJitter(long initialMillis, double multiplier, long maxMillis) {
BackoffPolicy base = exponential(initialMillis, multiplier, maxMillis);
return attempt -> {
long d = base.delayMillis(attempt);
return d == 0 ? 0 : ThreadLocalRandom.current().nextLong(d / 2, d + 1); // 50%~100%
};
}
}
// BackoffPolicy p = BackoffPolicy.exponential(100, 2, 5000);
// for (int a = 1; a <= 7; a++) System.out.print(p.delayMillis(a) + " ");
// 출력: 100 200 400 800 1600 3200 5000 ← 상한 5000 에서 멈춤
// exponentialWithJitter(100, 2, 5000): 73 151 322 512 1043 2890 4012 (매번 다름)public class RetryTemplate {
public static class RetryExhaustedException extends Exception {
public RetryExhaustedException(String msg, Throwable cause) { super(msg, cause); }
}
@FunctionalInterface
public interface RetryListener { void onRetry(int attempt, long delayMillis, Exception cause); }
private final int maxAttempts;
private final BackoffPolicy backoff;
private final Set<Class<? extends Exception>> retryable; // 화이트리스트
private final RetryListener listener;
private final AtomicLong totalRetries = new AtomicLong();
private final AtomicLong exhausted = new AtomicLong();
public RetryTemplate(int maxAttempts, BackoffPolicy backoff,
Set<Class<? extends Exception>> retryable, RetryListener listener) {
this.maxAttempts = maxAttempts; this.backoff = backoff; this.retryable = retryable;
this.listener = listener == null ? (a, d, e) -> {} : listener;
}
public boolean isRetryable(Exception e) {
for (Class<? extends Exception> c : retryable) if (c.isInstance(e)) return true;
return false;
}
public <T> T execute(Callable<T> action) throws Exception {
int attempt = 1;
while (true) {
try {
return action.call();
} catch (Exception e) {
if (!isRetryable(e)) throw e; // 영구 오류: 즉시 던짐
if (attempt >= maxAttempts) {
exhausted.incrementAndGet();
throw new RetryExhaustedException("재시도 " + maxAttempts + "회 모두 실패: " + e.getMessage(), e);
}
long delay = backoff.delayMillis(attempt);
listener.onRetry(attempt, delay, e);
totalRetries.incrementAndGet();
if (delay > 0) Thread.sleep(delay);
attempt++;
}
}
}
public long totalRetries() { return totalRetries.get(); }
public long exhaustedCount() { return exhausted.get(); }
}
// 사용: 30% 확률로 타임아웃 나는 API 를 20번 호출
FlakyApiClient api = new FlakyApiClient(1, 0.30, 0);
RetryTemplate retry = new RetryTemplate(4, BackoffPolicy.exponentialWithJitter(10, 2.0, 200),
Set.of(FlakyApiClient.TransientApiException.class),
(attempt, delay, cause) -> System.out.printf(" retry #%d after %dms: %s%n", attempt, delay, cause.getMessage()));
for (long id = 1; id <= 20; id++) {
final long fid = id;
try { retry.execute(() -> api.send(fid)); ok++; }
catch (RetryTemplate.RetryExhaustedException e) { failed++; }
}
// 출력:
// retry #1 after 7ms: 504 Gateway Timeout: id=2
// retry #1 after 9ms: 504 Gateway Timeout: id=5
// retry #2 after 14ms: 504 Gateway Timeout: id=5
// ...
// 20건 중 성공 20, 포기 0 / API 호출 28회, 재시도 8회public class SkipPolicy {
public static class SkipLimitExceededException extends RuntimeException {
public SkipLimitExceededException(String msg, Throwable cause) { super(msg, cause); }
}
private final int skipLimit;
private final Set<Class<? extends Throwable>> skippable;
private final AtomicLong skipCount = new AtomicLong();
public SkipPolicy(int skipLimit, Set<Class<? extends Throwable>> skippable) {
this.skipLimit = skipLimit; this.skippable = skippable;
}
public boolean isSkippable(Throwable t) {
// 재시도 소진은 원인(cause)으로 판별: Transient 를 skippable 에 넣으면 "재시도 다 해도 안 되면 스킵"
Throwable target = t instanceof RetryTemplate.RetryExhaustedException && t.getCause() != null ? t.getCause() : t;
for (Class<? extends Throwable> c : skippable) if (c.isInstance(target) || c.isInstance(t)) return true;
return false;
}
/** 스킵 가능하면 카운트 후 true, 한도 초과면 예외, 스킵 불가면 false */
public boolean shouldSkip(Throwable t) {
if (!isSkippable(t)) return false;
long n = skipCount.incrementAndGet();
if (n > skipLimit) throw new SkipLimitExceededException("스킵 한도 " + skipLimit + " 초과 (" + n + "번째 스킵)", t);
return true;
}
public long skipCount() { return skipCount.get(); }
}
// 사용: 잘못된 CSV 행 스킵 → 데드 레터
SkipPolicy skipPolicy = new SkipPolicy(rows / 50, Set.of(IllegalArgumentException.class)); // NumberFormatException 포함
try (BufferedReader r = Files.newBufferedReader(orders, UTF_8);
BatchReport report = new BatchReport("csvParseJob", Path.of("data/orders.error.csv"))) {
r.readLine();
String line; long lineNo = 1;
while ((line = r.readLine()) != null) {
lineNo++; report.read();
try { Order.parse(line); parsedOk++; }
catch (RuntimeException e) {
if (!skipPolicy.shouldSkip(e)) throw e; // 스킵 불가 → 중단
report.skip("line " + lineNo, e, line); // 데드 레터
}
}
report.success(parsedOk);
report.print(0);
}
// 출력:
// ┌─ csvParseJob 리포트 ─────────────────
// │ 총 건수 : 100,000
// │ 성공 : 99,012
// │ 스킵 : 988 (-> data\orders.error.csv)
// │ 재시도 : 0회
// │ 소요 시간 : 142 ms (704,225 items/s)
// └────────────────────────────────────
// 데드 레터 앞 3줄:
// lineOrId,reason,rawData
// line 137,NumberFormatException: For input string: "136abc",136;cust-412;136abc
// line 251,IllegalArgumentException: 필드 수 2 (3 필요),250;cust-87
// skipLimit=5 로 다시 실행하면:
// 배치 FAILED: 스킵 한도 5 초과 (6번째 스킵) / 원인: For input string: "612abc"public class BatchReport implements AutoCloseable {
private final String jobName;
private final long startMillis = System.currentTimeMillis();
private final AtomicLong total = new AtomicLong(), success = new AtomicLong(), skipped = new AtomicLong();
private final BufferedWriter deadLetter;
private final Path deadLetterFile;
public BatchReport(String jobName, Path deadLetterFile) throws IOException {
this.jobName = jobName; this.deadLetterFile = deadLetterFile;
Files.createDirectories(deadLetterFile.getParent());
this.deadLetter = Files.newBufferedWriter(deadLetterFile, StandardCharsets.UTF_8);
deadLetter.write("lineOrId,reason,rawData"); deadLetter.newLine();
}
public void read() { total.incrementAndGet(); }
public void read(long n) { total.addAndGet(n); }
public void success(long n) { success.addAndGet(n); }
public synchronized void skip(String idOrLine, Throwable reason, String raw) { // 사유 + 원본
skipped.incrementAndGet();
try {
deadLetter.write(idOrLine + "," + reason.getClass().getSimpleName() + ": "
+ String.valueOf(reason.getMessage()).replace(',', ';') + "," + raw.replace(',', ';'));
deadLetter.newLine();
} catch (IOException e) { throw new IllegalStateException("데드 레터 기록 실패", e); }
}
public void print(long retries) {
long elapsed = System.currentTimeMillis() - startMillis;
System.out.println(" ┌─ " + jobName + " 리포트 ─────────────────");
System.out.printf (" │ 총 건수 : %,d%n", total.get());
System.out.printf (" │ 성공 : %,d%n", success.get());
System.out.printf (" │ 스킵 : %,d (-> %s)%n", skipped.get(), deadLetterFile);
System.out.printf (" │ 재시도 : %,d회%n", retries);
System.out.printf (" │ 소요 시간 : %,d ms (%,d items/s)%n", elapsed, elapsed == 0 ? total.get() : total.get() * 1000 / elapsed);
System.out.println(" └────────────────────────────────────");
}
@Override public void close() throws IOException { deadLetter.close(); }
}데드 레터 기록이 실패하면(디스크 풀) 예외를 던져 배치를 중단시킵니다. "스킵했는데 기록은 못 했다"는 유실이기 때문입니다.
FlakyApiClient api3 = new FlakyApiClient(2, 0.10, 997); // 10% 일시 오류, id 997 배수는 영구 오류
RetryTemplate retry3 = new RetryTemplate(3, BackoffPolicy.exponential(1, 2.0, 4),
Set.of(FlakyApiClient.TransientApiException.class), null);
SkipPolicy skip3 = new SkipPolicy(1_000, Set.of(
IllegalArgumentException.class, // 파싱 오류
FlakyApiClient.PermanentApiException.class, // 영구 API 오류
FlakyApiClient.TransientApiException.class)); // 재시도 소진(cause) 도 스킵
try (BufferedReader r = Files.newBufferedReader(orders, UTF_8);
BatchReport report = new BatchReport("sendOrdersJob", Path.of("data/orders.deadletter.csv"))) {
r.readLine();
List<String> chunk = new ArrayList<>(200);
long readLines = 0; int chunkNo = 0;
while (true) {
chunk.clear();
String line;
while (chunk.size() < 200 && readLines < apiRows && (line = r.readLine()) != null) {
chunk.add(line); readLines++;
}
if (chunk.isEmpty()) break;
chunkNo++;
report.read(chunk.size());
List<String> sent = new ArrayList<>(chunk.size()); // 이 청크에서 성공한 것만
for (String raw : chunk) {
try {
Order o = Order.parse(raw); // 영구 오류 가능 → 스킵
String res = retry3.execute(() -> api3.send(o.id())); // 일시 오류 → 재시도, 소진 → 스킵
sent.add(res);
} catch (Exception e) {
if (!skip3.shouldSkip(e)) throw e; // 치명적 → 배치 중단
report.skip(raw.split(",")[0], e, raw);
}
}
report.success(sent.size()); // 커밋 포인트: 실패 건 제외하고 커밋
if (chunkNo % 5 == 0)
System.out.printf(" chunk#%d committed: %d/%d sent, skipped so far %d, retries so far %d%n",
chunkNo, sent.size(), chunk.size(), skip3.skipCount(), retry3.totalRetries());
}
report.print(retry3.totalRetries());
System.out.printf(" API 호출 %,d회, 재시도 소진 %d건%n", api3.calls(), retry3.exhaustedCount());
}
// 출력 (앞 5,000건):
// chunk#5 committed: 197/200 sent, skipped so far 12, retries so far 108
// chunk#10 committed: 199/200 sent, skipped so far 21, retries so far 221
// ...
// ┌─ sendOrdersJob 리포트 ─────────────────
// │ 총 건수 : 5,000
// │ 성공 : 4,946
// │ 스킵 : 54 (-> data\orders.deadletter.csv)
// │ 재시도 : 551회
// │ 소요 시간 : 1,830 ms (2,732 items/s)
// └────────────────────────────────────
// API 호출 5,502회 (성공건 + 재시도), 재시도 소진 5건
// 검산: 성공 4,946 + 스킵 54 = 5,000 ✓스킵 54건의 내역은 데드 레터에 있습니다. 파싱 오류 ~25건(1%의 절반씩 두 종류), 영구 API 오류 5건(997, 1994, 2991, 3988, 4985), 재시도 소진 5건(0.1³ ≈ 0.1% × 5,000)... 이런 식으로 숫자가 설명되어야 정상입니다.