공공부하자개발 · 영어 학습 노트
자바
실무 배치대용량 데이터 처리0/6 완료
  • 01배치 프로세스 개념과 아키텍처
  • 02대용량 파일 I/O (NIO, Buffered)
  • 03Chunk 단위 처리와 OOM 방지
  • 04트랜잭션: 커밋과 롤백 시뮬레이션
  • 05Skip과 Retry 로직 구현
  • 06메서드 활용 패턴 (배치 유틸)
사이트 소개개인정보처리방침연락처
© 2026 공부하자
홈 › 실무 배치 › 05 / 6

Skip과 Retry 로직 구현

섹션 7진행 0 / 6
1왜 배우는가2핵심 원리3코드 예제4응용 변형 예제5자주 하는 실수 (Tip)6연습 문제7정리‹ 이전다음 ›

4. 응용 변형 예제

3절에서는 RetryTemplate + SkipPolicy + BatchReport 를 하나의 통합 루프에 넣었습니다. 여기서는 같은 부품을 다른 각도에서 다시 씁니다.

백오프 정책의 실제 대기 시간을 재고, "재시도할 것인가"를 예외 클래스가 아니라 Predicate 로 판단하고, Writer 수준 실패를 Spring Batch 식 scan 모드로 격리하고, 백오프 대기 중 인터럽트로 배치를 정상 중단시키고, 스킵 한도가 마지막 건에서 걸릴 때 데드 레터에 무엇이 남는지 확인합니다.

변형 1: 백오프 정책별 실제 대기 시간 측정 — none / fixed / exponential / jitter

2.2 절의 표는 "계획된" 대기 시간입니다. 실제로 Thread.sleep 이 얼마나 걸리는지, 정책마다 총 대기가 어떻게 달라지는지를 같은 결정적 실패 시퀀스(3번 실패 후 4번째 성공)로 측정합니다. 지터는 Random(42) 시드를 고정해 재현 가능하게 했습니다. 재시도 총 대기 시간을 예측할 수 있어야 "배치가 왜 30분 더 걸렸는지"를 설명할 수 있습니다.

java
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 144ms

none 이 계획 0ms 인데도 수백 ms 가 나온 것은 첫 실행의 클래스 로딩·JIT 과 측정 시점의 CPU 경합 때문입니다 — 실측은 항상 이런 노이즈를 포함하므로 여러 번 중 최소값을 보는 습관이 필요합니다. 지터는 같은 지수 정책보다 대기가 짧게(50%~100%) 흩어지고, 시드를 고정하면 테스트에서도 재현됩니다.

변형 2: 재시도 판별을 클래스 집합 대신 Predicate<Exception> 으로

예제 2 의 RetryTemplate 은 예외 클래스로 재시도 여부를 정합니다. 그런데 실무에서는 같은 ApiException 이라도 HTTP 429/503 은 일시적, 400/404 는 영구적이고, DB 데드락은 벤더 메시지 문자열로만 구분되는 경우가 많습니다. 판별 로직을 Predicate<Exception> 으로 받으면 상태 코드·메시지·원인 체인 어떤 기준이든 or() 로 합성해 넘길 수 있습니다.

java
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 가 낫습니다.

변형 3: Writer 수준 실패와 scan 모드 — 청크 안의 범인 찾기

2.5 절 끝에서 언급한 Spring Batch 의 scan 모드입니다. DB 배치 INSERT 는 청크 안에 제약 위반 행이 하나만 있어도 통째로 거부하고 어느 행인지는 알려주지 않습니다. 그래서 청크를 롤백한 뒤 한 건씩 다시 써서 범인을 찾아 스킵하고 나머지는 건별 커밋합니다. 정상 청크는 writer 1회·커밋 1회로 빠르게, 문제 청크만 건별로 느리게 처리하는 것이 핵심입니다.

java
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 의 한도 검사를 붙여야 함은 물론입니다.

변형 4: 백오프 대기 중 인터럽트 — 운영자의 중단 요청을 정상 종료로

재시도의 Thread.sleep 은 배치가 가장 오래 멈춰 있는 지점이고, 운영자가 "지금 멈춰"라고 하면 바로 그 대기 중일 확률이 높습니다. InterruptedException 을 삼키지 않고 전파하면 진행 중 청크를 롤백하고 마지막 커밋 위치를 남긴 채 깨끗하게 끝낼 수 있습니다. 이 변형은 워커 스레드를 띄우고 백오프 대기 중에 interrupt() 를 보내 그 흐름을 확인합니다.

java
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() 가 모두 이 경로로 들어옵니다.

변형 5: 스킵 한도가 마지막 건에서 걸릴 때 — 데드 레터에 무엇이 남는가

skipLimit 초과는 "N+1 번째 스킵 시도" 순간에 발생합니다. 그 건이 입력의 마지막 행이면 배치는 99% 를 끝내고 FAILED 가 되고, 한도를 넘긴 그 행은 데드 레터에 기록되지 않습니다(shouldSkip 이 기록 전에 던지므로). 같은 입력(잘못된 행 3, 7, 10)을 한도 3 과 2 로 돌려 리포트와 데드 레터 파일 내용을 대조합니다.

java
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 레슨의 체크포인트·멱등성 위에서 이루어져야 합니다.

응용 변형 예제
  • 변형 1: 백오프 정책별 실제 대기 시간 측정 — none / fixed / exponential / jitter
  • 변형 2: 재시도 판별을 클래스 집합 대신 Predicate<Exception> 으로
  • 변형 3: Writer 수준 실패와 scan 모드 — 청크 안의 범인 찾기
  • 변형 4: 백오프 대기 중 인터럽트 — 운영자의 중단 요청을 정상 종료로
  • 변형 5: 스킵 한도가 마지막 건에서 걸릴 때 — 데드 레터에 무엇이 남는가
이전 섹션3 코드 예제4 / 7다음 섹션5 자주 하는 실수 (Tip)