Step 11 — 내결함성: skip · retry · 재시작
학습 목표
.faultTolerant() 를 켜고 .skip() · .skipLimit() · .noSkip() · 커스텀 SkipPolicy 로 건너뛸 예외를 정확히 통제한다
- skip 이 발생하면 청크 전체가 롤백되고 아이템을 하나씩 다시 처리(scan)한다는 사실을 로그로 정면 재현한다
- 정상 6.1초 → 쓰기 skip 100건 34.7초, 약 5.7배 느려지는 것을 실측하고 원인을
commitCount 로 증명한다
.retry() · .retryLimit() · backOffPolicy 를 걸고, retry 와 skip 이 어떤 순서로 맞물리는지 확인한다
- 실패한 Job 을 재시작해
BATCH_STEP_EXECUTION.read_count 로 중단 지점을 이어받았는지 검증한다
- Reader 가 상태를 저장하지 않으면 재시작이 처음부터 다시 읽는다는 조용한 사고를 재현한다
.allowStartIfComplete(true) · .startLimit(n) · .noRollback() 의 정확한 의미를 구분한다
선행 스텝: Step 10 — 흐름 제어와 조건 분기
예상 소요: 120분
배치가 10만 건을 처리하다 3만 건째에서 예외 하나를 만났습니다. 선택지는 셋입니다.
- 죽는다 — 3만 건은 커밋됐고 나머지는 안 됐습니다. 손으로 이어 붙여야 합니다.
- 건너뛴다(skip) — 그 한 건만 버리고 계속 갑니다. 대신 버렸다는 사실을 아무도 모르게 될 위험이 생깁니다.
- 다시 해 본다(retry) — 일시적 장애(데드락, 커넥션 끊김)라면 몇 초 뒤엔 성공할 수 있습니다.
Spring Batch 는 셋 다 지원합니다. 문제는 셋을 잘못 조합하면 에러 없이 조용히 데이터가 새거나, 배치가 6배 느려진다는 점입니다. 이 스텝은 그 대가를 전부 숫자로 확인합니다.
11-0. 실습 준비 — 불량 데이터 심기
정산 대상은 orders 의 COMPLETED 70,000건이고, 청크 1,000이면 정확히 70청크입니다(프로젝트 셋업 문서 참조).
여기에 두 종류의 사고를 인위적으로 심습니다.
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t <<'SQL'
-- (A) Processor 에서 터질 불량 주문: 금액이 음수. order_id 가 1000의 배수인 100건.
-- order_id % 10 = 0 이므로 이 100건은 전부 status='COMPLETED' 입니다.
UPDATE orders SET amount = -1.00 WHERE order_id % 1000 = 0;
-- (B) Writer 에서 터질 중복 정산: settlement 에 미리 100행을 넣어 UNIQUE 키를 충돌시킵니다.
TRUNCATE TABLE settlement;
INSERT INTO settlement (order_id, customer_id, settle_date, gross_amount, fee_rate, fee_amount, net_amount)
SELECT o.order_id, o.customer_id, DATE(o.ordered_at), 1000.00, 0.0300, 30.00, 970.00
FROM orders o WHERE o.status = 'COMPLETED' AND o.order_id % 1000 = 0;
SELECT (SELECT COUNT(*) FROM orders WHERE status='COMPLETED') AS target,
(SELECT COUNT(*) FROM orders WHERE amount < 0) AS bad_amount,
(SELECT COUNT(*) FROM settlement) AS preloaded;
SQL
결과
+--------+------------+-----------+
| target | bad_amount | preloaded |
+--------+------------+-----------+
| 70000 | 100 | 100 |
+--------+------------+-----------+
불량 100건이 70청크에 고르게 흩어져 있다는 점이 중요합니다. COMPLETED 700건마다 불량이 1건이므로, 청크 크기가 1,000이면 70청크 전부가 최소 1건의 불량을 품습니다. 뒤에서 이 배치가 왜 그렇게까지 느려지는지의 열쇠가 됩니다.
되돌리려면:
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -e "
UPDATE orders o JOIN (SELECT order_id, 1000 + (order_id % 977) * 100 AS amt FROM orders) s
ON o.order_id = s.order_id SET o.amount = s.amt WHERE o.amount < 0;
TRUNCATE TABLE settlement;"
11-1. 아무 대비 없는 Step 은 첫 예외에서 통째로 멈춘다
Step 07 에서 만든 정산 Step 을 그대로 돌립니다. Processor 는 금액이 0 이하면 예외를 던집니다.
@Bean
public ItemProcessor<Order, Settlement> settlementProcessor(JdbcTemplate jdbc) {
return order -> {
if (order.amount().signum() <= 0) {
throw new IllegalStateException(
"음수 금액 주문: order_id=" + order.orderId() + ", amount=" + order.amount());
}
BigDecimal rate = feeRateOf(jdbc, order.customerId());
BigDecimal fee = order.amount().multiply(rate).setScale(2, RoundingMode.HALF_UP);
return new Settlement(order.orderId(), order.customerId(),
order.orderedAt().toLocalDate(), order.amount(), rate, fee,
order.amount().subtract(fee));
};
}
@Bean
public Step settlementStep(JobRepository repo, PlatformTransactionManager tx,
ItemReader<Order> reader, ItemProcessor<Order, Settlement> processor,
ItemWriter<Settlement> writer) {
return new StepBuilder("settlementStep", repo)
.<Order, Settlement>chunk(1000, tx)
.reader(reader)
.processor(processor)
.writer(writer)
.build(); // 내결함성 없음
}
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-03-01'
결과
INFO 42117 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] launched with the following parameters: [{'date':'{value=2025-03-01, type=class java.lang.String, identifying=true}'}]
INFO 42117 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
ERROR 42117 --- [ main] o.s.batch.core.step.AbstractStep : Encountered an error executing step settlementStep in job settlementJob
java.lang.IllegalStateException: 음수 금액 주문: order_id=1000, amount=-1.00
at com.example.batch.step11.SettlementProcessor.process(SettlementProcessor.java:34) ~[main/:na]
at org.springframework.batch.core.step.item.SimpleChunkProcessor.doProcess(SimpleChunkProcessor.java:127) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.batch.core.step.item.SimpleChunkProcessor.transform(SimpleChunkProcessor.java:301) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.batch.core.step.item.SimpleChunkProcessor.process(SimpleChunkProcessor.java:210) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.batch.core.step.item.ChunkOrientedTasklet.execute(ChunkOrientedTasklet.java:75) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.batch.core.step.tasklet.TaskletStep$ChunkTransactionCallback.doInTransaction(TaskletStep.java:407) ~[spring-batch-core-5.1.1.jar:5.1.1]
...
INFO 42117 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 141ms
INFO 42117 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [FAILED] in 268ms
order_id=1000 은 첫 청크 안에 있습니다. 70,000건 중 999건도 못 쓰고 죽었습니다.
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t -e "
SELECT step_name, status, read_count, write_count, commit_count, rollback_count
FROM BATCH_STEP_EXECUTION ORDER BY step_execution_id DESC LIMIT 1;"
결과
+----------------+--------+------------+-------------+--------------+----------------+
| step_name | status | read_count | write_count | commit_count | rollback_count |
+----------------+--------+------------+-------------+--------------+----------------+
| settlementStep | FAILED | 0 | 0 | 0 | 1 |
+----------------+--------+------------+-------------+--------------+----------------+
read_count = 0 입니다. 1,000건을 읽긴 했지만 그 청크의 트랜잭션이 롤백되면서 카운터도 함께 되돌아갔습니다. StepExecution 의 카운터는 커밋된 것만 셉니다. 이 성질은 11-10 의 재시작 실습에서 그대로 쓰입니다.
11-2. .faultTolerant() 와 .skip()
.faultTolerant() 를 호출하면 빌더가 SimpleStepBuilder → FaultTolerantStepBuilder 로 바뀌고, 내부의 청크 처리기가 SimpleChunkProcessor → FaultTolerantChunkProcessor 로 교체됩니다. 이 교체가 성능 특성 전체를 바꿉니다(11-5).
return new StepBuilder("settlementStep", repo)
.<Order, Settlement>chunk(1000, tx)
.reader(reader)
.processor(processor)
.writer(writer)
.faultTolerant()
.skip(IllegalStateException.class)
.skipLimit(200)
.build();
결과
INFO 42214 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
WARN 42214 --- [ main] o.s.b.c.s.i.FaultTolerantChunkProcessor : Skipping item on process: java.lang.IllegalStateException: 음수 금액 주문: order_id=1000, amount=-1.00
WARN 42214 --- [ main] o.s.b.c.s.i.FaultTolerantChunkProcessor : Skipping item on process: java.lang.IllegalStateException: 음수 금액 주문: order_id=2000, amount=-1.00
...
INFO 42214 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 8s 312ms
INFO 42214 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [COMPLETED] in 8s 461ms
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t -e "
SELECT read_count, write_count, filter_count,
read_skip_count, process_skip_count, write_skip_count,
commit_count, rollback_count, status
FROM BATCH_STEP_EXECUTION ORDER BY step_execution_id DESC LIMIT 1;"
결과
+------------+-------------+--------------+-----------------+--------------------+------------------+--------------+----------------+-----------+
| read_count | write_count | filter_count | read_skip_count | process_skip_count | write_skip_count | commit_count | rollback_count | status |
+------------+-------------+--------------+-----------------+--------------------+------------------+--------------+----------------+-----------+
| 70000 | 69900 | 0 | 0 | 100 | 0 | 70 | 100 | COMPLETED |
+------------+-------------+--------------+-----------------+--------------------+------------------+--------------+----------------+-----------+
읽은 카운터 6종의 의미를 한 번 정리합니다.
| 카운터 | 증가 조건 |
|---|
read_count | Reader 가 아이템을 하나 돌려줄 때마다 |
write_count | Writer 에 실제로 넘어간 아이템 수 |
filter_count | Processor 가 null 을 반환해 걸러낸 수 (skip 아님) |
read_skip_count | 읽기 중 예외를 skip 한 수 |
process_skip_count | 처리 중 예외를 skip 한 수 |
write_skip_count | 쓰기 중 예외를 skip 한 수 |
commit_count | 청크 트랜잭션 커밋 횟수 |
rollback_count | 청크 트랜잭션 롤백 횟수 |
⚠️ 함정 — filter_count 와 process_skip_count 는 완전히 다릅니다
Processor 가 null 을 반환하면 filter_count 가 오르고, 아무 로그도 남지 않으며, 예외도 없습니다.
"정산 대상이 아니면 null 반환" 같은 코드는 정상 필터링이지만, null 반환을 예외 처리 대용으로 쓰면 조용히 데이터가 사라집니다.
버려도 되는 데이터면 null(filter), 버리면 안 되는데 어쩔 수 없이 버리는 데이터면 예외 + skip 입니다. skip 은 최소한 카운터와 WARN 로그를 남깁니다.
rollback_count = 100 에 주목하세요. 커밋은 70번인데 롤백이 100번입니다. skip 한 건마다 롤백이 한 번씩 일어났습니다. 이유는 11-5 에서 밝힙니다.
11-3. skipLimit 의 기본값과 예외 계층
기본값은 10 입니다
.skipLimit() 을 생략하면 어떻게 될까요?
.faultTolerant()
.skip(IllegalStateException.class) // skipLimit 생략
.build();
결과
WARN 42301 --- [ main] o.s.b.c.s.i.FaultTolerantChunkProcessor : Skipping item on process: java.lang.IllegalStateException: 음수 금액 주문: order_id=1000, amount=-1.00
...(10건)
ERROR 42301 --- [ main] o.s.batch.core.step.AbstractStep : Encountered an error executing step settlementStep in job settlementJob
org.springframework.batch.core.step.skip.SkipLimitExceededException: Skip limit of '10' exceeded
at org.springframework.batch.core.step.skip.LimitCheckingItemSkipPolicy.shouldSkip(LimitCheckingItemSkipPolicy.java:120) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.batch.core.step.item.FaultTolerantChunkProcessor$2.recover(FaultTolerantChunkProcessor.java:290) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.retry.support.RetryTemplate.handleRetryExhausted(RetryTemplate.java:557) ~[spring-retry-2.0.5.jar:na]
...
Caused by: java.lang.IllegalStateException: 음수 금액 주문: order_id=11000, amount=-1.00
at com.example.batch.step11.SettlementProcessor.process(SettlementProcessor.java:34) ~[main/:na]
Skip limit of '10' exceeded. Spring Batch 5.1 의 FaultTolerantStepBuilder 는 skipLimit 기본값이 10 입니다. 무제한이 아닙니다.
💡 실무 팁 — skipLimit 은 "허용치"가 아니라 "경보 임계값"입니다
skipLimit(Integer.MAX_VALUE) 로 두면 배치는 절대 안 죽지만, 10만 건 중 8만 건이 버려져도 COMPLETED 로 끝납니다. 그건 성공이 아닙니다.
정상 데이터 품질에서 나올 법한 skip 수의 2~3배 정도로 잡아 두면, 데이터가 갑자기 망가졌을 때 배치가 시끄럽게 실패합니다. 조용한 성공보다 시끄러운 실패가 낫습니다.
예외 계층과 .noSkip()
.skip() 은 하위 타입까지 포함합니다. .skip(Exception.class) 는 사실상 전부입니다.
.faultTolerant()
.skip(Exception.class) // 전부 건너뛴다
.noSkip(DataAccessResourceFailureException.class) // 단, DB 가 죽은 건 건너뛰지 않는다
.noSkip(OutOfMemoryError.class) // (Error 는 애초에 Exception 이 아니지만 명시)
.skipLimit(500)
판정은 BinaryExceptionClassifier 가 하며, 가장 구체적인(상속 거리가 가까운) 등록이 이깁니다. 등록 순서는 무관합니다.
| 던져진 예외 | .skip(Exception) + .noSkip(DataAccessResourceFailureException) |
|---|
IllegalStateException | skip |
DataIntegrityViolationException | skip (DataAccessException 하위지만 noSkip 대상 아님) |
DataAccessResourceFailureException | skip 안 함 → Step FAILED |
CannotAcquireLockException | skip |
⚠️ 함정 — .skip(Exception.class) 는 인프라 장애까지 건너뜁니다
DB 커넥션 풀이 고갈되면 아이템마다 CannotGetJdbcConnectionException 이 납니다. .skip(Exception.class).skipLimit(100000) 이면 배치는 7만 건을 전부 skip 하고 COMPLETED 로 끝납니다.
settlement 테이블은 비어 있는데 Job 은 성공. 다음 날 아침 정산 담당자가 발견합니다.
skip 대상은 "데이터가 나쁜 경우"로 한정하고, "인프라가 나쁜 경우"는 noSkip 으로 빼세요.
11-4. 커스텀 SkipPolicy — 조건을 코드로 쓴다
"예외 타입"만으로 판단이 안 되는 경우가 있습니다. SkipPolicy 를 직접 구현합니다.
public static class SettlementSkipPolicy implements SkipPolicy {
private final int limit;
public SettlementSkipPolicy(int limit) {
this.limit = limit;
}
@Override
public boolean shouldSkip(Throwable t, long skipCount) throws SkipLimitExceededException {
// (1) 인프라 장애는 절대 건너뛰지 않는다 — 즉시 Step 실패
if (t instanceof DataAccessResourceFailureException
|| t instanceof CannotGetJdbcConnectionException) {
return false;
}
// (2) 데이터 품질 문제는 한도까지 건너뛴다
if (t instanceof IllegalStateException || t instanceof DataIntegrityViolationException) {
if (skipCount >= limit) {
throw new SkipLimitExceededException(limit, t); // 한도 초과는 예외로 알린다
}
return true;
}
// (3) 나머지는 모르는 예외 → 건너뛰지 않는다
return false;
}
}
⚠️ 함정 — shouldSkip 의 두 번째 인자는 "해당 예외의 skip 수"가 아니라 "전체 skip 수"입니다
시그니처가 long skipCount 라서 "이 예외를 몇 번 건너뛰었나"로 오해하기 쉽습니다. 실제로는 읽기·처리·쓰기를 합친 Step 전체의 누적 skip 수입니다.
예외 종류별 한도를 두려면 정책 안에 Map<Class<?>, AtomicInteger> 를 직접 들고 있어야 합니다.
그리고 이 정책 객체는 Step 인스턴스와 수명을 같이하므로, @Bean 싱글턴으로 두면 두 번째 실행에서 카운터가 이어집니다. @StepScope 를 붙이거나 open() 에서 초기화하세요.
등록은 skipPolicy() 입니다.
.faultTolerant()
.skipPolicy(new SettlementSkipPolicy(200))
⚠️ .skipPolicy() 를 쓰면 .skip() / .skipLimit() 은 전부 무시됩니다.
둘을 섞어 쓰면 컴파일도 되고 실행도 되지만 .skip(...) 설정이 조용히 사라집니다. 하나만 고르세요.
INFO 42388 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 8s 297ms
+------------+-------------+--------------------+--------------+----------------+-----------+
| read_count | write_count | process_skip_count | commit_count | rollback_count | status |
+------------+-------------+--------------------+--------------+----------------+-----------+
| 70000 | 69900 | 100 | 70 | 100 | COMPLETED |
+------------+-------------+--------------------+--------------+----------------+-----------+
11-5. skip 의 진짜 비용 — 청크 롤백과 스캔 모드
여기가 이 스텝의 핵심입니다.
왜 롤백이 100번인가
청크는 하나의 트랜잭션입니다. 청크 중간에서 예외가 나면 트랜잭션은 되돌릴 수밖에 없습니다. 그런데 Spring Batch 는 "몇 번째 아이템이 문제인지" 미리 알지 못합니다. 그래서 이렇게 동작합니다.
[처리 단계 예외 — process skip]
청크 12 (아이템 11001..12000)
├ 1차 시도: 11001 처리 … 11700 에서 IllegalStateException ──▶ 트랜잭션 ROLLBACK (rollback_count +1)
├ 2차 시도: 캐시된 같은 1,000건을 처음부터 다시 처리.
│ 11700 은 "skip 대상"으로 기억해 두었으므로 건너뜀 → 999건 Writer 로 → COMMIT
└ 결과: rollback 1, commit 1, 처리 함수 호출 1,999회
처리 단계 skip 은 청크를 한 번 더 처리하는 비용입니다. 아이템 재처리는 메모리 안에서 일어나므로 상대적으로 쌉니다.
쓰기 단계는 완전히 다릅니다.
[쓰기 단계 예외 — write skip]
청크 12 (아이템 11001..12000)
├ 1차 시도: 1,000건을 batch INSERT → DuplicateKeyException ──▶ ROLLBACK (rollback_count +1)
├ ★ 스캔 모드 진입 ★
│ Spring Batch 는 "1,000건 중 누가 범인인지" 모릅니다.
│ 그래서 아이템을 하나씩, 각각 독립 트랜잭션으로 다시 씁니다.
│ 11001 → INSERT → COMMIT (commit_count +1)
│ 11002 → INSERT → COMMIT (commit_count +1)
│ ...
│ 11700 → INSERT → DuplicateKeyException → ROLLBACK → skip (rollback_count +1, write_skip_count +1)
│ ...
│ 12000 → INSERT → COMMIT (commit_count +1)
└ 결과: 청크 하나에 트랜잭션이 1,001개
청크 크기가 곧 트랜잭션 폭발 계수입니다. 청크 1,000에서 쓰기 skip 이 한 건만 나도 그 청크는 1,000번 커밋합니다.
실측
이제 (B) 시나리오, 즉 settlement 에 미리 넣어 둔 100행과 UNIQUE 키가 충돌하는 상황을 돌립니다.
.faultTolerant()
.skip(DataIntegrityViolationException.class) // DuplicateKeyException 의 상위
.skipLimit(200)
먼저 정상 기준선(불량 데이터 없이, settlement 비운 상태):
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -e "TRUNCATE TABLE settlement;"
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-03-01'
결과
INFO 42455 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 6s 108ms
+------------+-------------+------------------+--------------+----------------+-----------+
| read_count | write_count | write_skip_count | commit_count | rollback_count | status |
+------------+-------------+------------------+--------------+----------------+-----------+
| 70000 | 70000 | 0 | 70 | 0 | COMPLETED |
+------------+-------------+------------------+--------------+----------------+-----------+
6.108초, 커밋 70회. 이제 100행을 미리 심고 다시:
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb <<'SQL'
TRUNCATE TABLE settlement;
INSERT INTO settlement (order_id, customer_id, settle_date, gross_amount, fee_rate, fee_amount, net_amount)
SELECT order_id, customer_id, DATE(ordered_at), 1000.00, 0.0300, 30.00, 970.00
FROM orders WHERE status='COMPLETED' AND order_id % 1000 = 0;
SQL
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-03-02'
결과
WARN 42502 --- [ main] o.s.b.c.s.i.FaultTolerantChunkProcessor : Skipping item on write: org.springframework.dao.DuplicateKeyException: PreparedStatementCallback; SQL [INSERT INTO settlement ...]; Duplicate entry '1000' for key 'settlement.uk_settlement_order'
WARN 42502 --- [ main] o.s.b.c.s.i.FaultTolerantChunkProcessor : Skipping item on write: org.springframework.dao.DuplicateKeyException: ... Duplicate entry '2000' for key 'settlement.uk_settlement_order'
...
INFO 42502 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 34s 712ms
INFO 42502 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [COMPLETED] in 34s 856ms
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t -e "
SELECT read_count, write_count, write_skip_count, commit_count, rollback_count, status
FROM BATCH_STEP_EXECUTION ORDER BY step_execution_id DESC LIMIT 1;"
결과
+------------+-------------+------------------+--------------+----------------+-----------+
| read_count | write_count | write_skip_count | commit_count | rollback_count | status |
+------------+-------------+------------------+--------------+----------------+-----------+
| 70000 | 69900 | 100 | 69900 | 170 | COMPLETED |
+------------+-------------+------------------+--------------+----------------+-----------+
6.108초 → 34.712초. 약 5.7배 느려졌습니다.
증거는 commit_count 에 그대로 있습니다.
| 지표 | 정상 | 쓰기 skip 100건 | 배수 |
|---|
| 소요 시간 | 6.108초 | 34.712초 | 약 5.7배 |
commit_count | 70 | 69,900 | 999배 |
rollback_count | 0 | 170 | — |
| Writer 호출 횟수 | 70 | 70 + 70,000 | — |
rollback_count = 170 의 내역: 각 청크의 1차 시도 실패 70회 + 스캔 중 실제 skip 된 아이템 100회 = 170.
commit_count = 69,900 은 스캔 모드에서 아이템 하나씩 커밋한 횟수입니다(70,000건 중 skip 100건 제외).
⚠️ 함정 — skip 은 "그 한 건만 버리는" 값싼 기능이 아닙니다
많은 사람이 .skip(Exception.class).skipLimit(10000) 을 "안전장치"라고 생각하고 걸어 둡니다. 데이터가 깨끗할 땐 아무 일도 없습니다. 문제는 데이터가 나빠진 날입니다.
불량률이 0.14%(7만 건 중 100건) 올라간 것만으로 배치가 6초에서 35초가 됩니다. 불량 100건이 모든 청크에 하나씩 흩어져 있으면 전체 청크가 스캔 모드로 떨어지기 때문입니다.
실제 운영에서 "어제까지 20분이던 배치가 오늘 3시간 걸린다"의 상당수가 이 현상입니다. 코드는 그대로고, 에러도 없고, Job 은 COMPLETED 로 끝납니다.
불량이 뭉쳐 있으면 어떻게 다른가
같은 100건이라도 앞쪽 1청크에 몰려 있으면 이야기가 달라집니다.
# 불량을 order_id 1..1000 안에만 심는다 (COMPLETED 700건 중 100건)
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb <<'SQL'
TRUNCATE TABLE settlement;
INSERT INTO settlement (order_id, customer_id, settle_date, gross_amount, fee_rate, fee_amount, net_amount)
SELECT order_id, customer_id, DATE(ordered_at), 1000.00, 0.0300, 30.00, 970.00
FROM orders WHERE status='COMPLETED' AND order_id <= 1000 LIMIT 100;
SQL
결과
INFO 42571 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 6s 594ms
+------------+-------------+------------------+--------------+----------------+
| read_count | write_count | write_skip_count | commit_count | rollback_count |
+------------+-------------+------------------+--------------+----------------+
| 70000 | 69900 | 100 | 1069 | 101 |
+------------+-------------+------------------+--------------+----------------+
6.594초. skip 수는 똑같이 100건인데 34.7초가 아니라 6.6초입니다. 스캔 모드에 떨어진 청크가 70개가 아니라 1개이기 때문입니다(commit_count 1,069 = 정상 69청크 + 스캔 1,000회).
skip 비용은 "skip 건수"가 아니라 "skip 이 몇 개의 청크에 퍼져 있는가"에 비례합니다.
11-6. skip 비용을 줄이는 네 가지 방법
| 방법 | 내용 | 효과 |
|---|
| ① Writer 이전으로 옮기기 | 쓰기에서 터질 조건을 Processor 에서 미리 판정해 예외를 던지거나 null 을 반환 | 스캔 모드 자체가 사라짐 |
| ② Reader 쿼리에서 제외 | WHERE amount > 0 AND NOT EXISTS (SELECT 1 FROM settlement s WHERE s.order_id = o.order_id) | skip 이 0건이 됨 |
| ③ 청크 크기 줄이기 | 스캔 비용은 청크 크기에 비례. 1,000 → 200 이면 스캔 트랜잭션이 1/5 | 부분 완화 |
| ④ 멱등 쓰기 | INSERT ... ON DUPLICATE KEY UPDATE 로 중복 자체를 예외로 만들지 않음 | 예외가 안 나므로 최선 |
②를 적용해 다시 재봅니다.
// Reader 의 SQL 에 한 줄 추가
String sql = """
SELECT o.order_id, o.customer_id, o.amount, o.status, o.ordered_at
FROM orders o
WHERE o.status = 'COMPLETED'
AND o.amount > 0
AND NOT EXISTS (SELECT 1 FROM settlement s WHERE s.order_id = o.order_id)
ORDER BY o.order_id
""";
결과
INFO 42640 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 6s 233ms
+------------+-------------+------------------+--------------+----------------+
| read_count | write_count | write_skip_count | commit_count | rollback_count |
+------------+-------------+------------------+--------------+----------------+
| 69900 | 69900 | 0 | 70 | 0 |
+------------+-------------+------------------+--------------+----------------+
34.712초 → 6.233초. 약 5.6배 빨라졌습니다. skip 을 잘 쓰는 최고의 방법은 skip 이 안 나게 하는 것입니다.
💡 실무 팁 — skip 은 "예상 못 한 소수"를 위한 것입니다
예상되는 제외 조건(취소 주문, 중복 정산, 음수 금액)은 Reader 의 WHERE 절에 넣는 것이 정답입니다.
skip 은 "쿼리로는 못 거르는, 처리 중에만 알 수 있는, 드물게 발생하는" 예외를 위한 안전망입니다. 안전망으로 주행하면 안 됩니다.
11-7. .retry() 와 backoff
skip 은 "포기"이고 retry 는 "다시 해 보기"입니다. 일시적(transient) 장애에만 의미가 있습니다.
.faultTolerant()
.retry(DeadlockLoserDataAccessException.class)
.retry(CannotAcquireLockException.class)
.retryLimit(3)
.backOffPolicy(exponentialBackOff())
.skip(DataIntegrityViolationException.class)
.skipLimit(200)
private static BackOffPolicy exponentialBackOff() {
ExponentialBackOffPolicy policy = new ExponentialBackOffPolicy();
policy.setInitialInterval(200L); // 첫 재시도 전 200ms
policy.setMultiplier(2.0); // 200 → 400 → 800
policy.setMaxInterval(5_000L);
return policy;
}
Writer 가 세 번째 청크에서 데드락을 던지도록 흉내 내면:
결과
INFO 42711 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
DEBUG 42711 --- [ main] o.s.retry.support.RetryTemplate : Retry: count=0
WARN 42711 --- [ main] o.s.b.c.s.item.FaultTolerantChunkProcessor: Retryable exception on write: org.springframework.dao.DeadlockLoserDataAccessException: Deadlock found when trying to get lock; try restarting transaction
DEBUG 42711 --- [ main] o.s.retry.backoff.ExponentialBackOffPolicy: Sleeping for 200
DEBUG 42711 --- [ main] o.s.retry.support.RetryTemplate : Retry: count=1
WARN 42711 --- [ main] o.s.b.c.s.item.FaultTolerantChunkProcessor: Retryable exception on write: org.springframework.dao.DeadlockLoserDataAccessException: Deadlock found when trying to get lock; try restarting transaction
DEBUG 42711 --- [ main] o.s.retry.backoff.ExponentialBackOffPolicy: Sleeping for 400
DEBUG 42711 --- [ main] o.s.retry.support.RetryTemplate : Retry: count=2
INFO 42711 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 6s 913ms
+------------+-------------+------------------+--------------+----------------+-----------+
| read_count | write_count | write_skip_count | commit_count | rollback_count | status |
+------------+-------------+------------------+--------------+----------------+-----------+
| 70000 | 70000 | 0 | 70 | 2 | COMPLETED |
+------------+-------------+------------------+--------------+----------------+-----------+
Retry: count=0 → count=1 → count=2 로 총 3회 시도했고, 3회차에서 성공했습니다. rollback_count = 2 는 실패한 두 번의 시도입니다. skip 은 0건이므로 데이터는 온전합니다.
⚠️ 함정 — .retryLimit(3) 은 "재시도 3번"이 아니라 "총 시도 3회"입니다
내부적으로 SimpleRetryPolicy(maxAttempts = retryLimit) 이 만들어집니다. 최초 시도가 1회에 포함되므로 실제 재시도는 2번입니다.
"3번 더 시도하겠지" 하고 SLA 를 계산하면 어긋납니다. 재시도를 N번 하고 싶으면 .retryLimit(N + 1) 입니다.
⚠️ 함정 — .retry() 만 쓰고 .retryLimit() 을 빠뜨리면 아무 일도 안 일어납니다
FaultTolerantStepBuilder 의 retryLimit 기본값은 0 입니다. 이 상태에서는 등록한 예외가 재시도 없이 그대로 전파됩니다.
컴파일 에러도, 경고 로그도 없습니다. "retry 를 걸었는데 왜 안 되지?"의 절반은 이것입니다. .retry() 와 .retryLimit() 은 항상 짝으로 쓰세요.
💡 실무 팁 — backoff 없는 retry 는 장애를 키웁니다
backOffPolicy 를 생략하면 기본이 NoBackOffPolicy 라 즉시 재시도합니다. DB 가 부하로 죽어 가는 중이라면 재시도가 부하를 더합니다.
데드락·락 타임아웃·커넥션 고갈은 전부 "잠깐 기다리면 풀리는" 종류입니다. ExponentialBackOffPolicy 를 기본으로 쓰세요.
11-8. retry 와 skip 이 만나면
같은 예외를 .retry() 와 .skip() 에 모두 등록하면 retry 를 다 소진한 뒤 skip 으로 넘어갑니다.
.faultTolerant()
.retry(DataIntegrityViolationException.class).retryLimit(3)
.skip(DataIntegrityViolationException.class).skipLimit(200)
한 아이템의 생애:
아이템 11700 (중복 키)
├ 시도 1 → DuplicateKeyException → 롤백 → 재시도
├ 시도 2 → DuplicateKeyException → 롤백 → 재시도
├ 시도 3 → DuplicateKeyException → retry 소진
└ RetryTemplate 이 recover() 호출 → SkipPolicy 판정 → skip (write_skip_count +1)
결과
INFO 42780 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 1m 42s 337ms
+------------+-------------+------------------+--------------+----------------+-----------+
| read_count | write_count | write_skip_count | commit_count | rollback_count | status |
+------------+-------------+------------------+--------------+----------------+-----------+
| 70000 | 69900 | 100 | 69900 | 370 | COMPLETED |
+------------+-------------+------------------+--------------+----------------+-----------+
34.712초 → 102.337초. 약 2.9배 더 느려졌습니다. 절대 성공할 수 없는 예외(중복 키)를 3번씩 재시도한 뒤에야 skip 하기 때문입니다. rollback_count 가 170 → 370 으로 늘어난 것이 그 흔적입니다.
| 예외 성격 | retry | skip |
|---|
| 데드락 / 락 타임아웃 / 일시적 커넥션 끊김 | O | 최후 수단 |
| 외부 API 5xx · 타임아웃 | O | 상황에 따라 |
| 중복 키 / 제약 위반 | X (재시도해도 같음) | O |
| 데이터 형식 오류 · 파싱 실패 | X | O |
| NPE · 로직 버그 | X | X — 고쳐야 합니다 |
⚠️ 함정 — "일단 둘 다 걸어 두자"가 가장 흔한 실수입니다
.retry(Exception.class).skip(Exception.class) 조합은 문법적으로 완벽하고, 정상 데이터에서는 아무 차이가 없습니다.
불량 데이터가 들어온 날 배치는 정상 대비 17배(6.1초 → 102.3초) 느려지고, 여전히 COMPLETED 로 끝납니다.
재시도해서 성공할 가능성이 있는 예외에만 retry 를 거세요.
11-9. noRollback — 롤백시키지 않을 예외
기본적으로 청크에서 나간 예외는 전부 트랜잭션을 롤백시킵니다. 하지만 "예외로 알리되 DB 작업은 유지"하고 싶을 때가 있습니다. 대표적으로 검증 실패입니다.
.faultTolerant()
.skip(ValidationException.class)
.skipLimit(500)
.noRollback(ValidationException.class) // 이 예외는 트랜잭션을 롤백시키지 않는다
noRollback 없이 ValidationException 100건이 발생한 경우와 비교합니다.
결과 (noRollback 없음)
INFO 42841 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 8s 271ms
+------------+-------------+--------------------+--------------+----------------+
| read_count | write_count | process_skip_count | commit_count | rollback_count |
+------------+-------------+--------------------+--------------+----------------+
| 70000 | 69900 | 100 | 70 | 100 |
+------------+-------------+--------------------+--------------+----------------+
결과 (noRollback 적용)
INFO 42868 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 6s 402ms
+------------+-------------+--------------------+--------------+----------------+
| read_count | write_count | process_skip_count | commit_count | rollback_count |
+------------+-------------+--------------------+--------------+----------------+
| 70000 | 69900 | 100 | 70 | 0 |
+------------+-------------+--------------------+--------------+----------------+
8.271초 → 6.402초 (약 1.3배). rollback_count 가 100 → 0 이 되었습니다. 청크 재처리가 사라졌기 때문입니다.
⚠️ 함정 — noRollback 을 DB 예외에 걸면 데이터가 깨집니다
noRollback(DataIntegrityViolationException.class) 같은 설정은 절대 하면 안 됩니다.
DB 가 제약 위반을 던진 시점에 그 트랜잭션은 이미 오염됐고, MySQL 커넥션은 롤백을 기대합니다. 롤백을 건너뛰면 이후 문장이 Transaction is marked rollback-only 로 터지거나, 최악의 경우 반쯤 쓰인 청크가 커밋됩니다.
noRollback 은 DB 를 건드리지 않은, 순수 애플리케이션 검증 예외에만 쓰세요. Processor 가 던지는 ValidationException 이 정확히 그 자리입니다.
11-10. 재시작 — 중단 지점을 이어받는다
이제 skip 없이, 40,001번째 아이템에서 확실히 죽는 Step 을 만듭니다.
@Bean
@StepScope
public ItemProcessor<Order, Settlement> failingProcessor(
@Value("#{jobParameters['failAt']}") Long failAt) {
AtomicLong seen = new AtomicLong();
return order -> {
if (failAt != null && seen.incrementAndGet() == failAt) {
throw new IllegalStateException("의도적 실패: " + failAt + "번째 아이템");
}
return toSettlement(order);
};
}
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -e "TRUNCATE TABLE settlement;"
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-04-01 failAt=40001'
결과
INFO 42911 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
ERROR 42911 --- [ main] o.s.batch.core.step.AbstractStep : Encountered an error executing step settlementStep in job settlementJob
java.lang.IllegalStateException: 의도적 실패: 40001번째 아이템
at com.example.batch.step11.RestartConfig.lambda$failingProcessor$0(RestartConfig.java:58) ~[main/:na]
at org.springframework.batch.core.step.item.SimpleChunkProcessor.doProcess(SimpleChunkProcessor.java:127) ~[spring-batch-core-5.1.1.jar:5.1.1]
...
INFO 42911 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [{'date':'{value=2025-04-01,...}', 'failAt':'{value=40001,...}'}] and the following status: [FAILED] in 3s 981ms
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t -e "
SELECT step_execution_id, status, read_count, write_count, commit_count, rollback_count
FROM BATCH_STEP_EXECUTION ORDER BY step_execution_id DESC LIMIT 1;
SELECT COUNT(*) AS settled FROM settlement;"
결과
+-------------------+--------+------------+-------------+--------------+----------------+
| step_execution_id | status | read_count | write_count | commit_count | rollback_count |
+-------------------+--------+------------+-------------+--------------+----------------+
| 17 | FAILED | 40000 | 40000 | 40 | 1 |
+-------------------+--------+------------+-------------+--------------+----------------+
+---------+
| settled |
+---------+
| 40000 |
+---------+
40청크(40,000건)까지 커밋됐습니다. 41번째 청크는 롤백되어 카운터에 반영되지 않았습니다.
Reader 가 남긴 상태를 봅니다.
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -e "
SELECT short_context FROM BATCH_STEP_EXECUTION_CONTEXT WHERE step_execution_id = 17\G"
결과
*************************** 1. row ***************************
short_context: {"@class":"java.util.HashMap","orderReader.read.count":40000,"batch.taskletType":"org.springframework.batch.core.step.item.ChunkOrientedTasklet","batch.stepType":"org.springframework.batch.core.step.tasklet.TaskletStep"}
orderReader.read.count: 40000. 이것이 재시작의 전부입니다. Reader 가 ItemStream.update() 에서 청크 커밋마다 저장한 값입니다.
같은 파라미터로 다시 실행 = 재시작
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-04-01'
failAt 을 뺐습니다. 하지만 failAt 은 identifying=true 인 파라미터라 JobInstance 가 달라져 버립니다. 재시작하려면 식별 파라미터가 완전히 같아야 합니다.
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-04-01 failAt=999999'
failAt=999999 도 다른 인스턴스입니다. 실무에서는 실패를 유발한 조건만 고치고 같은 파라미터로 다시 던집니다. 여기서는 failAt 을 identifying=false 로 선언해 두었다고 가정하고 그대로 재실행합니다.
./gradlew bootRun --args='--spring.batch.job.name=settlementJob date=2025-04-01'
결과
INFO 42977 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] launched with the following parameters: [{'date':'{value=2025-04-01, type=class java.lang.String, identifying=true}'}]
INFO 42977 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
INFO 42977 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 2s 617ms
INFO 42977 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [COMPLETED] in 2s 744ms
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t -e "
SELECT step_execution_id, status, read_count, write_count, commit_count
FROM BATCH_STEP_EXECUTION ORDER BY step_execution_id DESC LIMIT 2;
SELECT COUNT(*) AS settled FROM settlement;"
결과
+-------------------+-----------+------------+-------------+--------------+
| step_execution_id | status | read_count | write_count | commit_count |
+-------------------+-----------+------------+-------------+--------------+
| 18 | COMPLETED | 30000 | 30000 | 30 |
| 17 | FAILED | 40000 | 40000 | 40 |
+-------------------+-----------+------------+-------------+--------------+
+---------+
| settled |
+---------+
| 70000 |
+---------+
두 번째 실행의 read_count 는 70,000 이 아니라 30,000 입니다. 40,000건을 건너뛰고 나머지 30,000건만 처리했습니다. settlement 는 정확히 70,000행. 중복도 누락도 없습니다.
| 1차 실행 | 2차 실행(재시작) | 합계 |
|---|
read_count | 40,000 | 30,000 | 70,000 |
write_count | 40,000 | 30,000 | 70,000 |
commit_count | 40 | 30 | 70 |
status | FAILED | COMPLETED | |
11-11. 함정 — Reader 가 상태를 저장하지 않으면
여기가 이 스텝에서 가장 조용하고 가장 비싼 사고입니다.
Reader 를 직접 만들어 봅니다. 아주 자연스러운 코드입니다.
// ⚠️ 이 코드는 컴파일도 되고 정상 실행도 되지만, 재시작이 망가집니다.
@Bean
public ItemReader<Order> naiveReader(JdbcTemplate jdbc) {
List<Order> all = jdbc.query(SQL, new DataClassRowMapper<>(Order.class));
Iterator<Order> it = all.iterator();
return () -> it.hasNext() ? it.next() : null; // ItemStream 을 구현하지 않음
}
Writer 는 중복 사고가 시끄럽게 나지 않도록 upsert 라고 가정합니다(실무에서 흔한 선택입니다).
String sql = """
INSERT INTO settlement (order_id, customer_id, settle_date, gross_amount, fee_rate, fee_amount, net_amount)
VALUES (:orderId, :customerId, :settleDate, :grossAmount, :feeRate, :feeAmount, :netAmount)
ON DUPLICATE KEY UPDATE gross_amount = VALUES(gross_amount), net_amount = VALUES(net_amount)
""";
40,001에서 실패시키고 재시작합니다.
결과 (1차 — 실패)
+-------------------+--------+------------+-------------+--------------+
| step_execution_id | status | read_count | write_count | commit_count |
+-------------------+--------+------------+-------------+--------------+
| 21 | FAILED | 40000 | 40000 | 40 |
+-------------------+--------+------------+-------------+--------------+
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -e "
SELECT short_context FROM BATCH_STEP_EXECUTION_CONTEXT WHERE step_execution_id = 21\G"
결과
*************************** 1. row ***************************
short_context: {"@class":"java.util.HashMap","batch.taskletType":"org.springframework.batch.core.step.item.ChunkOrientedTasklet","batch.stepType":"org.springframework.batch.core.step.tasklet.TaskletStep"}
read.count 키가 없습니다. Reader 가 ItemStream 이 아니므로 Step 이 저장할 상태 자체가 없었습니다.
결과 (2차 — 재시작)
INFO 43044 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 6s 155ms
+-------------------+-----------+------------+-------------+--------------+
| step_execution_id | status | read_count | write_count | commit_count |
+-------------------+-----------+------------+-------------+--------------+
| 22 | COMPLETED | 70000 | 70000 | 70 |
+-------------------+-----------+------------+-------------+--------------+
+---------+
| settled |
+---------+
| 70000 |
+---------+
Job 은 COMPLETED 이고 settlement 도 70,000행으로 정확합니다. 아무 문제 없어 보입니다.
그런데 2차 실행의 read_count 가 30,000 이 아니라 70,000 입니다. 처음부터 다시 읽은 것입니다. 정상 재시작(11-10)의 2.617초가 여기서는 6.155초, 약 2.4배입니다.
⚠️ 함정 — 재시작이 "동작하는 것처럼 보이는" 것이 가장 위험합니다
결과 테이블이 맞으니 아무도 눈치채지 못합니다. 하지만 실제로 벌어진 일은 이렇습니다.
| 정상 Reader | 상태 없는 Reader |
|---|
재시작 시 read_count | 30,000 | 70,000 |
| Processor 호출 횟수 | 30,000 | 70,000 (4만 건은 두 번째) |
| 소요 시간 | 2.617초 | 6.155초 (약 2.4배) |
| 결과 테이블 | 정확 | (upsert 덕에) 정확 |
Processor 에 부작용이 하나라도 있으면 그 순간 사고가 됩니다. 40,000건에 대해 알림 메일이 두 번 나가고, 포인트가 두 번 적립되고, 외부 정산 API 가 두 번 호출됩니다.
upsert 가 결과 테이블만 지켜 줄 뿐 부작용은 못 막습니다.
해결책 세 가지
ItemStreamReader 를 구현하거나, JdbcCursorItemReader / JdbcPagingItemReader / FlatFileItemReader 등 ItemStream 을 구현한 기본 Reader 를 쓴다(전부 saveState = true 가 기본).
- Reader 쿼리를 "아직 처리 안 된 것만" 으로 만든다(
NOT EXISTS). 상태를 DB 가 대신 들고 있게 하는 방식이며 가장 견고합니다.
- Processor 를 멱등하게 만든다. 부작용은 Step 밖으로 빼거나, 처리 이력 테이블로 이중 발사를 막습니다.
관련해서 하나 더:
⚠️ 함정 — .saveState(false) 와 비유니크 sortKey
JdbcCursorItemReader, JdbcPagingItemReader 는 saveState(false) 로 끄면 위와 똑같은 상황이 됩니다. 멀티스레드 Step(Step 13)에서 상태 저장이 무의미해 끄는 경우가 있는데, 그 순간 재시작 능력을 포기하는 것입니다.
JdbcPagingItemReader 의 sortKey 가 유니크하지 않으면 페이지 경계에서 같은 값이 잘리며 행이 중복되거나 누락됩니다. ORDER BY ordered_at 처럼 중복 가능한 컬럼 단독은 금지입니다. ordered_at, order_id 처럼 PK 를 뒤에 붙이세요.
둘 다 에러가 나지 않습니다. 카운터만 이상해집니다.
롤백 후 재처리 — Processor 는 두 번 호출된다
11-5 에서 본 "청크 재처리"에는 부수 효과가 하나 더 있습니다.
@Bean
public ItemProcessor<Order, Settlement> sideEffectProcessor(NotificationClient client) {
return order -> {
if (order.amount().signum() <= 0) throw new IllegalStateException("음수 금액");
client.notifySettled(order.orderId()); // ← 외부 호출. 부작용.
return toSettlement(order);
};
}
불량 100건이 70청크에 흩어진 상태로 돌리면:
결과
+------------+-------------+--------------------+----------------+
| read_count | write_count | process_skip_count | rollback_count |
+------------+-------------+--------------------+----------------+
| 70000 | 69900 | 100 | 100 |
+------------+-------------+--------------------+----------------+
[NotificationClient] 총 호출 횟수: 139,251
69,900건을 처리했는데 알림은 139,251번 나갔습니다. 롤백된 70청크가 통째로 재처리되면서 그 안의 정상 아이템들이 두 번씩 Processor 를 통과했기 때문입니다.
기본값 processorTransactional = true 는 "Processor 출력은 트랜잭션과 함께 버려지고 다시 만든다"는 뜻입니다. .processorNonTransactional() 을 붙이면 출력을 캐시해 재처리를 건너뜁니다.
.faultTolerant()
.skip(IllegalStateException.class).skipLimit(200)
.processorNonTransactional()
결과
[NotificationClient] 총 호출 횟수: 69,900
⚠️ .processorNonTransactional() 은 Processor 가 순수 함수일 때만 안전합니다.
캐시된 출력을 재사용하므로, Processor 가 내부 상태나 외부 조회 결과에 의존하면 낡은 값을 쓰게 됩니다.
근본 해법은 Processor 에서 부작용을 없애는 것입니다. 알림·API 호출은 별도 Step 이나 SkipListener/ItemWriteListener(Step 12)로 옮기세요.
11-12. allowStartIfComplete 와 startLimit
allowStartIfComplete(true)
COMPLETED 로 끝난 Step 은 재시작 시 건너뜁니다.
INFO 43111 --- [ main] o.s.batch.core.job.SimpleStepHandler : Step already complete or not restartable, so no action to take: StepExecution: id=17, version=4, name=prepareStep, status=COMPLETED, exitStatus=COMPLETED
INFO 43111 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
그런데 "매번 다시 해야 하는" Step 이 있습니다. 임시 테이블 정리, 작업 디렉터리 비우기 같은 것입니다.
@Bean
public Step prepareStep(JobRepository repo, PlatformTransactionManager tx) {
return new StepBuilder("prepareStep", repo)
.tasklet(truncateStagingTasklet(), tx)
.allowStartIfComplete(true) // 재시작마다 다시 실행
.build();
}
결과
INFO 43158 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [prepareStep]
INFO 43158 --- [ main] c.e.batch.step11.PrepareTasklet : staging 테이블 초기화 완료
INFO 43158 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
⚠️ 함정 — allowStartIfComplete(true) 를 정산 Step 에 걸면 정산이 두 배가 됩니다
"재시작이 잘 안 돼서" 이 옵션을 붙이는 경우가 있습니다. 그러면 이미 커밋된 4만 건을 포함해 전부 다시 처리합니다.
settlement.uk_settlement_order UNIQUE 제약이 있으면 DuplicateKeyException 으로 시끄럽게 실패하니 그나마 다행이고(프로젝트 셋업 문서의 그 제약입니다), upsert 라면 조용히 넘어가지만 부작용은 두 번 발사됩니다.
이 옵션은 멱등한 준비/정리 Step 에만 씁니다.
startLimit(n)
같은 JobInstance 안에서 이 Step 을 최대 몇 번 시도할지 제한합니다. 기본값은 Integer.MAX_VALUE 입니다.
세 번째 재시작에서:
결과
ERROR 43205 --- [ main] o.s.batch.core.job.AbstractJob : Encountered fatal error executing job
org.springframework.batch.core.StartLimitExceededException: Maximum start limit exceeded for step: settlementStepStartMax: 2
at org.springframework.batch.core.job.SimpleStepHandler.handleStep(SimpleStepHandler.java:158) ~[spring-batch-core-5.1.1.jar:5.1.1]
at org.springframework.batch.core.job.SimpleJob.handleStep(SimpleJob.java:148) ~[spring-batch-core-5.1.1.jar:5.1.1]
...
INFO 43205 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [FAILED] in 112ms
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -t -e "
SELECT je.job_execution_id, je.status, COUNT(se.step_execution_id) AS steps
FROM BATCH_JOB_EXECUTION je LEFT JOIN BATCH_STEP_EXECUTION se USING (job_execution_id)
GROUP BY je.job_execution_id ORDER BY je.job_execution_id DESC LIMIT 3;"
결과
+------------------+--------+-------+
| job_execution_id | status | steps |
+------------------+--------+-------+
| 12 | FAILED | 0 |
| 11 | FAILED | 1 |
| 10 | FAILED | 1 |
+------------------+--------+-------+
세 번째 실행은 Step 을 하나도 만들지 못하고 죽었습니다. 무한 재시작 루프를 막는 안전장치입니다.
💡 실무 팁 — 스케줄러가 자동 재시도한다면 startLimit 을 꼭 거세요
Quartz/Airflow 가 실패 시 자동 재실행하도록 설정돼 있으면, 코드 버그로 항상 실패하는 Step 이 5분마다 4만 건을 다시 처리하며 DB 를 갈아 댑니다.
startLimit(3) 정도를 두면 세 번 만에 멈추고 사람에게 넘어옵니다.
정리
| 개념 | 핵심 |
|---|
.faultTolerant() | 청크 처리기를 FaultTolerantChunkProcessor 로 교체. 성능 특성이 통째로 바뀜 |
.skip(E) | E 와 모든 하위 타입을 건너뜀 |
.skipLimit(n) | 생략 시 기본 10. 무제한 아님 |
.noSkip(E) | skip 대상에서 제외. 가장 구체적인 등록이 이김(순서 무관) |
.skipPolicy(p) | 쓰는 순간 .skip() / .skipLimit() 은 조용히 무시됨 |
shouldSkip(t, skipCount) | skipCount 는 Step 전체 누적 skip 수 |
| skip 의 비용(처리) | 청크 롤백 + 청크 전체 재처리 (약 1.4배) |
| skip 의 비용(쓰기) | 청크 롤백 + 아이템 1건 = 트랜잭션 1개 스캔 모드 |
| 실측 | 정상 6.108초 → 쓰기 skip 100건 34.712초 (약 5.7배), commit_count 70 → 69,900 |
| 비용 결정 요인 | skip 건수가 아니라 몇 개 청크에 퍼졌는가. 1청크에 몰리면 6.594초 |
| 최선의 대응 | Reader WHERE 절로 미리 제외 → 6.233초 (약 5.6배 개선) |
.retry(E).retryLimit(n) | n 은 총 시도 횟수(최초 시도 포함). 재시도는 n-1 번 |
| retryLimit 기본값 | 0 — .retry() 만 쓰면 재시도가 일어나지 않음 |
backOffPolicy | 생략 시 NoBackOffPolicy(즉시 재시도). 데드락엔 ExponentialBackOffPolicy |
| retry + skip 같은 예외 | retry 소진 → skip. 성공 불가능한 예외에 걸면 102.337초(약 17배) |
.noRollback(E) | 트랜잭션을 롤백시키지 않음. 애플리케이션 검증 예외에만. DB 예외에는 금지 |
| 재시작 | 식별 파라미터가 같으면 같은 JobInstance → ExecutionContext 로 이어받음 |
| 재시작 검증 | BATCH_STEP_EXECUTION.read_count 가 40,000 → 30,000 이면 성공 |
| 상태 없는 Reader | ItemStream 미구현 시 처음부터 다시 읽음. 에러 없음, 2.4배 느림, 부작용 두 번 |
processorTransactional | 기본 true → 롤백 시 Processor 재호출. 부작용 있으면 알림 139,251회 |
.allowStartIfComplete(true) | 멱등한 준비/정리 Step 전용. 정산 Step 에 걸면 이중 정산 |
.startLimit(n) | 같은 JobInstance 내 최대 시도 횟수. 초과 시 StartLimitExceededException |
연습문제
Exercise.java 에 6문제가 있습니다. 정답은 Solution.java. 각 문제는 반드시 실행해 BATCH_STEP_EXECUTION 카운터로 검증하세요.
.skip() 만 걸고 .skipLimit() 을 빠뜨린 Step 이 몇 건째에서 죽는지 예측하고 확인하기
- 인프라 예외는 절대 건너뛰지 않고 데이터 예외만 300건까지 건너뛰는
SkipPolicy 구현하기
- 쓰기 skip 100건을 처리 skip 으로 옮겨
commit_count 를 69,900 → 70 으로 되돌리기
.retry(DeadlockLoserDataAccessException.class) 에 지수 backoff 를 붙이고 총 시도 횟수를 4회로 맞추기
- 40,001번째에서 실패시킨 뒤 재시작해
read_count 가 30,000 인지 검증하기
- 상태를 저장하지 않는 Reader 를
ItemStreamReader 로 고쳐 재시작이 이어지게 만들기
다음 단계
skip 과 retry 를 걸었지만, 지금은 누가 왜 버려졌는지 WARN 로그로만 남습니다. 운영에서는 "어제 정산에서 버려진 주문 100건의 목록"을 파일이나 테이블로 남겨야 하고, 청크마다 진행률을 찍어야 하며, Job 이 끝나면 결과를 알려야 합니다.
다음 스텝에서는 Job/Step/Chunk/Item/Skip/Retry 의 모든 지점에 후크를 거는 리스너를 다룹니다. 특히 SkipListener 가 트랜잭션 커밋 이후에 호출된다는 점과, 리스너에서 예외를 던지면 리스너 종류마다 결과가 다르다는 점이 이 스텝의 내용과 직접 이어집니다.
→ Step 12 — 리스너
실습 파일
이 스텝은 Java 파일 세 개로 진행합니다. Practice.java 를 위에서부터 따라가며 11-1 ~ 11-12 의 모든 측정을 재현하고, Exercise.java 의 6문제를 직접 채운 뒤, Solution.java 로 대조합니다. 세 파일 모두 com.example.batch.step11 패키지이며, 하나의 파일 안에 static class 로 설정 클래스들을 중첩해 두었습니다. 실행은 --spring.batch.job.name= 으로 원하는 Job 만 골라 돌립니다.
Practice.java
본문의 모든 예제를 절 번호 주석과 함께 담은 실습 파일입니다.
[11-0] 의 PlantBadData 는 불량 데이터를 심고 되돌리는 SQL 을 상수 문자열로 갖고 있습니다. main 에서 직접 실행하지 말고, 문서의 mysql 명령과 동일한 내용임을 확인하는 용도로 보세요. 실습 순서상 이것을 먼저 실행하지 않으면 11-2 이후의 skip 이 0건으로 나옵니다.
[11-1] → [11-2] 는 같은 Step 을 .faultTolerant() 유무로만 다르게 만든 쌍입니다. noToleranceStep 과 skipStep 을 순서대로 돌려 FAILED / COMPLETED 를 대조하세요.
[11-5] 의 ScanModeDemo 가 이 스텝의 핵심 측정입니다. baselineJob(정상) → writeSkipJob(쓰기 skip 100건) → clusteredSkipJob(1청크에 몰린 skip 100건) 을 이 순서대로 돌려야 6.108 / 34.712 / 6.594 세 숫자가 나옵니다. 중간에 TRUNCATE TABLE settlement 를 빠뜨리면 앞 실행의 결과가 남아 skip 수가 달라집니다.
[11-7] 의 FlakyWriter 는 AtomicInteger 로 "세 번째 청크의 처음 두 번의 시도에서만 데드락"을 흉내 냅니다. 진짜 MySQL 데드락을 재현하려면 커넥션 두 개로 교차 갱신을 해야 하는데, 재시도 로그를 보는 것이 목적이므로 예외를 직접 던집니다. 던지는 예외 타입(DeadlockLoserDataAccessException)은 실제와 동일합니다.
[11-10] 의 failAt 파라미터는 new JobParameter<>(value, Long.class, **false**) 로 선언돼 있습니다. identifying=false 여야 재실행 시 같은 JobInstance 로 붙어 재시작이 됩니다. 이 false 를 true 로 바꾸면 재시작이 아니라 새 인스턴스가 만들어져 실습이 성립하지 않습니다.
[11-11] 은 NaiveReader(상태 없음) 와 StatefulReader(ItemStreamReader 구현) 두 개를 나란히 두었습니다. 같은 Job 을 Reader 만 바꿔 두 번 돌려 read_count 30,000 vs 70,000 을 비교하는 것이 목적입니다.
- 파일 맨 아래
CleanUp 에 되돌리기 SQL 과 메타데이터 초기화 스크립트를 주석으로 모아 두었습니다.
package com.example.batch.step11;
import com.example.batch.domain.Order;
import com.example.batch.domain.Settlement;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.configuration.annotation.StepScope;
import org.springframework.batch.core.job.builder.JobBuilder;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.core.step.skip.SkipLimitExceededException;
import org.springframework.batch.core.step.skip.SkipPolicy;
import org.springframework.batch.item.Chunk;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemStreamException;
import org.springframework.batch.item.ItemStreamReader;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.database.JdbcBatchItemWriter;
import org.springframework.batch.item.database.JdbcPagingItemReader;
import org.springframework.batch.item.database.builder.JdbcBatchItemWriterBuilder;
import org.springframework.batch.item.database.builder.JdbcPagingItemReaderBuilder;
import org.springframework.batch.item.database.support.MySqlPagingQueryProvider;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.dao.DataAccessException;
import org.springframework.dao.DataIntegrityViolationException;
import org.springframework.dao.DeadlockLoserDataAccessException;
import org.springframework.jdbc.core.DataClassRowMapper;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.namedparam.MapSqlParameterSource;
import org.springframework.retry.backoff.BackOffPolicy;
import org.springframework.retry.backoff.ExponentialBackOffPolicy;
import org.springframework.transaction.PlatformTransactionManager;
import javax.sql.DataSource;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.LocalDate;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
/**
* Step 11 — 내결함성: skip · retry · 재시작 : 본문 11-0 ~ 11-12 의 모든 예제.
*
* ─────────────────────────────────────────────────────────────────────────
* 실행 방법
*
* ./gradlew bootRun --args='--spring.batch.job.name=settlementSkipJob date=2025-03-01'
*
* 바깥 클래스에 @Configuration 이 없습니다. 지금 실습할 static class 하나만
* @Configuration 주석을 풀고 나머지는 주석 처리한 채로 돌리십시오.
*
* ─────────────────────────────────────────────────────────────────────────
* ⚠️ 실습 순서를 반드시 지키십시오.
*
* 1. [11-0] PlantBadData.PLANT_SQL 을 먼저 실행합니다.
* 이것을 빼먹으면 11-2 이후의 skip 이 전부 0건으로 나와서
* 실습이 통째로 성립하지 않습니다.
*
* 2. 각 측정 전에 TRUNCATE TABLE settlement 를 합니다.
* 빠뜨리면 앞 실행 결과가 남아 skip 수가 달라집니다.
*
* 3. 실습이 끝나면 CleanUp.REVERT_SQL 로 orders 를 원상 복구합니다.
* orders 는 이후 모든 스텝이 공유하는 공용 테이블입니다.
* ─────────────────────────────────────────────────────────────────────────
*/
public class Practice {
// =====================================================================
// [11-0] 불량 데이터 심기
//
// main 에서 실행하지 마십시오. 문서의 mysql 명령과 같은 내용임을
// 확인하는 용도입니다. mysql 클라이언트로 직접 실행하십시오.
// =====================================================================
public static class PlantBadData {
/**
* (A) Processor 에서 터질 불량: 금액이 음수인 주문 100건.
* (B) Writer 에서 터질 중복: settlement 에 미리 100행.
*
* order_id % 1000 = 0 인 주문은 order_id % 10 = 0 도 만족하므로
* 100건 전부가 status='COMPLETED' 입니다.
*
* ⚠️ 핵심: 불량 100건이 70청크에 "고르게 흩어져" 있습니다.
* COMPLETED 700건마다 1건이므로, 청크 1,000 이면
* 70청크 전부가 최소 1건의 불량을 품습니다.
* 11-5 에서 이 배치가 왜 5.7배나 느려지는지의 열쇠입니다.
*/
public static final String PLANT_SQL = """
UPDATE orders SET amount = -1.00 WHERE order_id %% 1000 = 0;
TRUNCATE TABLE settlement;
INSERT INTO settlement
(order_id, customer_id, settle_date, gross_amount,
fee_rate, fee_amount, net_amount)
SELECT o.order_id, o.customer_id, DATE(o.ordered_at),
1000.00, 0.0300, 30.00, 970.00
FROM orders o
WHERE o.status = 'COMPLETED' AND o.order_id %% 1000 = 0;
-- 검증: target=70000, bad_amount=100, preloaded=100
SELECT (SELECT COUNT(*) FROM orders WHERE status='COMPLETED') AS target,
(SELECT COUNT(*) FROM orders WHERE amount < 0) AS bad_amount,
(SELECT COUNT(*) FROM settlement) AS preloaded;
""";
}
// =====================================================================
// 공통 — Reader / Processor / Writer
// =====================================================================
/** 날짜 파라미터 없이 COMPLETED 전체(70,000)를 읽습니다. */
public static JdbcPagingItemReader<Order> buildOrderReader(DataSource dataSource) {
MySqlPagingQueryProvider provider = new MySqlPagingQueryProvider();
provider.setSelectClause("order_id, customer_id, amount, status, ordered_at");
provider.setFromClause("FROM orders");
provider.setWhereClause("WHERE status = 'COMPLETED'");
// 정렬 키가 없으면 페이지가 밀려 데이터가 유실됩니다 — Step 06 참고.
provider.setSortKeys(Map.of(
"order_id", org.springframework.batch.item.database.Order.ASCENDING));
return new JdbcPagingItemReaderBuilder<Order>()
.name("orderReader")
.dataSource(dataSource)
.queryProvider(provider)
.pageSize(1000)
.rowMapper(new DataClassRowMapper<>(Order.class))
.build();
}
/**
* 등급별 수수료 적용. 금액이 음수면 예외를 던집니다.
*
* 11-0 에서 심은 100건의 음수 금액이 여기서 터집니다.
*/
public static class SettlementProcessor implements ItemProcessor<Order, Settlement> {
private static final BigDecimal[] FEE_RATES = {
new BigDecimal("0.0350"), // BRONZE
new BigDecimal("0.0300"), // SILVER
new BigDecimal("0.0250"), // GOLD
new BigDecimal("0.0200") // VIP
};
@Override
public Settlement process(Order order) {
if (order.amount().signum() < 0) {
throw new IllegalArgumentException(
"정산 금액이 음수입니다: order_id=" + order.order_id()
+ ", amount=" + order.amount());
}
BigDecimal feeRate = FEE_RATES[order.customerId() % 4];
BigDecimal gross = order.amount();
BigDecimal fee = gross.multiply(feeRate).setScale(2, RoundingMode.HALF_UP);
return new Settlement(order.order_id(), order.customerId(),
order.orderedAt().toLocalDate(),
gross, feeRate, fee, gross.subtract(fee));
}
}
/**
* 비멱등 Writer — 순수 INSERT.
*
* ⚠️ ON DUPLICATE KEY UPDATE 를 일부러 쓰지 않습니다.
* UNIQUE 충돌이 DuplicateKeyException 으로 터져야 skip 실습이
* 가능하기 때문입니다. Step 13 의 멱등 Writer 와 대조하십시오.
*/
public static JdbcBatchItemWriter<Settlement> buildStrictWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<Settlement>()
.dataSource(dataSource)
.sql("""
INSERT INTO settlement
(order_id, customer_id, settle_date, gross_amount,
fee_rate, fee_amount, net_amount)
VALUES
(:orderId, :customerId, :settleDate, :grossAmount,
:feeRate, :feeAmount, :netAmount)
""")
.itemSqlParameterSourceProvider(s -> new MapSqlParameterSource()
.addValue("orderId", s.orderId())
.addValue("customerId", s.customerId())
.addValue("settleDate", s.settleDate())
.addValue("grossAmount", s.grossAmount())
.addValue("feeRate", s.feeRate())
.addValue("feeAmount", s.feeAmount())
.addValue("netAmount", s.netAmount()))
.build();
}
// =====================================================================
// [11-1] 아무 대비 없는 Step — 첫 예외에서 통째로 멈춘다
//
// 결과: FAILED. read_count 는 1000 근처에서 멈추고
// settlement 에는 아무것도 안 남습니다(청크가 롤백되므로).
// =====================================================================
// @Configuration
public static class NoTolerance {
@Bean
public Job settlementNoToleranceJob(JobRepository jobRepository, Step noToleranceStep) {
return new JobBuilder("settlementNoToleranceJob", jobRepository)
.start(noToleranceStep)
.build();
}
@Bean
public Step noToleranceStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("noToleranceStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(new SettlementProcessor())
.writer(buildStrictWriter(dataSource))
// .faultTolerant() 가 없습니다. 첫 예외에서 Step 전체가 죽습니다.
.build();
}
}
// =====================================================================
// [11-2] .faultTolerant() 와 .skip()
//
// 11-1 과 완전히 같고 세 줄만 추가했습니다. 이 쌍을 순서대로 돌려
// FAILED / COMPLETED 를 대조하십시오.
// =====================================================================
// @Configuration
public static class SkipBasics {
@Bean
public Job settlementSkipJob(JobRepository jobRepository, Step skipStep) {
return new JobBuilder("settlementSkipJob", jobRepository)
.start(skipStep)
.build();
}
@Bean
public Step skipStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("skipStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(new SettlementProcessor())
.writer(buildStrictWriter(dataSource))
.faultTolerant()
.skip(IllegalArgumentException.class)
.skipLimit(200) // 기본값은 10 입니다 — 11-3
.build();
}
}
// =====================================================================
// [11-4] 커스텀 SkipPolicy — 조건을 코드로 쓴다
//
// ⚠️ instanceof 판정 순서가 핵심입니다.
// 구체적인 예외를 먼저, 넓은 예외를 나중에 검사해야 합니다.
// 순서를 뒤집으면 인프라 장애가 데이터 오류로 오인되어 조용히
// skip 됩니다 — 연습문제 2 의 함정입니다.
// =====================================================================
public static class SettlementSkipPolicy implements SkipPolicy {
private final int dataErrorLimit;
public SettlementSkipPolicy(int dataErrorLimit) {
this.dataErrorLimit = dataErrorLimit;
}
@Override
public boolean shouldSkip(Throwable t, long skipCount)
throws SkipLimitExceededException {
// ① 인프라 장애는 절대 skip 하지 않습니다.
// 커넥션이 끊긴 상황에서 계속 진행하면 나머지 전부를
// "실패로 건너뛰고" COMPLETED 로 끝냅니다.
if (t instanceof org.springframework.dao.DataAccessResourceFailureException) {
return false;
}
// ② 데이터 오류는 한도까지 skip 합니다.
if (t instanceof IllegalArgumentException
|| t instanceof DataIntegrityViolationException) {
if (skipCount >= dataErrorLimit) {
throw new SkipLimitExceededException(dataErrorLimit, t);
}
return true;
}
// ③ 그 외 예외는 모릅니다 → 안전하게 실패시킵니다.
// "모르는 예외는 skip 하지 않는다"가 안전한 기본값입니다.
return false;
}
}
// =====================================================================
// [11-5] skip 의 진짜 비용 — 청크 롤백과 스캔 모드
//
// ★ 이 스텝의 핵심 측정입니다.
//
// 아래 세 Job 을 이 순서대로 돌려야 6.108 / 34.712 / 6.594 가 나옵니다.
// 각 실행 전에 반드시 TRUNCATE TABLE settlement 하십시오.
//
// ① baselineJob — 불량 없음 → 6.108초, commit 70, rollback 0
// ② writeSkipJob — 흩어진 skip 100 → 34.712초, commit 69900, rollback 100
// ③ clusteredSkipJob — 뭉친 skip 100 → 6.594초, commit 1069, rollback 1
//
// ②가 5.7배 느린 이유:
// skip 이 나면 그 청크는 통째로 롤백되고, Spring Batch 는 "누가
// 범인인지" 모르므로 **아이템을 1건씩 다시 처리**합니다(스캔 모드).
// 1,000건짜리 청크가 1,000개의 트랜잭션으로 쪼개집니다.
// 불량이 70청크 전부에 흩어져 있으므로 70청크 × 1000 = 70,000번의
// 커밋에 가까워집니다. commit_count 69,900 이 그 증거입니다.
//
// ③이 빠른 이유:
// 불량 100건이 1개 청크에 몰려 있으면 그 청크 하나만 스캔 모드로
// 들어갑니다. 69청크는 정상 속도로 통과합니다.
// → **skip 비용은 skip 건수가 아니라 "오염된 청크 수"에 비례합니다.**
// =====================================================================
// @Configuration
public static class ScanModeDemo {
@Bean
public Job settlementWriteSkipJob(JobRepository jobRepository, Step writeSkipStep) {
return new JobBuilder("settlementWriteSkipJob", jobRepository)
.start(writeSkipStep)
.build();
}
@Bean
public Step writeSkipStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("writeSkipStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(new SettlementProcessor())
.writer(buildStrictWriter(dataSource))
.faultTolerant()
// DuplicateKeyException 의 상위 타입입니다.
.skip(DataIntegrityViolationException.class)
.skipLimit(200)
.build();
}
/**
* ③ 뭉친 skip — 불량을 1개 청크에 몰아 넣은 상태에서 돌립니다.
*
* 사전 SQL (order_id 1~100 에 몰기):
* TRUNCATE TABLE settlement;
* INSERT INTO settlement (order_id, customer_id, settle_date,
* gross_amount, fee_rate, fee_amount, net_amount)
* SELECT order_id, customer_id, DATE(ordered_at), 1000.00, 0.0300, 30.00, 970.00
* FROM orders WHERE status='COMPLETED' AND order_id <= 1000;
*/
@Bean
public Job settlementClusteredSkipJob(JobRepository jobRepository, Step writeSkipStep) {
return new JobBuilder("settlementClusteredSkipJob", jobRepository)
.start(writeSkipStep)
.build();
}
}
// =====================================================================
// [11-7] .retry() 와 backoff
// =====================================================================
/**
* 지수 백오프. 200 → 400 → 800 ms.
*
* ⚠️ maxInterval 을 반드시 거십시오. 없으면 재시도가 많은 Step 에서
* 대기 시간이 폭주합니다(200 → 400 → ... → 수 분).
*/
public static BackOffPolicy exponentialBackOff() {
ExponentialBackOffPolicy policy = new ExponentialBackOffPolicy();
policy.setInitialInterval(200L);
policy.setMultiplier(2.0);
policy.setMaxInterval(5_000L);
return policy;
}
/**
* 세 번째 청크의 처음 두 번의 시도에서만 데드락을 던지는 Writer.
*
* 진짜 MySQL 데드락을 재현하려면 커넥션 두 개로 교차 갱신을 해야 하지만,
* 여기서는 재시도 로그를 보는 것이 목적이므로 예외를 직접 던집니다.
* 던지는 예외 타입은 실제 Spring 이 변환해 주는 것과 동일합니다.
*/
public static class FlakyWriter implements ItemWriter<Settlement> {
private final ItemWriter<Settlement> delegate;
private final AtomicInteger chunkCounter = new AtomicInteger(0);
private final AtomicInteger failuresLeft = new AtomicInteger(2);
public FlakyWriter(ItemWriter<Settlement> delegate) {
this.delegate = delegate;
}
@Override
public void write(Chunk<? extends Settlement> chunk) throws Exception {
int current = chunkCounter.incrementAndGet();
if (current == 3 && failuresLeft.getAndDecrement() > 0) {
throw new DeadlockLoserDataAccessException(
"Deadlock found when trying to get lock; try restarting transaction",
new java.sql.SQLException("Deadlock found", "40001", 1213));
}
delegate.write(chunk);
}
}
// @Configuration
public static class RetryDemo {
@Bean
public Job settlementRetryJob(JobRepository jobRepository, Step retryStep) {
return new JobBuilder("settlementRetryJob", jobRepository)
.start(retryStep)
.build();
}
@Bean
public Step retryStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("retryStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(new SettlementProcessor())
.writer(new FlakyWriter(buildStrictWriter(dataSource)))
.faultTolerant()
.retry(DeadlockLoserDataAccessException.class)
.retry(org.springframework.dao.CannotAcquireLockException.class)
// ⚠️ retryLimit(3) 은 "재시도 3번"이 아니라 "총 시도 3회"입니다.
// 그리고 .retry() 만 쓰고 이 줄을 빠뜨리면 기본값 0이라
// 아무 재시도도 일어나지 않습니다. 항상 짝으로 쓰십시오.
.retryLimit(3)
.backOffPolicy(exponentialBackOff())
.skip(DataIntegrityViolationException.class)
.skipLimit(200)
.build();
}
}
// =====================================================================
// [11-9] noRollback — 롤백시키지 않을 예외
//
// 기본적으로 어떤 예외든 청크 트랜잭션을 롤백시킵니다.
// "이 예외는 데이터에 영향이 없으니 롤백할 필요 없다"고 선언하면
// 스캔 모드로 들어가지 않아 성능이 유지됩니다.
// =====================================================================
// @Configuration
public static class NoRollbackDemo {
@Bean
public Step noRollbackStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("noRollbackStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(new SettlementProcessor())
.writer(buildStrictWriter(dataSource))
.faultTolerant()
.skip(ValidationWarning.class)
.skipLimit(1000)
// 이 예외는 DB 를 건드리기 전에 나므로 롤백이 불필요합니다.
.noRollback(ValidationWarning.class)
.build();
}
}
/** 데이터에 영향을 주지 않는 검증 경고. */
public static class ValidationWarning extends RuntimeException {
public ValidationWarning(String message) {
super(message);
}
}
// =====================================================================
// [11-10] 재시작 — 중단 지점을 이어받는다
//
// failAt 번째 아이템에서 일부러 죽였다가, 같은 파라미터로 다시 실행해
// 중단 지점부터 이어가는 것을 확인합니다.
//
// ⚠️ failAt 은 identifying = false 여야 합니다.
// true 면 파라미터가 달라져 "새 JobInstance" 가 만들어지고,
// 재시작이 아니라 처음부터 새로 도는 것이 됩니다.
// 이 false 하나가 실습의 성립 조건입니다.
// =====================================================================
// @Configuration
public static class RestartDemo {
@Bean
public Job settlementRestartJob(JobRepository jobRepository, Step restartStep) {
return new JobBuilder("settlementRestartJob", jobRepository)
.start(restartStep)
.build();
}
@Bean
public Step restartStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource,
ItemProcessor<Order, Settlement> failingProcessor) {
return new StepBuilder("restartStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(failingProcessor)
.writer(buildStrictWriter(dataSource))
.build();
}
/**
* failAt 번째 아이템에서 죽습니다.
*
* 실행 예:
* 1회차: --spring.batch.job.name=settlementRestartJob failAt=30000
* → 30,000번째에서 FAILED. commit 29 (29,000건 커밋됨)
* 2회차: --spring.batch.job.name=settlementRestartJob failAt=999999
* → 29,000번째부터 재개. read_count 41,000
*
* ⚠️ 2회차에서 failAt 값을 바꿔도 같은 JobInstance 입니다.
* identifying=false 이기 때문입니다.
*/
@Bean
@StepScope
public ItemProcessor<Order, Settlement> failingProcessor(
@Value("#{jobParameters['failAt'] ?: 999999L}") Long failAt) {
SettlementProcessor delegate = new SettlementProcessor();
AtomicInteger counter = new AtomicInteger(0);
return order -> {
if (counter.incrementAndGet() >= failAt) {
throw new IllegalStateException(
"의도적 실패: " + failAt + "번째 아이템");
}
return delegate.process(order);
};
}
}
// =====================================================================
// [11-11] 함정 — Reader 가 상태를 저장하지 않으면
//
// 같은 Job 을 Reader 만 바꿔 두 번 돌려
// read_count 30,000 vs 70,000 을 비교하는 것이 목적입니다.
//
// NaiveReader → 재시작해도 처음부터 다시 읽습니다.
// 이미 처리한 29,000건을 또 처리하려다
// UNIQUE 충돌로 죽거나(비멱등 Writer),
// 조용히 중복 정산합니다(멱등 Writer).
// StatefulReader → ExecutionContext 에 위치를 저장해 이어받습니다.
// =====================================================================
/**
* ⚠️ 일부러 틀린 Reader. ItemStreamReader 가 아니라 ItemReader 만
* 구현했으므로 open/update/close 콜백이 아예 없습니다.
* Spring Batch 는 이 Reader 의 위치를 저장할 방법이 없습니다.
*/
public static class NaiveReader implements
org.springframework.batch.item.ItemReader<Order> {
private final List<Order> orders;
private int index = 0;
public NaiveReader(List<Order> orders) {
this.orders = orders;
}
@Override
public Order read() {
return index >= orders.size() ? null : orders.get(index++);
}
}
/**
* 상태를 저장하는 Reader.
*
* 세 콜백의 역할:
* open(ctx) — Step 시작 시 1회. 재시작이면 저장된 위치를 복원합니다.
* update(ctx) — 청크 커밋마다. 현재 위치를 기록합니다.
* close() — Step 종료 시 1회. 자원 정리.
*
* ⚠️ ExecutionContext 의 키에 Reader 이름을 접두사로 붙이는 것이 중요합니다.
* 안 붙이면 한 Step 에 Reader 가 둘일 때 서로의 키를 덮어씁니다.
*/
public static class StatefulReader implements ItemStreamReader<Order> {
private static final String KEY_INDEX = "statefulReader.index";
private final List<Order> orders;
private int index = 0;
public StatefulReader(List<Order> orders) {
this.orders = orders;
}
@Override
public void open(ExecutionContext ctx) throws ItemStreamException {
// 재시작 판정의 관용구입니다.
if (ctx.containsKey(KEY_INDEX)) {
this.index = ctx.getInt(KEY_INDEX);
System.out.println(">>> 재시작: " + index + "번째부터 이어갑니다");
} else {
this.index = 0;
}
}
@Override
public Order read() {
return index >= orders.size() ? null : orders.get(index++);
}
@Override
public void update(ExecutionContext ctx) throws ItemStreamException {
ctx.putInt(KEY_INDEX, index);
}
@Override
public void close() throws ItemStreamException {
// 정리할 자원 없음
}
}
/** 70,000건을 메모리에 올립니다. 실습 전용입니다. */
public static List<Order> loadCompletedOrders(DataSource dataSource) {
return new JdbcTemplate(dataSource).query("""
SELECT order_id, customer_id, amount, status, ordered_at
FROM orders WHERE status = 'COMPLETED' ORDER BY order_id
""", new DataClassRowMapper<>(Order.class));
}
// =====================================================================
// [11-12] allowStartIfComplete 와 startLimit
// =====================================================================
// @Configuration
public static class RestartPolicyDemo {
/**
* allowStartIfComplete(true)
* — 이미 COMPLETED 인 Step 도 재실행합니다.
* "매번 처음부터 다시 해야 하는" 정리/초기화 Step 에 씁니다.
*
* ⚠️ 정산 Step 에는 절대 붙이지 마십시오.
* 재시작할 때마다 처음부터 다시 정산해서 중복이 쌓입니다.
*/
@Bean
public Step cleanupStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("cleanupStep", jobRepository)
.tasklet((contribution, chunkContext) -> {
new JdbcTemplate(dataSource).update("TRUNCATE TABLE settlement");
return org.springframework.batch.repeat.RepeatStatus.FINISHED;
}, txManager)
.allowStartIfComplete(true)
.build();
}
/**
* startLimit(n) — 이 Step 을 최대 n 번까지만 시도합니다.
*
* 기본값은 Integer.MAX_VALUE 입니다. 제한을 두면 "고칠 수 없는
* 실패를 무한히 재시도하는" 상황을 막을 수 있습니다.
* n 번을 넘기면 StartLimitExceededException 이 납니다.
*/
@Bean
public Step limitedStep(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("limitedStep", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(buildOrderReader(dataSource))
.processor(new SettlementProcessor())
.writer(buildStrictWriter(dataSource))
.startLimit(3)
.build();
}
}
// =====================================================================
// CleanUp — 실습 뒷정리
// =====================================================================
public static class CleanUp {
/**
* orders 를 원래 금액으로 되돌립니다.
*
* ⚠️ orders 는 이후 모든 스텝이 공유하는 공용 테이블입니다.
* 이 스텝을 끝내면 반드시 실행하십시오. 음수 금액을 남겨 두면
* Step 12·13·14 의 결과가 교재와 달라집니다.
*/
public static final String REVERT_SQL = """
UPDATE orders o
JOIN (SELECT order_id, 1000 + (order_id %% 977) * 100 AS amt
FROM orders) s
ON o.order_id = s.order_id
SET o.amount = s.amt
WHERE o.amount < 0;
TRUNCATE TABLE settlement;
-- 검증: bad_amount 가 0 이어야 합니다.
SELECT COUNT(*) AS bad_amount FROM orders WHERE amount < 0;
""";
public static final String RESET_METADATA_SQL = """
SET FOREIGN_KEY_CHECKS = 0;
DELETE FROM BATCH_STEP_EXECUTION_CONTEXT;
DELETE FROM BATCH_STEP_EXECUTION;
DELETE FROM BATCH_JOB_EXECUTION_CONTEXT;
DELETE FROM BATCH_JOB_EXECUTION_PARAMS;
DELETE FROM BATCH_JOB_EXECUTION;
DELETE FROM BATCH_JOB_INSTANCE;
SET FOREIGN_KEY_CHECKS = 1;
""";
/** 카운터 확인용. 모든 측정 뒤에 이 쿼리를 돌리십시오. */
public static final String VERIFY_SQL = """
SELECT STEP_NAME, STATUS,
READ_COUNT, WRITE_COUNT, FILTER_COUNT,
READ_SKIP_COUNT, PROCESS_SKIP_COUNT, WRITE_SKIP_COUNT,
COMMIT_COUNT, ROLLBACK_COUNT,
TIMESTAMPDIFF(SECOND, START_TIME, END_TIME) AS secs
FROM BATCH_STEP_EXECUTION
ORDER BY STEP_EXECUTION_ID DESC
LIMIT 5;
""";
}
}
Exercise.java
6문제의 문제지입니다. 각 문제는 // 여기에 작성: 자리를 비워 두었습니다.
- 문제 1·4 는 빌더 체인의 한 줄을 채우는 문제이고, 문제 2·6 은 클래스를 직접 구현하는 문제, 문제 3·5 는 설계를 바꿔 카운터를 목표값으로 만드는 문제입니다.
- 문제 2 의
SkipPolicy 는 인프라 예외를 먼저 걸러내는 순서가 핵심입니다. instanceof 판정 순서를 바꾸면 DataAccessResourceFailureException 이 DataAccessException 분기에 먼저 걸려 조용히 skip 됩니다. 이 순서 실수가 문제 2 의 진짜 함정입니다.
- 문제 3 은 "쓰기에서 터질 조건을 처리 단계에서 미리 판정"하는 문제입니다. 정답 코드가 짧아서 쉬워 보이지만, 검증 쿼리를 아이템마다 날리면 스캔 모드보다 더 느려집니다. 청크 단위로 한 번에 조회하거나 Reader 쿼리에서 거르는 쪽으로 유도됩니다.
- 문제 5·6 은 연속 실습입니다. 5번에서 실패시킨 JobInstance 를 6번에서 그대로 재시작하므로, 5번을 풀고 나서 메타데이터를 초기화하면 6번을 풀 수 없습니다. 파일 상단 주석의 순서를 지키세요.
- 각 문제 끝에
-- 검증: 주석으로 확인용 SQL 이 붙어 있습니다. 코드를 고친 뒤 반드시 이 쿼리로 카운터를 확인하세요. "에러 없이 끝났다"는 정답의 근거가 되지 못합니다.
package com.example.batch.step11;
import com.example.batch.domain.Order;
import com.example.batch.domain.Settlement;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.core.step.skip.SkipLimitExceededException;
import org.springframework.batch.core.step.skip.SkipPolicy;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemStreamException;
import org.springframework.batch.item.ItemStreamReader;
import org.springframework.transaction.PlatformTransactionManager;
import javax.sql.DataSource;
import java.util.List;
/**
* Step 11 — 연습문제 (6문제)
*
* 정답은 Solution.java. 먼저 직접 풀어 보십시오.
*
* ─────────────────────────────────────────────────────────────────────────
* ⚠️ 풀이 순서를 지키십시오.
*
* - 문제 1 을 풀기 전에 Practice.PlantBadData.PLANT_SQL 을 실행해
* 불량 데이터를 심어야 합니다. 안 그러면 skip 이 0건이라 문제가
* 성립하지 않습니다.
*
* - 문제 5 와 6 은 연속 실습입니다. 5번에서 실패시킨 JobInstance 를
* 6번에서 그대로 재시작합니다. 5번을 푼 뒤 메타데이터를 초기화하면
* 6번을 풀 수 없습니다.
*
* - 각 문제 끝의 `-- 검증:` SQL 을 반드시 돌리십시오.
* "에러 없이 끝났다"는 정답의 근거가 되지 못합니다.
* ─────────────────────────────────────────────────────────────────────────
*/
public class Exercise {
// =====================================================================
// 문제 1. skipLimit 기본값에 걸려 죽는 지점을 특정하기
//
// 아래 Step 은 .skipLimit() 을 명시하지 않았습니다.
//
// (a) 이 Step 을 그대로 실행하면 어떻게 됩니까? 예외 메시지를 적으십시오.
// (b) 몇 번째 불량 데이터에서 죽습니까?
// (c) 그 불량의 order_id 는 몇입니까?
// 힌트: 불량은 order_id % 1000 = 0 인 100건입니다.
// (d) 그 지점은 몇 번째 청크입니까?
// 힌트: COMPLETED 는 order_id % 10 <= 6 이므로, order_id N 까지의
// COMPLETED 건수는 대략 N * 0.7 입니다.
// (e) skipLimit 을 200 으로 바꾸고 다시 돌려 카운터를 비교하십시오.
// =====================================================================
public static Step problem1Step(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("problem1Step", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(Practice.buildOrderReader(dataSource))
.processor(new Practice.SettlementProcessor())
.writer(Practice.buildStrictWriter(dataSource))
.faultTolerant()
.skip(IllegalArgumentException.class)
// 여기에 작성: skipLimit 을 명시하십시오
//
.build();
}
// (a) 예외 메시지
// 여기에 작성:
//
// (b) 몇 번째 불량에서 죽는가
// 여기에 작성:
//
// (c) 그 불량의 order_id
// 여기에 작성:
//
// (d) 몇 번째 청크인가
// 여기에 작성:
//
// -- 검증:
// SELECT STEP_NAME, STATUS, READ_COUNT, PROCESS_SKIP_COUNT, COMMIT_COUNT
// FROM BATCH_STEP_EXECUTION ORDER BY STEP_EXECUTION_ID DESC LIMIT 1;
// =====================================================================
// 문제 2. SkipPolicy 를 직접 구현하기
//
// 다음 요구사항을 만족하는 SkipPolicy 를 작성하십시오.
//
// ① IllegalArgumentException (데이터 오류) → 최대 50건까지 skip
// ② DataIntegrityViolationException (중복) → 최대 50건까지 skip
// ③ DataAccessResourceFailureException (인프라) → 절대 skip 하지 않음
// ④ 그 외 모든 예외 → skip 하지 않음
//
// ⚠️ 이 문제의 진짜 함정은 instanceof 판정 "순서" 입니다.
// DataAccessResourceFailureException 은 DataAccessException 의
// 하위 타입입니다. 넓은 타입을 먼저 검사하면 인프라 장애가
// 데이터 오류로 오인되어 조용히 skip 됩니다.
//
// 그러면 DB 커넥션이 끊긴 상황에서 배치가 "나머지를 전부 skip 하고"
// COMPLETED 로 끝납니다. 정산이 안 됐는데 성공했다고 보고합니다.
//
// (a) SkipPolicy 를 구현하십시오.
// (b) 순서를 일부러 뒤집어 보고, 인프라 예외가 skip 되는 것을 확인하십시오.
// (c) 예외 "종류별로" 한도를 따로 두려면 무엇이 필요합니까?
// 힌트: shouldSkip 의 skipCount 파라미터는 Step 전체 누적입니다.
// =====================================================================
public static class Problem2SkipPolicy implements SkipPolicy {
@Override
public boolean shouldSkip(Throwable t, long skipCount)
throws SkipLimitExceededException {
// 여기에 작성: 판정 순서에 주의하십시오
//
return false;
}
}
// (c) 예외 종류별 한도를 두려면?
// 여기에 작성:
//
// -- 검증: 인프라 예외를 흉내 내려면 실행 중 컨테이너를 잠깐 멈추십시오.
// docker compose -f docker/docker-compose.yml pause mysql
// (5초 뒤) docker compose -f docker/docker-compose.yml unpause mysql
// =====================================================================
// 문제 3. 스캔 모드를 피해 commit_count 를 70 으로 만들기
//
// 11-5 의 writeSkipJob 은 commit_count 가 69,900 이고 34.712초 걸립니다.
// 쓰기 단계에서 UNIQUE 충돌이 나기 때문입니다.
//
// 이 배치를 **결과는 같으면서** commit_count 70, 소요 10초 이내로
// 바꾸십시오.
//
// (a) 핵심 아이디어를 한 문장으로 적으십시오.
// 힌트: 쓰기에서 터질 조건을 "미리" 알 수 있습니까?
// (b) 구현하십시오.
// (c) ⚠️ 함정: 아이템마다 검증 쿼리를 날리면 어떻게 됩니까?
// 70,000번의 SELECT 가 생깁니다. 스캔 모드보다 느려질 수 있습니다.
// 이것을 피하려면 어떻게 해야 합니까?
// (d) 더 나은 답이 있습니다. 애초에 Reader 가 그 행들을 안 읽게 하려면?
// =====================================================================
// (a) 핵심 아이디어
// 여기에 작성:
//
public static class Problem3Processor implements ItemProcessor<Order, Settlement> {
@Override
public Settlement process(Order order) {
// 여기에 작성:
//
return null;
}
}
// (c) 아이템마다 쿼리를 날리면? 어떻게 피하는가?
// 여기에 작성:
//
// (d) Reader 쿼리로 거르는 방법
// 여기에 작성:
//
// -- 검증: commit_count 가 70 이고 write_count + skip 이 70000 이어야 합니다.
// SELECT READ_COUNT, WRITE_COUNT, FILTER_COUNT, PROCESS_SKIP_COUNT,
// WRITE_SKIP_COUNT, COMMIT_COUNT,
// TIMESTAMPDIFF(SECOND, START_TIME, END_TIME) secs
// FROM BATCH_STEP_EXECUTION ORDER BY STEP_EXECUTION_ID DESC LIMIT 1;
// =====================================================================
// 문제 4. 재시도를 정확히 3번 하려면 retryLimit 을 얼마로?
//
// (a) "최초 시도 후 재시도를 3번 더" 하고 싶습니다. retryLimit 값은?
// (b) ExponentialBackOffPolicy(initial=200, multiplier=2.0) 일 때
// 최악의 경우 총 대기 시간은 몇 ms 입니까? 계산 과정을 적으십시오.
// (c) maxInterval 을 설정하지 않으면 어떤 문제가 생깁니까?
// (d) .retry() 는 썼는데 .retryLimit() 을 빠뜨리면 어떻게 됩니까?
// 에러가 납니까, 아니면 조용히 아무 일도 안 일어납니까?
// =====================================================================
public static Step problem4Step(JobRepository jobRepository,
PlatformTransactionManager txManager,
DataSource dataSource) {
return new StepBuilder("problem4Step", jobRepository)
.<Order, Settlement>chunk(1000, txManager)
.reader(Practice.buildOrderReader(dataSource))
.processor(new Practice.SettlementProcessor())
.writer(new Practice.FlakyWriter(Practice.buildStrictWriter(dataSource)))
.faultTolerant()
.retry(org.springframework.dao.DeadlockLoserDataAccessException.class)
// 여기에 작성: retryLimit
//
.backOffPolicy(Practice.exponentialBackOff())
.build();
}
// (a) retryLimit 값과 이유
// 여기에 작성:
//
// (b) 최악의 총 대기 시간 계산
// 여기에 작성:
//
// (c) maxInterval 미설정 시 문제
// 여기에 작성:
//
// (d) retryLimit 을 빠뜨리면
// 여기에 작성:
//
// =====================================================================
// 문제 5. 중단과 재시작 (문제 6 과 연속)
//
// (a) settlementRestartJob 을 failAt=30000 으로 실행해 실패시키십시오.
// 실행 후 read_count / commit_count / status 를 기록하십시오.
//
// read_count = ____
// commit_count= ____
// status = ____
//
// (b) settlement 테이블에는 몇 행이 있습니까? 왜 그 숫자입니까?
// 힌트: 실패한 청크는 롤백됩니다.
//
// (c) 같은 파라미터로 다시 실행하면 새 JobInstance 가 생깁니까,
// 아니면 같은 JobInstance 에 붙습니까? 왜입니까?
//
// ⚠️ 여기서 메타데이터를 지우지 마십시오. 문제 6 에서 이어서 씁니다.
// =====================================================================
// (a) 기록
// 여기에 작성:
//
// (b) settlement 행 수와 이유
// 여기에 작성:
//
// (c) JobInstance 판정
// 여기에 작성:
//
// -- 검증:
// SELECT ji.JOB_INSTANCE_ID, je.JOB_EXECUTION_ID, je.STATUS,
// se.READ_COUNT, se.WRITE_COUNT, se.COMMIT_COUNT
// FROM BATCH_JOB_INSTANCE ji
// JOIN BATCH_JOB_EXECUTION je USING (JOB_INSTANCE_ID)
// JOIN BATCH_STEP_EXECUTION se USING (JOB_EXECUTION_ID)
// WHERE ji.JOB_NAME = 'settlementRestartJob'
// ORDER BY je.JOB_EXECUTION_ID;
// =====================================================================
// 문제 6. 상태를 저장하는 Reader 직접 만들기 (문제 5 에서 이어짐)
//
// Practice.NaiveReader 는 재시작해도 처음부터 다시 읽습니다.
//
// (a) 문제 5 의 Job 을 NaiveReader 로 바꿔 재시작하십시오.
// read_count 가 얼마입니까? 왜 그렇습니까?
//
// (b) ItemStreamReader 를 구현해 재시작이 되게 만드십시오.
// open / update / close 를 각각 언제 호출하는지도 적으십시오.
//
// (c) ⚠️ ExecutionContext 의 키 이름을 그냥 "index" 로 하면
// 어떤 문제가 생깁니까? 한 Step 에 Reader 가 둘이라면?
//
// (d) 재시작 여부는 어떻게 판정합니까? open() 안의 관용구를 적으십시오.
//
// (e) 마지막 질문: 이걸 직접 만들어야 합니까?
// 실무에서는 어떻게 하는 것이 옳습니까?
// =====================================================================
// (a) NaiveReader 의 read_count 와 이유
// 여기에 작성:
//
public static class Problem6Reader implements ItemStreamReader<Order> {
private final List<Order> orders;
private int index = 0;
public Problem6Reader(List<Order> orders) {
this.orders = orders;
}
@Override
public void open(ExecutionContext ctx) throws ItemStreamException {
// 여기에 작성: 재시작이면 저장된 위치를 복원하십시오
//
}
@Override
public Order read() {
// 여기에 작성:
//
return null;
}
@Override
public void update(ExecutionContext ctx) throws ItemStreamException {
// 여기에 작성: 현재 위치를 기록하십시오
//
}
@Override
public void close() throws ItemStreamException {
// 여기에 작성:
//
}
}
// (b) open / update / close 호출 시점
// 여기에 작성:
//
// (c) 키 이름을 "index" 로 하면 생기는 문제
// 여기에 작성:
//
// (d) 재시작 판정 관용구
// 여기에 작성:
//
// (e) 직접 만들어야 하는가?
// 여기에 작성:
//
// -- 검증: 재시작 후 read_count 가 41,000 (= 70,000 - 29,000) 이어야 합니다.
// SELECT se.STEP_EXECUTION_ID, se.READ_COUNT, se.WRITE_COUNT, se.STATUS
// FROM BATCH_STEP_EXECUTION se
// JOIN BATCH_JOB_EXECUTION je USING (JOB_EXECUTION_ID)
// JOIN BATCH_JOB_INSTANCE ji USING (JOB_INSTANCE_ID)
// WHERE ji.JOB_NAME = 'settlementRestartJob'
// ORDER BY se.STEP_EXECUTION_ID;
}
Solution.java
6문제의 정답과, "왜 그 답인가"를 설명하는 긴 주석이 들어 있습니다. 풀어 본 뒤에 여세요.
- 정답 1 은
Skip limit of '10' exceeded 가 11번째 불량에서 터진다는 것을 계산으로 보여줍니다. 불량이 order_id % 1000 = 0 이므로 11번째 불량은 order_id = 11000 이고, 이는 COMPLETED 기준 7,700번째 아이템, 즉 8번째 청크 안입니다. "10건 skip 후 죽는다"를 실제 order_id 까지 특정하는 것이 학습 포인트입니다.
- 정답 2 는
instanceof 체인 대신 BinaryExceptionClassifier 를 조합해 쓰는 대안도 함께 제시합니다. 그리고 skipCount 가 Step 전체 누적이라 예외 종류별 한도를 두려면 정책이 직접 카운터를 들어야 한다는 점을, ConcurrentHashMap<Class<?>, LongAdder> 구현으로 보여줍니다.
- 정답 3 의 결론은
commit_count 69,900 → 70 입니다. 처리 단계 skip 은 청크를 한 번 더 처리할 뿐 트랜잭션을 쪼개지 않기 때문입니다. 34.712초 → 8.271초로 줄고, Reader WHERE 절로 아예 거르면 6.233초까지 갑니다. 세 단계의 개선폭을 표로 비교합니다.
- 정답 4 는
.retryLimit(4) 입니다. "재시도 3번"이 아니라 "총 시도 4회"라는 것과, ExponentialBackOffPolicy 의 initialInterval × multiplier^(n-1) 로 최악의 대기 시간이 200+400+800 = 1.4초임을 계산해 둡니다. maxInterval 을 안 걸면 재시도가 많은 Step 에서 대기가 폭주한다는 경고도 있습니다.
- 정답 6 이 가장 깁니다.
ItemStreamReader 의 open / update / close 를 각각 언제 호출하는지, ExecutionContext 의 키에 왜 Reader 이름을 접두사로 붙여야 하는지(setName() 을 안 부르면 두 Reader 가 같은 키를 덮어씁니다), 그리고 open() 에서 context.containsKey(...) 로 재시작 여부를 판정하는 관용구를 설명합니다. 마지막에 "직접 구현하지 말고 JdbcPagingItemReader 를 쓰라"는 결론이 붙습니다 — 직접 만들 수 있게 된 다음에야 안 만드는 선택이 의미가 있기 때문입니다.
package com.example.batch.step11;
import com.example.batch.domain.Order;
import com.example.batch.domain.Settlement;
import org.springframework.batch.core.step.skip.SkipLimitExceededException;
import org.springframework.batch.core.step.skip.SkipPolicy;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemStreamException;
import org.springframework.batch.item.ItemStreamReader;
import org.springframework.dao.DataAccessResourceFailureException;
import org.springframework.dao.DataIntegrityViolationException;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.LongAdder;
/**
* Step 11 — 연습문제 정답과 해설
*
* 문제를 직접 풀어 본 뒤에 여십시오.
*/
public class Solution {
// =====================================================================
// 정답 1. skipLimit 기본값에 걸려 죽는 지점 특정
// =====================================================================
/*
* (a) 예외 메시지
*
* org.springframework.batch.core.step.skip.SkipLimitExceededException:
* Skip limit of '10' exceeded
*
* 실제 로그:
*
* ERROR 43318 --- [main] o.s.batch.core.step.AbstractStep :
* Encountered an error executing step problem1Step in job settlementSkipJob
* org.springframework.batch.core.step.skip.SkipLimitExceededException:
* Skip limit of '10' exceeded
* at org.springframework.batch.core.step.skip.LimitCheckingItemSkipPolicy
* .shouldSkip(LimitCheckingItemSkipPolicy.java:122)
*
* ⚠️ 핵심: skipLimit 의 기본값은 **10** 입니다.
* .skip(...) 만 쓰고 .skipLimit(...) 을 안 쓰면 10건까지만
* 봐주고 11번째에서 Step 전체를 실패시킵니다.
*
* "skip 을 걸었으니 괜찮겠지" 하고 넘어가면, 불량이 11건 이상인
* 날에 배치가 통째로 죽습니다. 그리고 그런 날은 반드시 옵니다.
*
* (b) 몇 번째 불량에서 죽는가 → **11번째**
*
* LimitCheckingItemSkipPolicy 는 `skipCount >= skipLimit` 일 때
* 예외를 던집니다. 10건까지는 skip 되고, 11번째 시도에서 터집니다.
*
* (c) 그 불량의 order_id → **11000**
*
* 불량은 order_id % 1000 = 0 인 주문입니다.
* 1번째 불량 = order_id 1000
* 2번째 불량 = order_id 2000
* ...
* 11번째 불량 = order_id 11000
*
* (d) 몇 번째 청크인가 → **8번째 청크**
*
* 계산:
* COMPLETED 조건은 order_id % 10 <= 6 이므로 10개 중 7개입니다.
* order_id 11000 까지의 COMPLETED 건수 = 11000 × 0.7 = 7,700건
* 청크 크기가 1,000 이므로 7,700번째 아이템은
* 7700 / 1000 = 7.7 → **8번째 청크**(7,001~8,000번째 아이템) 안에 있습니다.
*
* 실제 메타데이터로 확인하면 commit_count = 7 입니다.
* 7청크(7,000건)까지는 정상 커밋됐고 8번째 청크에서 죽었기 때문입니다.
*
* +---------------+--------+------------+--------------------+--------------+
* | STEP_NAME | STATUS | READ_COUNT | PROCESS_SKIP_COUNT | COMMIT_COUNT |
* +---------------+--------+------------+--------------------+--------------+
* | problem1Step | FAILED | 8000 | 10 | 7 |
* +---------------+--------+------------+--------------------+--------------+
*
* "10건 skip 후 죽는다"를 아는 것과, 그게 order_id 11000 이고
* 8번째 청크라는 것까지 계산해 내는 것은 다릅니다. 후자를 할 수 있어야
* 장애 대응 때 "어디까지 처리됐나"를 즉시 답할 수 있습니다.
*
* (e) skipLimit(200) 으로 바꾸면
*
* +---------------+-----------+------------+--------------------+--------------+
* | STEP_NAME | STATUS | READ_COUNT | PROCESS_SKIP_COUNT | COMMIT_COUNT |
* +---------------+-----------+------------+--------------------+--------------+
* | problem1Step | COMPLETED | 70000 | 100 | 100 |
* +---------------+-----------+------------+--------------------+--------------+
*
* 100건 전부 skip 되고 COMPLETED 로 끝납니다.
* write_count 는 69,900 입니다 (70,000 - 100).
*
* ⚠️ 그런데 commit_count 가 70 이 아니라 100 입니다.
* 불량이 든 청크가 스캔 모드로 쪼개졌기 때문입니다. 11-5 의 주제입니다.
*/
// =====================================================================
// 정답 2. SkipPolicy 직접 구현
// =====================================================================
/** (a) 판정 순서가 생명입니다. 구체적인 예외를 먼저. */
public static class CorrectSkipPolicy implements SkipPolicy {
private static final int DATA_ERROR_LIMIT = 50;
@Override
public boolean shouldSkip(Throwable t, long skipCount)
throws SkipLimitExceededException {
// ① 인프라 장애를 "가장 먼저" 걸러냅니다.
// DataAccessResourceFailureException 은 DataAccessException 의
// 하위 타입이므로, 이 검사가 아래로 내려가면 절대 도달하지 못합니다.
if (t instanceof DataAccessResourceFailureException) {
return false; // 절대 skip 하지 않음 → Step 을 실패시킴
}
// ② 데이터 오류는 한도까지 허용
if (t instanceof IllegalArgumentException
|| t instanceof DataIntegrityViolationException) {
if (skipCount >= DATA_ERROR_LIMIT) {
throw new SkipLimitExceededException(DATA_ERROR_LIMIT, t);
}
return true;
}
// ③ 모르는 예외는 skip 하지 않습니다.
// "모르면 멈춘다"가 안전한 기본값입니다.
return false;
}
}
/*
* (b) 순서를 뒤집으면
*
* if (t instanceof DataAccessException) { ... return true; } // 먼저
* if (t instanceof DataAccessResourceFailureException) { ... } // 도달 불가
*
* DataAccessResourceFailureException 은 DataAccessException 을 상속하므로
* 첫 번째 분기에 걸려 **skip 됩니다.**
*
* 실제로 벌어지는 일:
* - DB 커넥션이 끊깁니다.
* - 남은 모든 아이템이 전부 같은 예외로 실패합니다.
* - 정책이 전부 skip 이라고 답합니다.
* - skipLimit 에 걸릴 때까지 skip 하다가... 한도가 크면 끝까지 갑니다.
* - Job 이 **COMPLETED 로 끝납니다.**
*
* 정산이 하나도 안 됐는데 배치는 성공했다고 보고합니다.
* 다음 날 아침 정산 담당자가 발견합니다.
*
* 컨테이너를 잠깐 멈춰 재현할 수 있습니다:
* docker compose pause mysql (5초 뒤) unpause
*
* ⚠️ 일반 원칙: **skip 대상은 화이트리스트로 좁게 지정하십시오.**
* "이 예외만 skip" 이라고 열거하는 것이, "이건 빼고 다 skip" 보다
* 항상 안전합니다. 인프라 예외는 그 자체로 "지금 배치를 계속하면
* 안 된다"는 신호입니다.
*
* (c) 예외 종류별 한도를 두려면
*
* shouldSkip 의 `skipCount` 파라미터는 **Step 전체의 누적 skip 수**입니다.
* 예외 종류를 구분하지 않습니다. 따라서 종류별 한도를 두려면
* **정책이 직접 카운터를 들어야** 합니다.
*/
/** (c) 예외 종류별 한도를 갖는 정책. */
public static class PerTypeSkipPolicy implements SkipPolicy {
private final Map<Class<?>, Integer> limits = Map.of(
IllegalArgumentException.class, 50,
DataIntegrityViolationException.class, 200
);
private final Set<Class<?>> neverSkip = Set.of(
DataAccessResourceFailureException.class
);
// 예외 타입별 누적 카운터.
// ⚠️ ConcurrentHashMap + LongAdder 를 쓰는 이유는 멀티스레드 Step
// (Step 13)에서도 이 정책이 안전해야 하기 때문입니다.
// SkipPolicy 는 여러 스레드가 동시에 호출할 수 있습니다.
private final Map<Class<?>, LongAdder> counters = new ConcurrentHashMap<>();
@Override
public boolean shouldSkip(Throwable t, long skipCount)
throws SkipLimitExceededException {
for (Class<?> never : neverSkip) {
if (never.isInstance(t)) {
return false;
}
}
for (Map.Entry<Class<?>, Integer> e : limits.entrySet()) {
if (e.getKey().isInstance(t)) {
LongAdder counter = counters.computeIfAbsent(
e.getKey(), k -> new LongAdder());
counter.increment();
if (counter.sum() > e.getValue()) {
throw new SkipLimitExceededException(e.getValue(), t);
}
return true;
}
}
return false;
}
}
/*
* 대안 — BinaryExceptionClassifier 조합
*
* instanceof 체인 대신 Spring Retry 의 분류기를 쓸 수도 있습니다.
*
* BinaryExceptionClassifier skippable = new BinaryExceptionClassifier(
* Map.of(IllegalArgumentException.class, true,
* DataAccessResourceFailureException.class, false),
* false); // 기본값 false = 모르면 skip 안 함
*
* BinaryExceptionClassifier 는 **가장 가까운 상위 타입**을 찾아 매칭하므로
* 선언 순서에 의존하지 않습니다. 즉 (b) 의 순서 함정이 원천적으로
* 없습니다. 예외 종류가 많아지면 이쪽이 안전합니다.
*/
// =====================================================================
// 정답 3. 스캔 모드를 피해 commit_count 를 70 으로
// =====================================================================
/*
* (a) 핵심 아이디어
*
* **쓰기 단계에서 터질 조건을 처리 단계에서 미리 판정해 걸러낸다.**
*
* 왜 이게 효과가 있는가:
* - 쓰기(write) 에서 예외가 나면 청크 전체가 롤백되고 스캔 모드로
* 들어갑니다. 1,000건짜리 청크가 1,000개의 트랜잭션이 됩니다.
* - 처리(process) 에서 null 을 반환해 필터링하면 **예외가 아닙니다.**
* 롤백도 없고 스캔 모드도 없습니다. 그냥 그 아이템이 빠질 뿐입니다.
*
* 같은 "100건을 제외한다"는 결과인데, 예외로 하느냐 필터로 하느냐에
* 따라 commit_count 가 69,900 과 70 으로 갈립니다.
*/
/** (b) 구현 — 청크 단위로 미리 조회해 판정. */
public static class PreCheckProcessor implements ItemProcessor<Order, Settlement> {
private final Set<Long> alreadySettled; // 미리 로드한 기존 정산 order_id
private final Practice.SettlementProcessor delegate =
new Practice.SettlementProcessor();
public PreCheckProcessor(Set<Long> alreadySettled) {
this.alreadySettled = alreadySettled;
}
@Override
public Settlement process(Order order) {
if (alreadySettled.contains(order.order_id())) {
return null; // 필터링. 예외가 아닙니다.
}
return delegate.process(order);
}
}
/*
* (c) ⚠️ 함정 — 아이템마다 검증 쿼리를 날리면
*
* 가장 먼저 떠오르는 구현은 이것입니다.
*
* if (jdbcTemplate.queryForObject(
* "SELECT COUNT(*) FROM settlement WHERE order_id = ?",
* Integer.class, order.order_id()) > 0) {
* return null;
* }
*
* 동작은 합니다. 그런데 **70,000번의 SELECT** 가 발생합니다.
* 측정하면 48.3초입니다. 스캔 모드(34.712초)보다 오히려 느립니다.
* "고쳤는데 더 느려졌다"는 전형적인 사례입니다.
*
* 해결책 두 가지:
*
* ① Step 시작 시 한 번에 로드 (위 PreCheckProcessor)
*
* @Bean @StepScope
* public Set<Long> alreadySettled(DataSource ds) {
* return new HashSet<>(new JdbcTemplate(ds).queryForList(
* "SELECT order_id FROM settlement", Long.class));
* }
*
* 쿼리 1회. 100건이면 메모리도 무시할 수준입니다.
* ⚠️ 다만 기존 정산이 수백만 건이면 이 Set 이 메모리를 먹습니다.
* 그 경우 ②를 쓰십시오.
*
* ② 청크 단위로 묶어서 조회 (ItemProcessor 대신 ItemWriter 앞단에서)
*
* 1,000건의 order_id 를 모아 IN 절로 한 번에 조회합니다.
* 쿼리 70회. 메모리는 청크 크기에 비례할 뿐입니다.
*
* (d) 더 나은 답 — Reader 가 아예 안 읽게 한다
*
* 가장 좋은 해법은 애초에 그 행들을 읽지 않는 것입니다.
*
* provider.setFromClause("""
* FROM orders o
* LEFT JOIN settlement s ON s.order_id = o.order_id
* """);
* provider.setWhereClause("""
* WHERE o.status = 'COMPLETED' AND s.order_id IS NULL
* """);
*
* 안티 조인으로 "아직 정산되지 않은 주문"만 읽습니다.
*
* 세 방식의 비교:
*
* | 방식 | 소요 | commit | read | 비고 |
* |---|---|---|---|---|
* | 쓰기 skip (원본) | 34.712초 | 69,900 | 70,000 | 스캔 모드 |
* | 아이템별 검증 쿼리 | 48.300초 | 70 | 70,000 | **더 느려짐** |
* | 미리 로드 + 필터 | 8.271초 | 70 | 70,000 | filter_count 100 |
* | Reader 에서 제외 | 6.233초 | 70 | 69,900 | **최선** |
*
* 34.712초 → 6.233초. 약 5.6배입니다.
*
* 일반 원칙: **거를 수 있으면 최대한 앞에서 거르십시오.**
* Reader > Processor > Writer 순으로 앞일수록 비용이 쌉니다.
* Writer 에서 예외로 거르는 것이 가장 비쌉니다.
*/
// =====================================================================
// 정답 4. retryLimit 계산
// =====================================================================
/*
* (a) retryLimit = **4**
*
* .retryLimit(n) 은 내부적으로 SimpleRetryPolicy(maxAttempts = n) 을
* 만듭니다. maxAttempts 는 **최초 시도를 포함한 총 시도 횟수**입니다.
*
* retryLimit(1) → 재시도 0번 (최초 시도만)
* retryLimit(3) → 재시도 2번
* retryLimit(4) → 재시도 3번 ← 정답
*
* 즉 재시도를 N번 하고 싶으면 retryLimit(N + 1) 입니다.
*
* ⚠️ 이 off-by-one 이 실무에서 SLA 계산을 어긋나게 합니다.
* "3번 재시도하니까 최대 3 × 200ms 대기"라고 계산했는데
* 실제로는 2번만 재시도하고 포기하는 식입니다.
*
* (b) 최악의 총 대기 시간
*
* ExponentialBackOffPolicy(initialInterval=200, multiplier=2.0) 의
* n번째 재시도 전 대기 = initialInterval × multiplier^(n-1)
*
* 1번째 재시도 전: 200 × 2^0 = 200ms
* 2번째 재시도 전: 200 × 2^1 = 400ms
* 3번째 재시도 전: 200 × 2^2 = 800ms
* ─────────────────────────────────
* 총 대기 = 1,400ms
*
* retryLimit(4) = 최초 1회 + 재시도 3회이므로 대기는 3번 발생합니다.
* **최악의 경우 1.4초**입니다.
*
* (c) maxInterval 미설정 시
*
* ExponentialBackOffPolicy 의 maxInterval 기본값은 30,000ms(30초)입니다.
* 완전히 무제한은 아니지만 충분히 위험합니다.
*
* retryLimit 을 10 으로 두면:
* 200 → 400 → 800 → 1600 → 3200 → 6400 → 12800 → 25600 → 30000 → 30000
* 총 111초입니다. 한 청크에서.
*
* 그리고 이것이 **청크마다** 일어납니다. 70청크 전부에서 재시도가
* 발생하면 배치가 2시간을 넘깁니다. 야간 배치 윈도를 넘겨
* 아침 서비스에 영향을 줍니다.
*
* maxInterval 을 명시적으로 짧게(5초 정도) 잡으십시오.
* 그리고 재시도로 안 풀리는 문제는 재시도로 풀리지 않습니다.
*
* (d) .retryLimit() 을 빠뜨리면 → **조용히 아무 일도 안 일어납니다**
*
* FaultTolerantStepBuilder 의 retryLimit 필드 기본값은 **0** 입니다.
* 이 상태에서는 등록한 예외가 재시도 없이 그대로 전파됩니다.
*
* - 컴파일 에러 없음
* - 경고 로그 없음
* - .retry(...) 를 쓴 흔적은 코드에 남아 있음
*
* "retry 를 걸어 뒀으니 데드락은 알아서 복구되겠지"라고 믿고 있는데
* 실제로는 첫 데드락에서 배치가 죽습니다.
*
* ⚠️ .retry() 와 .retryLimit() 은 **항상 짝으로** 쓰십시오.
* skip 의 기본값이 10 인 것과 대조적입니다 — skip 은 조금이라도
* 동작하는데 retry 는 아예 동작하지 않습니다.
*/
// =====================================================================
// 정답 5. 중단과 재시작
// =====================================================================
/*
* (a) failAt=30000 으로 실행한 결과
*
* +------------------+--------+------------+-------------+--------------+
* | JOB_EXECUTION_ID | STATUS | READ_COUNT | WRITE_COUNT | COMMIT_COUNT |
* +------------------+--------+------------+-------------+--------------+
* | 1 | FAILED | 30000 | 29000 | 29 |
* +------------------+--------+------------+-------------+--------------+
*
* read_count = 30,000 (30,000번째를 읽고 나서 처리 중 죽음)
* write_count = 29,000
* commit_count = 29
* status = FAILED
*
* (b) settlement 행 수 → **29,000**
*
* 29개 청크(29,000건)가 정상 커밋됐고, 30번째 청크는 처리 도중
* 예외가 나서 **통째로 롤백**됐습니다. 그 청크의 29,001~30,000번째
* 아이템은 하나도 저장되지 않았습니다.
*
* 이것이 "청크 = 트랜잭션 = 롤백 단위"의 의미입니다.
* 부분적으로 저장되는 일은 없습니다. 전부 아니면 전무입니다.
*
* (c) 같은 JobInstance 에 붙습니다
*
* 두 가지 조건이 모두 만족되기 때문입니다.
*
* ① 직전 실행이 FAILED 입니다.
* COMPLETED 였다면 JobInstanceAlreadyCompleteException 이 나서
* 아예 실행되지 않습니다(Step 03).
* FAILED 는 "재시작 가능" 상태입니다.
*
* ② failAt 이 identifying = false 입니다.
* JOB_KEY 는 identifying 파라미터만으로 계산됩니다. failAt 이
* 제외되므로 failAt=30000 이든 999999 든 **JOB_KEY 가 같습니다.**
* 따라서 같은 JobInstance 로 인식되어 재시작이 됩니다.
*
* ⚠️ 만약 identifying=true 였다면 failAt 값이 달라지는 순간
* JOB_KEY 가 바뀌어 **새 JobInstance** 가 만들어집니다.
* 재시작이 아니라 처음부터 새로 도는 것이고, 이미 정산된
* 29,000건과 UNIQUE 충돌이 나서 즉시 죽습니다.
* Practice 의 주석이 이 false 를 강조한 이유입니다.
*
* 메타데이터로 확인하면 JOB_INSTANCE_ID 가 같고 JOB_EXECUTION_ID 만
* 늘어난 것이 보입니다.
*
* +-----------------+------------------+-----------+------------+
* | JOB_INSTANCE_ID | JOB_EXECUTION_ID | STATUS | READ_COUNT |
* +-----------------+------------------+-----------+------------+
* | 1 | 1 | FAILED | 30000 |
* | 1 | 2 | COMPLETED | 41000 |
* +-----------------+------------------+-----------+------------+
*
* 1:N 관계가 눈에 보입니다. Step 01 의 다이어그램 그대로입니다.
*/
// =====================================================================
// 정답 6. 상태를 저장하는 Reader — 이 스텝의 결론
// =====================================================================
/*
* (a) NaiveReader 로 재시작하면 read_count = **70,000**
*
* 재시작인데 처음부터 70,000건을 다시 읽었습니다.
*
* 왜냐하면 NaiveReader 는 ItemReader 만 구현했고 ItemStreamReader 가
* 아니기 때문입니다. open/update/close 콜백이 아예 없으므로
* Spring Batch 는 이 Reader 의 위치를 저장할 방법도, 복원할 방법도
* 없습니다. 매번 index = 0 에서 시작합니다.
*
* 그리고 결과가 갈립니다.
* - 비멱등 Writer(순수 INSERT): 이미 정산된 29,000건과 UNIQUE 가
* 충돌해 즉시 DuplicateKeyException. 시끄럽게 실패합니다.
* - 멱등 Writer(ON DUPLICATE KEY UPDATE): **조용히 29,000건을
* 다시 정산합니다.** 결과 데이터는 우연히 맞지만, 41,000건이면
* 될 일을 70,000건 처리했습니다. 그리고 만약 수수료율이 그 사이
* 바뀌었다면 이미 정산된 건이 새 요율로 덮어써집니다.
*
* 후자가 훨씬 위험합니다. 이 코스가 반복하는 주제입니다.
*
* (b) open / update / close 호출 시점
*
* | 콜백 | 호출 시점 | 횟수 | 용도 |
* |---|---|---|---|
* | open(ctx) | Step 시작 직후 | 1회 | 자원 열기 + **재시작 위치 복원** |
* | update(ctx) | **매 청크 커밋 직전** | 청크 수만큼 | 현재 위치 기록 |
* | close() | Step 종료 시 | 1회 | 자원 정리 (성공/실패 무관) |
*
* ⚠️ update() 가 "커밋 직전"이라는 것이 중요합니다.
* 청크 트랜잭션과 **같은 트랜잭션 안에서** ExecutionContext 가
* 저장되므로, 데이터 저장과 위치 기록이 원자적으로 함께 커밋됩니다.
* 롤백되면 둘 다 롤백됩니다. 그래서 "29,000건 저장 + 위치 29,000"이
* 항상 일관됩니다. 프로젝트 셋업에서 메타데이터와 업무 데이터를
* 같은 DataSource 에 둔 이유가 이것입니다.
*/
/** (b)(c)(d) 정답 구현. */
public static class RestartableReader implements ItemStreamReader<Order> {
// (c) 키에 Reader 이름을 접두사로 붙입니다.
private static final String KEY_INDEX = "restartableReader.index";
private final List<Order> orders;
private int index = 0;
public RestartableReader(List<Order> orders) {
this.orders = orders;
}
@Override
public void open(ExecutionContext ctx) throws ItemStreamException {
// (d) 재시작 판정 관용구
if (ctx.containsKey(KEY_INDEX)) {
this.index = ctx.getInt(KEY_INDEX);
} else {
this.index = 0;
}
}
@Override
public Order read() {
return index >= orders.size() ? null : orders.get(index++);
}
@Override
public void update(ExecutionContext ctx) throws ItemStreamException {
ctx.putInt(KEY_INDEX, index);
}
@Override
public void close() throws ItemStreamException {
// 정리할 자원 없음. 파일이나 커넥션을 열었다면 여기서 닫습니다.
}
}
/*
* (c) 키 이름을 "index" 로 하면
*
* ExecutionContext 는 **Step 전체가 공유하는 하나의 Map** 입니다.
* Reader 별로 분리된 공간이 아닙니다.
*
* 한 Step 에 Reader 가 둘이면(예: 두 파일을 번갈아 읽는 구성이나
* CompositeItemReader) 둘 다 "index" 키에 쓰게 되어 **서로를
* 덮어씁니다.** 재시작하면 A Reader 가 B Reader 의 위치에서
* 시작합니다.
*
* 그래서 Spring Batch 의 기본 Reader 들은 전부 setName() 을 요구하고,
* 그 이름을 키 접두사로 씁니다.
*
* new JdbcPagingItemReaderBuilder<Order>()
* .name("orderReader") // ← 이것이 키 접두사가 됩니다
*
* 실제 저장된 값을 보면 이렇게 생겼습니다.
*
* SELECT SHORT_CONTEXT FROM BATCH_STEP_EXECUTION_CONTEXT ...
* {"@class":"java.util.HashMap","orderReader.read.count":29000}
*
* ⚠️ .name() 을 빠뜨리면 ItemStreamSupport 가
* "ItemStream must have a name" 예외를 던집니다.
* 단, saveState(false) 면 이름이 필요 없어 통과합니다 —
* 상태를 저장하지 않으니 키도 필요 없기 때문입니다.
*
* (e) 이걸 직접 만들어야 하는가 → **아니오**
*
* 실무에서는 `JdbcPagingItemReader`, `JdbcCursorItemReader`,
* `FlatFileItemReader` 같은 기본 제공 Reader 를 쓰십시오.
* 전부 ItemStreamReader 를 이미 올바르게 구현하고 있습니다.
* 재시작, 이름 기반 키, 스레드 안전성 문서화까지 다 되어 있습니다.
*
* 직접 만들어야 하는 경우는 정말 드뭅니다.
* (외부 API 를 페이지네이션하며 읽는 Reader 정도)
*
* 그런데 왜 이 문제를 풀었는가:
*
* **직접 만들 수 있게 된 다음에야, 안 만드는 선택이 의미가 있기
* 때문입니다.**
*
* ItemStreamReader 의 계약을 모르면 다음을 이해할 수 없습니다.
* - 왜 Reader 에 .name() 을 줘야 하는지
* - 왜 saveState(false) 를 켜면 재시작이 안 되는지 (Step 13)
* - 왜 @Bean 이 아닌 Reader 는 콜백을 못 받는지 (Step 08, 09)
* - 재시작 시 read_count 가 왜 41,000 인지
*
* 이 네 가지는 전부 open/update/close 계약에서 나옵니다.
* 기본 Reader 를 "그냥 쓰는" 사람과 "왜 그렇게 동작하는지 아는"
* 사람의 차이가 장애 대응에서 갈립니다.
*/
}