Step 08 — ItemWriter
학습 목표
- Spring Batch 5.x 의
void write(Chunk<? extends T> chunk) 시그니처를 이해하고, 4.x 의 List<T> 에서 무엇이 왜 바뀌었는지 설명한다
JdbcBatchItemWriter 로 70,000건을 쓰고, record 에 beanMapped() 가 통하지 않는 이유를 확인한 뒤 itemSqlParameterSourceProvider 람다로 해결한다
- JDBC URL 의
rewriteBatchedStatements=true 유무로 70,000건 쓰기 시간을 실측해 약 8배 차이를 확인한다
assertUpdates 옵션이 무엇을 잡고 무엇을 못 잡는지 구분한다
FlatFileItemWriter 로 DelimitedLineAggregator · headerCallback · footerCallback 를 갖춘 CSV 를 만든다
CompositeItemWriter / ClassifierCompositeItemWriter 로 출력을 복제하거나 분기한다
ItemStream 의 open/update/close 생명주기를 이해하고, 콜백이 호출되지 않아 파일이 0바이트가 되는 사고를 재현한 뒤 .stream() 수동 등록으로 고친다
선행 스텝: Step 07 — ItemProcessor
예상 소요: 100분
8-0. 실습 준비
mysql -h127.0.0.1 -P3308 -uroot -proot1234 batchdb <<'SQL'
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;
TRUNCATE TABLE settlement;
SQL
mkdir -p ./out
이 스텝은 Step 07 의 파이프라인(리더 10만 건 → 프로세서가 3만 건 필터 → 70,000건 출력)을 그대로 이어받습니다. 바뀌는 것은 writer 뿐입니다.
8-1. 5.x 의 write(Chunk<? extends T>) — 무엇이 바뀌었나
Spring Batch 4.x 의 ItemWriter 는 이랬습니다.
// Spring Batch 4.x
public interface ItemWriter<T> {
void write(List<? extends T> items) throws Exception;
}
Spring Batch 5.0 부터는 이렇습니다.
// Spring Batch 5.x
@FunctionalInterface
public interface ItemWriter<T> {
void write(Chunk<? extends T> chunk) throws Exception;
}
java.util.List 가 org.springframework.batch.item.Chunk 로 바뀌었습니다. 소스 호환성이 깨지는 변경입니다. 4.x 예제를 그대로 붙여 넣으면 이렇게 됩니다.
결과
> Task :compileJava FAILED
/src/main/java/com/example/batch/step08/LegacyWriter.java:14: error: LegacyWriter is not abstract and does not override abstract method write(Chunk<? extends Settlement>) in ItemWriter
public class LegacyWriter implements ItemWriter<Settlement> {
^
/src/main/java/com/example/batch/step08/LegacyWriter.java:17: error: method does not override or implement a method from a supertype
@Override
^
2 errors
BUILD FAILED in 2s
이건 좋은 실패입니다. 컴파일러가 잡아 줬습니다. 고치는 법은 파라미터 타입만 바꾸는 것입니다.
public class SettlementLogWriter implements ItemWriter<Settlement> {
private static final Logger log = LoggerFactory.getLogger(SettlementLogWriter.class);
@Override
public void write(Chunk<? extends Settlement> chunk) {
log.debug("chunk size={} first={} last={}",
chunk.size(),
chunk.getItems().get(0).orderId(),
chunk.getItems().get(chunk.size() - 1).orderId());
}
}
Chunk 가 List 보다 나은 이유는 아이템 말고도 정보를 실을 수 있기 때문입니다.
| 메서드 | 설명 |
|---|
getItems() | List<T> 반환. 4.x 코드를 옮길 때 여기서 꺼내 쓰면 됩니다 |
size() / isEmpty() | 아이템 개수 |
iterator() | Chunk 는 Iterable<T> 이므로 for (Settlement s : chunk) 가 됩니다 |
getSkips() | 이 청크에서 스킵된 아이템 목록 (fault tolerant 모드) |
getErrors() | 이 청크에서 발생한 예외 목록 |
clear() / add(T) | 커스텀 writer 에서 청크를 조작할 때 |
@FunctionalInterface 가 붙은 것도 실질적인 변화입니다. 람다로 writer 를 쓸 수 있습니다.
ItemWriter<Settlement> writer = chunk -> chunk.forEach(System.out::println);
💡 실무 팁 — 4.x 코드를 옮길 때는 getItems() 한 줄만
// 4.x
public void write(List<? extends Settlement> items) { doWrite(items); }
// 5.x — 본문은 그대로 두고 어댑터만
public void write(Chunk<? extends Settlement> chunk) { doWrite(chunk.getItems()); }
대부분의 마이그레이션은 이 한 줄로 끝납니다. getSkips() 같은 새 기능은 필요할 때 나중에 쓰면 됩니다.
8-2. JdbcBatchItemWriter — 그리고 record 가 조용히 부수는 것
가장 흔한 writer 입니다. 스프링 문서와 대부분의 예제는 이렇게 씁니다.
// ❌ record 에는 통하지 않습니다
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)
""")
.beanMapped() // ← 여기
.build();
beanMapped() 는 내부적으로 BeanPropertyItemSqlParameterSourceProvider 를 씁니다. 이름 그대로 자바빈 규약(getOrderId())을 기대합니다. 그런데 record 의 접근자는 orderId() 입니다. get 이 없습니다.
실행하면 이렇습니다.
결과
INFO 61022 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
ERROR 61022 --- [ main] o.s.batch.core.step.AbstractStep : Encountered an error executing step settlementStep in job settlementJob
org.springframework.dao.InvalidDataAccessApiUsageException: No value supplied for the SQL parameter 'orderId': Invalid property 'orderId' of bean class [com.example.batch.domain.Settlement]: Bean property 'orderId' is not readable or has an invalid getter method: Does the return type of the getter match the parameter type of the setter?
at org.springframework.jdbc.core.namedparam.NamedParameterUtils.substituteNamedParameters(NamedParameterUtils.java:392)
at org.springframework.batch.item.database.JdbcBatchItemWriter.write(JdbcBatchItemWriter.java:189)
이 경우는 운이 좋습니다. 시끄럽게 실패했으니까요.
⚠️ 함정 — 컬럼이 전부 NULL 로 들어가는 조용한 실패
위 예외는 settlement 컬럼들이 NOT NULL 이라서 명확하게 터진 것입니다. NULL 을 허용하는 테이블이라면 예외 없이 전부 NULL 인 7만 행이 들어갑니다.
SELECT COUNT(*) 는 70000 이 나오고, 로그는 COMPLETED 이고, writeCount=70000 입니다. 모든 지표가 정상입니다. 정산 금액만 전부 비어 있습니다.
이 사고를 만드는 조합은 세 가지가 맞물릴 때입니다: ① 도메인이 record ② beanMapped() ③ 대상 테이블이 nullable.
세 번째를 없애는 것이 가장 확실한 방어입니다. 결과 테이블의 컬럼은 되도록 NOT NULL 로 만드세요. 제약은 성능을 위한 것이기 이전에 조용한 실패를 시끄럽게 만드는 장치입니다.
record 에 맞는 해법은 람다로 파라미터 소스를 직접 만드는 것입니다.
@Bean
public JdbcBatchItemWriter<Settlement> settlementJdbcWriter(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()))
.assertUpdates(true)
.build();
}
장황해 보이지만 컴파일 타임에 검증되는 코드입니다. record 의 필드 이름을 바꾸면 여기서 컴파일 에러가 납니다. beanMapped() 는 리팩터링을 따라오지 못합니다.
| 방식 | record 지원 | 오타 검출 | 비고 |
|---|
.beanMapped() | ❌ | 런타임 | 자바빈(POJO + getter)에만 |
.itemSqlParameterSourceProvider(람다) | ✅ | 컴파일 타임 | 권장 |
.columnMapped() | Map<String,Object> 아이템 전용 | 런타임 | 리더가 Map 을 낼 때 |
💡 읽기 쪽도 대칭입니다. Step 06 에서 본 것처럼 읽기는 DataClassRowMapper, 쓰기는 람다가 record 의 정석 조합입니다.
BeanPropertyRowMapper 는 읽기 쪽에서 같은 이유로 record 를 못 채웁니다.
8-3. rewriteBatchedStatements=true — 70,000건 쓰기 실측
이제 이 코스에서 가장 가성비 좋은 한 줄입니다.
JdbcBatchItemWriter 는 청크의 1,000건을 PreparedStatement.addBatch() 로 쌓고 executeBatch() 를 한 번 부릅니다. 자바 코드상으로는 "한 번"입니다. 그런데 MySQL JDBC 드라이버가 그것을 어떻게 보내느냐는 별개의 문제입니다.
rewriteBatchedStatements=false (기본값)
executeBatch()
→ INSERT INTO settlement VALUES (...) ← 왕복 1
→ INSERT INTO settlement VALUES (...) ← 왕복 2
→ ... ← 왕복 1000
(청크 하나에 1,000번 왕복. 70청크면 70,000번)
rewriteBatchedStatements=true
executeBatch()
→ INSERT INTO settlement VALUES (...),(...),(...), ... ,(...) ← 왕복 1
(청크 하나에 1번 왕복. 70청크면 70번)
측정합니다. 먼저 옵션을 뺀 URL 로.
spring:
datasource:
url: jdbc:mysql://127.0.0.1:3308/batchdb?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Seoul
mysql -h127.0.0.1 -P3308 -uroot -proot1234 batchdb -e "TRUNCATE TABLE settlement;"
./gradlew bootRun --args='--spring.batch.job.name=settlementJob'
결과
INFO 62110 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] launched with the following parameters: [{'run.id':'{value=1, type=class java.lang.Long, identifying=true}'}]
INFO 62110 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
INFO 62110 --- [ main] c.e.batch.step08.CountLogger : readCount=100000, writeCount=70000, filterCount=30000, skipCount=0
INFO 62110 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 1m1s412ms
INFO 62110 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [COMPLETED] in 1m1s688ms
61.4초. 이제 URL 끝에 옵션 하나만 붙입니다.
url: jdbc:mysql://127.0.0.1:3308/batchdb?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Seoul&rewriteBatchedStatements=true
mysql -h127.0.0.1 -P3308 -uroot -proot1234 batchdb -e "TRUNCATE TABLE settlement;"
./gradlew bootRun --args='--spring.batch.job.name=settlementJob'
결과
INFO 62344 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [settlementStep]
INFO 62344 --- [ main] c.e.batch.step08.CountLogger : readCount=100000, writeCount=70000, filterCount=30000, skipCount=0
INFO 62344 --- [ main] o.s.batch.core.step.AbstractStep : Step: [settlementStep] executed in 7s702ms
INFO 62344 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=settlementJob]] completed with the following parameters: [...] and the following status: [COMPLETED] in 7s934ms
61.4초 → 7.7초. 약 8배 빨라졌습니다.
| 항목 | rewriteBatchedStatements 없음 | 있음 |
|---|
| Step 실행 시간 | 1m 1s 412ms | 7s 702ms |
| INSERT 왕복 횟수 | 70,000 | 70 |
writeCount | 70000 | 70000 |
| 결과 행 수 | 70000 | 70000 |
| Job 상태 | COMPLETED | COMPLETED |
아래 네 줄이 완전히 같다는 점을 보세요. 결과도, 카운터도, 로그 상태도 동일합니다. 다른 것은 시간뿐입니다.
서버 쪽에서도 확인할 수 있습니다.
mysql -h127.0.0.1 -P3308 -uroot -proot1234 -e "
SELECT VARIABLE_NAME, VARIABLE_VALUE FROM performance_schema.global_status
WHERE VARIABLE_NAME IN ('Com_insert','Questions');"
결과 (옵션 없이 한 번 돌린 뒤)
+---------------+----------------+
| VARIABLE_NAME | VARIABLE_VALUE |
+---------------+----------------+
| Com_insert | 70000 |
| Questions | 70412 |
+---------------+----------------+
결과 (옵션을 켜고 한 번 돌린 뒤)
+---------------+----------------+
| VARIABLE_NAME | VARIABLE_VALUE |
+---------------+----------------+
| Com_insert | 70 |
| Questions | 482 |
+---------------+----------------+
Com_insert 가 70,000 에서 70 으로 줄었습니다. DB 가 한 일의 양이 아니라, 왕복 횟수가 줄어든 것입니다.
⚠️ 함정 — 이 옵션이 없으면 "느린 것"이 아니라 "느린 줄 모르는 것"입니다
이 옵션의 부재는 어떤 에러도, 어떤 경고도 남기지 않습니다. writeCount=70000 이 정확히 찍히고 Job 은 COMPLETED 입니다.
그래서 "우리 배치는 원래 한 시간 걸려요"가 됩니다. 청크 크기를 늘려도, 인덱스를 지워도, 서버를 키워도 8배가 안 나옵니다. 병목이 DB 가 아니라 네트워크 왕복이기 때문입니다.
진단법은 간단합니다. 위처럼 Com_insert 를 재 보세요. 처리 건수와 같은 숫자가 나오면 배치 쓰기가 배치가 아닌 것입니다.
💡 실무 팁 — rewriteBatchedStatements 는 MySQL 전용이며 주의점이 둘 있습니다
- PostgreSQL·Oracle 에는 이 옵션이 없습니다. 두 드라이버는 기본적으로 배치를 제대로 묶어 보냅니다. MySQL 만의 기본값 문제입니다.
- 켜면
max_allowed_packet 에 걸릴 수 있습니다. 청크 크기 × 행 길이가 패킷 한도(기본 64MB)를 넘으면 드라이버가 알아서 쪼개지만, 청크를 10만 단위로 키운다면 SHOW VARIABLES LIKE 'max_allowed_packet'; 를 확인하세요.
INSERT ... ON DUPLICATE KEY UPDATE 도 재작성 대상입니다. 다만 UPDATE/DELETE 배치는 멀티 VALUES 로 합칠 수 없어 여러 문장을 세미콜론으로 이어 보내는 방식이 되고, 이때 반환되는 갱신 건수가 달라집니다 — 다음 절의 주제입니다.
8-4. assertUpdates — 무엇을 잡고 무엇을 못 잡는가
JdbcBatchItemWriter 의 assertUpdates 는 기본값이 true 입니다. executeBatch() 가 돌려준 갱신 건수 배열을 훑어 0 인 항목이 있으면 예외를 던집니다.
// JdbcBatchItemWriter 내부 (5.1.1)
if (assertUpdates) {
for (int i = 0; i < updateCounts.length; i++) {
if (updateCounts[i] == 0) {
throw new EmptyResultDataAccessException(
"Item " + i + " of " + updateCounts.length
+ " did not update any rows: [" + chunk.getItems().get(i) + "]", 1);
}
}
}
이게 필요한 상황은 UPDATE writer 입니다. WHERE 절에 걸리는 행이 없으면 갱신 건수가 0 이고, 그건 대개 버그입니다.
.sql("UPDATE settlement SET net_amount = :netAmount WHERE order_id = :orderId")
.assertUpdates(true)
없는 order_id 를 하나 섞어 실행하면:
결과
ERROR 62901 --- [ main] o.s.batch.core.step.AbstractStep : Encountered an error executing step updateStep in job settlementJob
org.springframework.dao.EmptyResultDataAccessException: Item 137 of 1000 did not update any rows: [Settlement[orderId=999999, customerId=4242, settleDate=2025-03-01, grossAmount=5000.00, feeRate=0.0250, feeAmount=125.00, netAmount=4875.00]]
at org.springframework.batch.item.database.JdbcBatchItemWriter.write(JdbcBatchItemWriter.java:206)
assertUpdates(false) 로 두면 이 상황이 아무 흔적 없이 지나갑니다. writeCount 는 여전히 70000 입니다. 갱신되지 않은 행이 몇 개인지 아무도 모릅니다.
⚠️ 함정 — assertUpdates 는 "0건"만 잡습니다
두 가지 구멍이 있습니다.
(1) 드라이버가 SUCCESS_NO_INFO(-2) 를 돌려주면 검사가 무력화됩니다.
rewriteBatchedStatements=true 로 UPDATE/DELETE 배치가 다중 문장으로 재작성되면, MySQL 드라이버는 개별 갱신 건수를 알 수 없어 Statement.SUCCESS_NO_INFO(값 -2)를 채워 돌려줍니다.
위 루프는 == 0 만 봅니다. -2 는 통과합니다. 즉 성능 옵션을 켜는 순간 UPDATE writer 의 안전장치가 조용히 꺼집니다.
→ UPDATE 계열 writer 를 쓴다면 writeCount 를 믿지 말고, Step 뒤에 SELECT COUNT(*) 검증 Tasklet 을 붙이세요(Step 10 의 흐름 제어로 조건 분기까지 만들 수 있습니다).
(2) "1건 기대했는데 500건 갱신됨"은 잡지 못합니다.
WHERE 절을 빠뜨린 UPDATE settlement SET fee_rate = :feeRate 같은 실수는 갱신 건수가 0 이 아니라 70,000 입니다. assertUpdates 는 통과시킵니다. 정산 테이블 전체가 마지막 아이템의 값으로 덮입니다.
이건 옵션으로 막을 수 없습니다. UPDATE 문의 WHERE 절에 유니크 키가 들어 있는지 리뷰하는 것이 유일한 방어입니다.
| 상황 | assertUpdates(true) 가 잡는가 |
|---|
UPDATE 가 0건 갱신 (드라이버가 실제 건수 반환) | ✅ |
UPDATE 가 0건 갱신 (드라이버가 -2 반환) | ❌ |
UPDATE 가 의도보다 많이 갱신 | ❌ |
INSERT 중복 키 | 해당 없음 — DuplicateKeyException 이 먼저 남 |
8-5. FlatFileItemWriter — CSV 만들기
정산 결과를 파일로도 내보냅니다. 회계팀에 넘길 CSV 입니다.
@Bean
public FlatFileItemWriter<Settlement> settlementFileWriter() {
DelimitedLineAggregator<Settlement> aggregator = new DelimitedLineAggregator<>();
aggregator.setDelimiter(",");
aggregator.setFieldExtractor(s -> new Object[]{
s.orderId(), s.customerId(), s.settleDate(),
s.grossAmount(), s.feeRate(), s.feeAmount(), s.netAmount()
});
return new FlatFileItemWriterBuilder<Settlement>()
.name("settlementFileWriter")
.resource(new FileSystemResource("out/settlement.csv"))
.encoding("UTF-8")
.lineSeparator("\n")
.lineAggregator(aggregator)
.headerCallback(w -> w.write(
"order_id,customer_id,settle_date,gross_amount,fee_rate,fee_amount,net_amount"))
.footerCallback(w -> w.write("# generated by settlementJob"))
.shouldDeleteIfExists(true)
.build();
}
FieldExtractor 를 람다로 쓴 것에 주목하세요. BeanWrapperFieldExtractor 는 beanMapped() 와 같은 이유로 record 에 통하지 않습니다. 8-2 와 완전히 같은 함정입니다.
실행하고 결과를 봅니다.
./gradlew bootRun --args='--spring.batch.job.name=fileJob'
wc -l out/settlement.csv && ls -lh out/settlement.csv
결과
70002 out/settlement.csv
-rw-r--r-- 1 julong staff 3.0M Jul 20 14:12 out/settlement.csv
70,000 데이터 행 + 헤더 1 + 푸터 1 = 70,002줄입니다.
head -3 out/settlement.csv && echo '...' && tail -2 out/settlement.csv
결과
order_id,customer_id,settle_date,gross_amount,fee_rate,fee_amount,net_amount
1,1,2025-01-01,1100.00,0.0300,33.00,1067.00
2,2,2025-01-01,1200.00,0.0250,30.00,1170.00
...
99996,996,2025-06-29,1400.00,0.0350,49.00,1351.00
# generated by settlementJob
Step 07 에서 검증한 앞 5건과 같은 값입니다.
주요 옵션을 정리합니다.
| 옵션 | 의미 | 주의 |
|---|
.resource(...) | 출력 경로 | 디렉터리가 없으면 예외. 미리 mkdir |
.encoding("UTF-8") | 인코딩 | 기본은 플랫폼 기본값. 명시하세요 |
.lineSeparator("\n") | 줄 구분자 | 기본은 System.lineSeparator(). 윈도우에서 만든 파일이 리눅스에서 \r 섞여 나오는 사고 방지 |
.append(true) | 이어쓰기 | shouldDeleteIfExists 와 배타적 |
.shouldDeleteIfExists(true) | 기존 파일 삭제 후 새로 씀 | 기본값 true |
.transactional(true) | 청크 커밋 시점까지 버퍼링 | 기본값 true. 롤백되면 그 청크는 파일에도 안 씀 |
.headerCallback(...) | 파일 열 때 1회 | open() 안에서 호출됩니다 — 8-7 의 핵심 |
.footerCallback(...) | 파일 닫을 때 1회 | close() 안에서 호출됩니다 |
💡 실무 팁 — .transactional(true) 의 의미
FlatFileItemWriter 는 기본적으로 청크 커밋 직전까지 출력을 메모리 버퍼에 담아 뒀다가, 트랜잭션이 커밋될 때 실제로 파일에 씁니다.
덕분에 청크가 롤백되면 파일에도 그 청크가 남지 않습니다. DB 와 파일의 내용이 어긋나지 않게 하는 장치입니다.
반대로 말하면 버퍼를 비우는 시점이 커밋 시점이라는 뜻이고, 그래서 close() 가 호출되지 않으면 마지막 버퍼가 통째로 사라집니다.
8-6. CompositeItemWriter / ClassifierCompositeItemWriter
CompositeItemWriter — 같은 청크를 여러 곳에
DB 에도 쓰고 파일에도 쓰고 싶을 때입니다. 모든 델리게이트가 같은 청크를 받습니다.
@Bean
public ItemWriter<Settlement> dualWriter(JdbcBatchItemWriter<Settlement> settlementJdbcWriter,
FlatFileItemWriter<Settlement> settlementFileWriter) {
CompositeItemWriter<Settlement> writer = new CompositeItemWriter<>();
writer.setDelegates(List.of(settlementJdbcWriter, settlementFileWriter));
return writer;
}
결과
INFO 63180 --- [ main] c.e.batch.step08.CountLogger : readCount=100000, writeCount=70000, filterCount=30000, skipCount=0
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -N -e "SELECT COUNT(*) FROM settlement;"
wc -l out/settlement.csv
결과
70000
70002 out/settlement.csv
writeCount 는 70000 한 번만 올라갑니다. 두 곳에 썼다고 140000 이 되지 않습니다. writeCount 는 "writer 를 통과한 아이템 수"이지 "쓰기 연산 수"가 아닙니다.
ClassifierCompositeItemWriter — 아이템마다 한 곳으로
VIP 정산(15,000건)만 별도 파일로 빼고, 나머지 55,000건은 DB 로 보냅니다.
@Bean
public ItemWriter<Settlement> routingWriter(JdbcBatchItemWriter<Settlement> settlementJdbcWriter,
FlatFileItemWriter<Settlement> vipFileWriter) {
ClassifierCompositeItemWriter<Settlement> writer = new ClassifierCompositeItemWriter<>();
BigDecimal vipRate = new BigDecimal("0.0150");
writer.setClassifier((Classifier<Settlement, ItemWriter<? super Settlement>>)
s -> vipRate.compareTo(s.feeRate()) == 0 ? vipFileWriter : settlementJdbcWriter);
return writer;
}
내부적으로는 청크를 분류자로 그룹핑한 뒤, 그룹별로 델리게이트의 write() 를 한 번씩 부릅니다. 아이템마다 한 번씩 부르는 것이 아닙니다 — 배치 성능이 유지됩니다.
결과
INFO 63455 --- [ main] c.e.batch.step08.CountLogger : readCount=100000, writeCount=70000, filterCount=30000, skipCount=0
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -N -e "SELECT COUNT(*) FROM settlement;"
wc -l out/vip-settlement.csv
결과
55000
0 out/vip-settlement.csv
파일이 0줄입니다. DB 는 55,000건으로 정확한데 파일만 비어 있습니다. Job 은 COMPLETED 이고 writeCount=70000 입니다.
이것이 이 스텝의 본론입니다.
8-7. ItemStream — open / update / close
파일에 쓰는 writer 는 상태를 가집니다. 열려 있는 파일 핸들, 지금까지 쓴 줄 수, 버퍼. 이 상태를 프레임워크가 관리하도록 하는 계약이 ItemStream 입니다.
public interface ItemStream {
void open(ExecutionContext executionContext) throws ItemStreamException;
void update(ExecutionContext executionContext) throws ItemStreamException;
void close() throws ItemStreamException;
}
FlatFileItemWriter 는 ItemStreamWriter<T>(= ItemWriter<T> + ItemStream)를 구현합니다. 각 콜백에서 무슨 일을 하는지가 핵심입니다.
| 콜백 | 호출 시점 | FlatFileItemWriter 가 하는 일 |
|---|
open(ctx) | Step 시작 시 1회 | 파일 열기, headerCallback 실행, 재시작이면 ctx 의 written 위치로 truncate |
update(ctx) | 청크 커밋마다 | 버퍼 flush, ctx 에 settlementFileWriter.written=NNNNN 저장 |
close() | Step 종료 시 1회 | footerCallback 실행, flush, 파일 핸들 반납 |
그림으로 보면 이렇습니다.
Step 시작
└── open(ctx) ← 파일 열기 + 헤더 기록
├── [청크 1] write() ... update(ctx) ← flush + "1000줄 썼음" 저장
├── [청크 2] write() ... update(ctx) ← flush + "2000줄 썼음" 저장
└── ...
└── close() ← 푸터 기록 + 마지막 flush + 파일 닫기
Step 종료
세 콜백이 각각 다른 것을 책임집니다.
open 이 안 불리면 → 파일이 안 열리고 헤더가 없습니다.
update 가 안 불리면 → 재시작 시 어디부터 이어 쓸지 모릅니다.
close 가 안 불리면 → 푸터가 없고, 마지막 버퍼가 통째로 사라집니다.
그러면 누가 이 콜백들을 불러 줄까요? StepBuilder 입니다.
// SimpleStepBuilder.registerAsStreams() (5.1.1)
protected void registerStepListenerAsItemListener() { ... }
// build() 과정에서
if (reader instanceof ItemStream) stream((ItemStream) reader);
if (processor instanceof ItemStream) stream((ItemStream) processor);
if (writer instanceof ItemStream) stream((ItemStream) writer);
.reader(...) / .processor(...) / .writer(...) 로 직접 넘긴 객체가 ItemStream 이면 자동으로 등록합니다.
문제는 직접 넘기지 않았을 때입니다.
8-8. 함정 — 파일이 0바이트가 되는 이유
8-6 의 ClassifierCompositeItemWriter 를 다시 봅니다.
.writer(routingWriter) // ← StepBuilder 가 받은 것은 이것뿐
StepBuilder 는 routingWriter 만 봅니다. 그리고 ClassifierCompositeItemWriter 는 ItemStream 을 구현하지 않습니다.
public class ClassifierCompositeItemWriter<T> implements ItemWriter<T> { // ItemStream 없음
private Classifier<T, ItemWriter<? super T>> classifier;
...
}
그래서 stream() 에 아무것도 등록되지 않고, 안에 숨어 있는 vipFileWriter 의 open() / update() / close() 는 한 번도 호출되지 않습니다.
open() 이 없는데 write() 가 왜 예외를 안 냈을까요? 냈어야 정상입니다. 확인해 봅니다.
// FlatFileItemWriter.write() 내부
if (!getOutputState().isInitialized()) {
throw new WriterNotOpenException("Writer must be open before it can be written to");
}
실제로 이렇게 나옵니다.
결과
ERROR 63455 --- [ main] o.s.batch.core.step.AbstractStep : Encountered an error executing step routingStep in job fileJob
org.springframework.batch.item.WriterNotOpenException: Writer must be open before it can be written to
at org.springframework.batch.item.support.AbstractFileItemWriter.write(AbstractFileItemWriter.java:249)
at org.springframework.batch.item.support.ClassifierCompositeItemWriter.write(ClassifierCompositeItemWriter.java:66)
이 경우는 시끄럽게 실패했으니 다행입니다. 그런데 8-6 에서는 예외 없이 0줄짜리 파일이 나왔습니다. 차이가 뭘까요?
vipFileWriter 를 @Bean 으로 등록했기 때문입니다. FlatFileItemWriter 는 InitializingBean 이므로 스프링이 afterPropertiesSet() 을 불러 주고, 그 과정에서 shouldDeleteIfExists(true) 에 따라 빈 파일이 만들어집니다. 그 뒤 open() 이 안 불리면...
경우가 갈립니다. 그래서 표로 못 박습니다.
| 구성 | open/update/close | 증상 |
|---|
.writer(flatFileWriter) 직접 | ✅ 자동 등록 | 정상 |
.writer(compositeItemWriter) | ✅ CompositeItemWriter 가 ItemStream 이라 델리게이트에 전파 | 정상 |
.writer(classifierCompositeItemWriter) | ❌ ItemStream 아님 | 파일 0바이트 / WriterNotOpenException |
| 프로세서·리스너 안에 숨긴 writer | ❌ | 파일 0바이트 |
@StepScope 프록시 뒤의 writer | ⚠️ 반환 타입이 ItemWriter<T> 면 instanceof ItemStream 이 실패 | 파일 0바이트 |
⚠️ 함정 — CompositeItemWriter 는 전파하고 ClassifierCompositeItemWriter 는 전파하지 않습니다
이름이 비슷하고 역할도 비슷한 두 클래스가 정반대로 동작합니다. CompositeItemWriter 는 ItemStreamWriter 를 구현해 open/update/close 를 모든 델리게이트에 넘깁니다. ClassifierCompositeItemWriter 는 그냥 ItemWriter 입니다.
그래서 "CompositeItemWriter 로 잘 되던 코드"를 분기가 필요해져 ClassifierCompositeItemWriter 로 바꾸는 순간 파일이 조용히 비기 시작합니다. 코드 리뷰에서는 클래스 이름 한 단어만 바뀐 것으로 보입니다.
재시작 쪽 피해는 더 깊습니다. update() 가 안 불리면 ExecutionContext 에 written 위치가 남지 않습니다. 실패 후 재시작하면 writer 는 파일 처음부터 다시 씁니다. 리더는 중단 지점부터 이어가는데 파일은 처음부터라, 파일에는 뒤쪽 데이터만 남습니다. DB 는 맞고 파일은 틀린 상태입니다.
해법: .stream() 으로 수동 등록합니다.
@Bean
public Step routingStep(JobRepository jobRepository, PlatformTransactionManager tx,
ItemReader<Order> allOrdersReader,
ItemProcessor<Order, Settlement> settlementPipeline,
ItemWriter<Settlement> routingWriter,
FlatFileItemWriter<Settlement> vipFileWriter,
JdbcBatchItemWriter<Settlement> settlementJdbcWriter) {
return new StepBuilder("routingStep", jobRepository)
.<Order, Settlement>chunk(1000, tx)
.reader(allOrdersReader)
.processor(settlementPipeline)
.writer(routingWriter)
.stream(vipFileWriter) // ← 숨어 있는 ItemStream 을 손으로 등록
.build();
}
JdbcBatchItemWriter 는 ItemStream 이 아니므로 등록할 필요가 없습니다(파일 핸들 같은 상태가 없습니다).
다시 실행합니다.
결과
INFO 63980 --- [ main] o.s.batch.core.job.SimpleStepHandler : Executing step: [routingStep]
INFO 63980 --- [ main] c.e.batch.step08.CountLogger : readCount=100000, writeCount=70000, filterCount=30000, skipCount=0
INFO 63980 --- [ main] o.s.batch.core.step.AbstractStep : Step: [routingStep] executed in 8s119ms
INFO 63980 --- [ main] o.s.b.c.l.s.TaskExecutorJobLauncher : Job: [SimpleJob: [name=fileJob]] completed with the following parameters: [...] and the following status: [COMPLETED] in 8s344ms
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -N -e "SELECT COUNT(*) FROM settlement;"
wc -l out/vip-settlement.csv
head -2 out/vip-settlement.csv
tail -1 out/vip-settlement.csv
결과
55000
15002 out/vip-settlement.csv
order_id,customer_id,settle_date,gross_amount,fee_rate,fee_amount,net_amount
3,3,2025-01-01,1300.00,0.0150,19.50,1280.50
# vip settlements: generated by fileJob
55,000 + 15,000 = 70,000. 헤더와 푸터도 살아났습니다. .stream() 한 줄이 만든 차이입니다.
등록이 제대로 됐는지는 ExecutionContext 로 확인하는 것이 가장 확실합니다.
mysql -h127.0.0.1 -P3308 -ubatch -pbatch1234 batchdb -e "
SELECT SHORT_CONTEXT FROM BATCH_STEP_EXECUTION_CONTEXT
ORDER BY STEP_EXECUTION_ID DESC LIMIT 1\G"
결과
*************************** 1. row ***************************
SHORT_CONTEXT: {"@class":"java.util.HashMap","vipFileWriter.written":15000,"allOrdersReader.read.count":100000}
vipFileWriter.written 키가 보이면 update() 가 불렸다는 증거입니다. 이 키가 없으면 등록이 안 된 것이고, 재시작이 깨져 있는 상태입니다.
💡 실무 팁 — "파일 writer 를 만들면 .stream() 을 의심한다"를 습관으로
점검 절차 세 줄로 압축합니다.
- Step 을 만들 때
ItemStream 구현체가 .reader()/.writer() 에 직접 들어갔는지 본다.
- 아니라면(합성·분기·프록시 뒤에 있다면)
.stream(...) 으로 손수 등록한다.
- 한 번 돌리고
BATCH_STEP_EXECUTION_CONTEXT 에 그 writer 의 .written 키가 있는지 확인한다.
3번을 하는 데 10초 걸립니다. 안 하면 장애가 나기 전까지, 즉 재시작이 필요해지는 그날까지 아무도 모릅니다.
8-9. 커스텀 ItemWriter 를 직접 만들 때
프레임워크가 제공하지 않는 대상(외부 API, 큐 등)에 쓸 때는 직접 구현합니다. 지켜야 할 규칙이 셋 있습니다.
public class SettlementApiWriter implements ItemStreamWriter<Settlement> {
private final RestClient restClient;
private long sent = 0;
public SettlementApiWriter(RestClient restClient) {
this.restClient = restClient;
}
@Override
public void open(ExecutionContext ctx) {
// (1) 재시작이면 이전 위치를 복원한다
this.sent = ctx.getLong("settlementApiWriter.sent", 0L);
}
@Override
public void write(Chunk<? extends Settlement> chunk) {
// (2) 아이템 하나씩이 아니라 청크 통째로 한 번에 보낸다
restClient.post().uri("/settlements/bulk")
.body(chunk.getItems())
.retrieve().toBodilessEntity();
sent += chunk.size();
}
@Override
public void update(ExecutionContext ctx) {
// (3) 청크 커밋마다 위치를 남긴다
ctx.putLong("settlementApiWriter.sent", sent);
}
@Override
public void close() { }
}
| 규칙 | 이유 |
|---|
write() 안에서 아이템 단위 루프로 원격 호출 금지 | 청크 크기만큼 왕복이 늘어납니다. 8-3 의 8배 차이와 같은 문제입니다 |
open() 에서 복원, update() 에서 저장 | 이걸 빼면 재시작 시 중복 전송됩니다 |
| writer 는 멱등하게 | 청크가 롤백되면 같은 청크가 다시 옵니다. 외부 시스템에는 같은 데이터가 두 번 갑니다 |
⚠️ 함정 — 외부 시스템 쓰기는 트랜잭션에 참여하지 않습니다
DB writer 는 청크가 롤백되면 함께 롤백됩니다. REST 호출은 롤백되지 않습니다. 이미 보낸 것은 보낸 것입니다.
청크의 900번째 아이템에서 예외가 나면 DB 는 깨끗이 되돌아가지만 API 쪽에는 이미 1,000건이 들어가 있습니다. 재시도하면 2,000건이 됩니다.
방어책은 멱등 키입니다. order_id 처럼 아이템마다 고유한 값을 요청에 실어 보내고, 수신 측이 중복을 무시하게 만드세요. 이 코스의 settlement.order_id 에 걸린 UNIQUE 제약이 바로 그 역할을 하는 DB 버전입니다.
정리
| 개념 | 핵심 |
|---|
| 5.x 시그니처 | void write(Chunk<? extends T> chunk). 4.x 의 List<T> 에서 깨지는 변경 |
Chunk | Iterable. getItems() / size() / getSkips() / getErrors() |
beanMapped() | record 에 통하지 않음. nullable 테이블이면 전 컬럼 NULL 로 조용히 성공 |
itemSqlParameterSourceProvider | 람다로 직접 작성. 컴파일 타임에 검증됨 — record 의 정답 |
rewriteBatchedStatements=true | 70,000건 쓰기 61.4초 → 7.7초, 약 8배. Com_insert 70000 → 70 |
| 진단법 | Com_insert 가 처리 건수와 같으면 배치 쓰기가 배치가 아님 |
assertUpdates | 갱신 건수 0 만 잡음. -2(SUCCESS_NO_INFO)와 "과다 갱신"은 못 잡음 |
FlatFileItemWriter | DelimitedLineAggregator + 람다 FieldExtractor(record 대응) |
headerCallback | open() 에서 실행 |
footerCallback | close() 에서 실행 → close 가 안 불리면 푸터도 버퍼도 사라짐 |
CompositeItemWriter | 모든 델리게이트에 같은 청크. ItemStream 전파함 |
ClassifierCompositeItemWriter | 아이템별 분기. ItemStream 전파 안 함 |
ItemStream | open(1회) / update(청크마다) / close(1회) |
| 자동 등록 | .reader()/.processor()/.writer() 에 직접 넘긴 것만 |
.stream(writer) | 합성·분기·프록시 뒤에 숨은 ItemStream 을 수동 등록 |
| 등록 검증 | BATCH_STEP_EXECUTION_CONTEXT 에 <name>.written 키가 있는가 |
| 외부 시스템 writer | 트랜잭션에 참여하지 않음 → 멱등 키 필수 |
연습문제
Exercise.java 에 7문제가 있습니다. 정답은 Solution.java.
- 4.x 스타일
write(List<? extends T>) writer 를 5.x 시그니처로 마이그레이션하기
record 도메인에 맞는 JdbcBatchItemWriter 를 itemSqlParameterSourceProvider 로 완성하고, beanMapped() 버전이 왜 실패하는지 예외 메시지를 예측하기
rewriteBatchedStatements 를 껐다 켜며 70,000건 쓰기를 실측하고, Com_insert 값으로 왕복 횟수를 증명하기
assertUpdates(true) 가 잡지 못하는 두 가지 상황을 코드로 만들고, 각각의 대안 방어책을 적기
headerCallback / footerCallback 를 갖춘 FlatFileItemWriter 를 만들고, 출력이 70,002줄인 이유를 설명하기
ClassifierCompositeItemWriter 로 VIP 15,000건만 파일로 분기하되, .stream() 없이 먼저 실행해 실패를 관찰한 뒤 고치기
- 커스텀
ItemStreamWriter 를 만들고, open/update 를 뺐을 때 재시작에서 무엇이 깨지는지 답하기
다음 단계
8-7 에서 update(ExecutionContext) 가 청크 커밋마다 위치를 저장한다는 것을 봤습니다.
그 ExecutionContext 가 정확히 무엇이고, 어디에 어떤 형태로 저장되며, 재시작 시 어떻게 복원되는지 — 그리고 Job 레벨 컨텍스트와 Step 레벨 컨텍스트의 차이는 무엇인지를 다음 스텝에서 파고듭니다.
Step 07 에서 미뤄 둔 "누적 합계를 재시작에도 이어가는 법"도 여기서 답이 나옵니다.
→ Step 09 — ExecutionContext
실습 파일
이 스텝은 Java 파일 세 개로 진행합니다. Practice.java 로 8-1 ~ 8-9 의 writer 구현을 모두 훑고, Exercise.java 의 7문제를 채운 뒤, Solution.java 로 대조합니다. 세 파일 모두 com.example.batch.step08 패키지이고, 예제는 static class 중첩으로 담았습니다. 실행 전에 반드시 mkdir -p ./out 을 해 두세요 — FlatFileItemWriter 는 상위 디렉터리를 만들어 주지 않습니다.
Practice.java
본문의 모든 writer 와 Step 설정을 절 번호 주석(// [8-3])과 함께 모아 둔 참조 코드입니다.
Practice 는 @Configuration 이며 settlementJob(DB 만) / fileJob(파일 포함) / routingJob(분기) 세 Job 을 정의합니다. --spring.batch.job.name= 으로 고릅니다.
[8-2] 의 brokenBeanMappedWriter() 는 record 에 beanMapped() 를 쓴 실패 예시입니다. @Bean 이 아니므로 그냥 빌드해도 안전하고, 재현하려면 settlementJdbcWriter 자리에 손으로 바꿔 끼워야 합니다. 이때 나오는 InvalidDataAccessApiUsageException 이 settlement 테이블이 NOT NULL 이라서 나는 것이라는 점을 주석이 짚어 둡니다. nullable 테이블이었다면 조용히 성공했을 것입니다.
[8-3] 의 8배 실측은 코드가 아니라 application.yml 의 URL 을 고쳐야 재현됩니다. 파일 상단 주석에 두 URL 이 그대로 적혀 있으니 복사해 쓰세요. 측정 전마다 TRUNCATE TABLE settlement; 를 잊지 마세요 — 안 그러면 DuplicateKeyException 이 먼저 납니다(uk_settlement_order 제약).
[8-6]/[8-8] 의 routingStep 에는 .stream(vipFileWriter) 줄이 주석 처리된 채로 들어 있습니다. 먼저 주석인 상태로 돌려 0줄 파일(또는 WriterNotOpenException)을 관찰하고, 그다음 주석을 풀어 15,002줄이 되는 것을 확인하는 순서로 만들어 두었습니다. 순서를 지키세요. 고쳐진 코드만 보면 이 함정은 배워지지 않습니다.
[8-9] 의 SettlementApiWriter 는 실제 엔드포인트가 없으므로 그대로 실행하면 커넥션 예외가 납니다. 구조를 읽는 용도이며, open/update 에서 ExecutionContext 를 다루는 부분이 Step 09 의 예고편입니다.
package com.example.batch.step08;
import com.example.batch.domain.Order;
import com.example.batch.domain.Settlement;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.StepExecutionListener;
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.item.Chunk;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemStreamWriter;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.database.JdbcBatchItemWriter;
import org.springframework.batch.item.database.builder.JdbcBatchItemWriterBuilder;
import org.springframework.batch.item.database.builder.JdbcCursorItemReaderBuilder;
import org.springframework.batch.item.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.batch.item.file.transform.DelimitedLineAggregator;
import org.springframework.batch.item.support.ClassifierCompositeItemWriter;
import org.springframework.batch.item.support.CompositeItemWriter;
import org.springframework.classify.Classifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.io.FileSystemResource;
import org.springframework.jdbc.core.namedparam.MapSqlParameterSource;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.web.client.RestClient;
import javax.sql.DataSource;
import java.math.BigDecimal;
import java.util.List;
/**
* Step 08 — ItemWriter 실습 전체 코드.
*
* 실행 전 준비
* mkdir -p ./out
* mysql -h127.0.0.1 -P3308 -uroot -proot1234 batchdb -e "TRUNCATE TABLE settlement;"
*
* 실행
* ./gradlew bootRun --args='--spring.batch.job.name=settlementJob' # DB 만
* ./gradlew bootRun --args='--spring.batch.job.name=fileJob' # DB + 파일
* ./gradlew bootRun --args='--spring.batch.job.name=routingJob' # 분기
*
* ─────────────────────────────────────────────────────────────────────
* [8-3] 8배 실측을 재현하려면 application.yml 의 URL 을 아래 두 개로 번갈아 두세요.
*
* (A) 옵션 없음 — 70,000건 쓰기에 약 61.4초, Com_insert = 70000
* url: jdbc:mysql://127.0.0.1:3308/batchdb?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Seoul
*
* (B) 옵션 있음 — 약 7.7초, Com_insert = 70
* url: jdbc:mysql://127.0.0.1:3308/batchdb?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/Seoul&rewriteBatchedStatements=true
*
* 측정 전마다 TRUNCATE TABLE settlement; 를 하세요.
* 안 하면 uk_settlement_order 제약 때문에 DuplicateKeyException 이 먼저 납니다.
* ─────────────────────────────────────────────────────────────────────
*/
@Configuration
public class Practice {
private static final Logger log = LoggerFactory.getLogger(Practice.class);
// ======================================================================
// [8-1] 5.x 시그니처 — void write(Chunk<? extends T>)
// ======================================================================
/**
* 4.x 는 write(List<? extends T> items) 였습니다. 5.0 에서 Chunk 로 바뀌었고
* 이것은 소스 호환성이 깨지는 변경입니다. 4.x 코드를 붙여 넣으면
* error: ... does not override abstract method write(Chunk<? extends Settlement>)
* 로 컴파일이 실패합니다. 컴파일러가 잡아 주는 "좋은 실패" 입니다.
*/
public static class SettlementLogWriter implements ItemWriter<Settlement> {
@Override
public void write(Chunk<? extends Settlement> chunk) {
log.debug("chunk size={} first={} last={}",
chunk.size(),
chunk.getItems().get(0).orderId(),
chunk.getItems().get(chunk.size() - 1).orderId());
}
}
/** 4.x 코드를 옮길 때의 최소 어댑터. 본문은 그대로 두고 getItems() 만 끼웁니다. */
public static class MigratedWriter implements ItemWriter<Settlement> {
@Override
public void write(Chunk<? extends Settlement> chunk) {
doWrite(chunk.getItems());
}
/** 4.x 시절 그대로의 본문. */
private void doWrite(List<? extends Settlement> items) {
log.debug("{} items", items.size());
}
}
/** ItemWriter 는 5.x 에서 @FunctionalInterface 입니다. 람다로도 됩니다. */
static ItemWriter<Settlement> lambdaWriter() {
return chunk -> chunk.forEach(s -> log.trace("{}", s.orderId()));
}
// ======================================================================
// [8-2] JdbcBatchItemWriter — record 와 beanMapped()
// ======================================================================
/**
* ❌ record 에는 통하지 않습니다.
*
* beanMapped() 는 BeanPropertyItemSqlParameterSourceProvider 를 쓰고,
* 이것은 자바빈 규약(getOrderId())을 기대합니다. record 의 접근자는 orderId() 입니다.
*
* 실행 결과:
* org.springframework.dao.InvalidDataAccessApiUsageException:
* No value supplied for the SQL parameter 'orderId': Invalid property 'orderId'
* of bean class [com.example.batch.domain.Settlement]
*
* ⚠️ 이 예외가 나는 것은 settlement 의 컬럼들이 NOT NULL 이기 때문입니다.
* nullable 테이블이었다면 예외 없이 전 컬럼 NULL 인 70,000행이 들어가고
* writeCount=70000, STATUS=COMPLETED 로 모든 지표가 정상으로 보입니다.
* 결과 테이블을 NOT NULL 로 만드는 것은 "조용한 실패를 시끄럽게 만드는 장치" 입니다.
*
* @Bean 이 아니므로 그냥 빌드해도 안전합니다.
* 재현하려면 settlementJdbcWriter 자리에 손으로 바꿔 끼우세요.
*/
public JdbcBatchItemWriter<Settlement> brokenBeanMappedWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<Settlement>()
.dataSource(dataSource)
.sql(INSERT_SQL)
.beanMapped()
.build();
}
static final String INSERT_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)
""";
/**
* ✅ record 의 정답. 장황하지만 컴파일 타임에 검증됩니다.
* record 의 컴포넌트 이름을 바꾸면 여기서 컴파일 에러가 납니다.
* beanMapped() 는 리팩터링을 따라오지 못합니다.
*/
@Bean
public JdbcBatchItemWriter<Settlement> settlementJdbcWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<Settlement>()
.dataSource(dataSource)
.sql(INSERT_SQL)
.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()))
// [8-4] 기본값 true. UPDATE writer 에서 "0건 갱신"을 잡아 줍니다.
.assertUpdates(true)
.build();
}
// ======================================================================
// [8-4] assertUpdates — 무엇을 잡고 무엇을 못 잡는가
// ======================================================================
/**
* UPDATE writer. WHERE 에 걸리는 행이 없으면 갱신 건수가 0 이고,
* assertUpdates(true) 가 EmptyResultDataAccessException 을 던집니다.
*
* Item 137 of 1000 did not update any rows: [Settlement[orderId=999999, ...]]
*
* ⚠️ 구멍 두 개
* (1) rewriteBatchedStatements=true 로 UPDATE 배치가 다중 문장으로 재작성되면
* 드라이버가 Statement.SUCCESS_NO_INFO(-2) 를 돌려줍니다.
* JdbcBatchItemWriter 의 검사는 `updateCounts[i] == 0` 뿐이라 -2 는 통과합니다.
* 성능 옵션을 켜는 순간 안전장치가 조용히 꺼집니다.
* (2) "1건 기대했는데 70,000건 갱신됨"은 잡지 못합니다.
* WHERE 절을 빠뜨린 UPDATE 는 갱신 건수가 0 이 아니기 때문입니다.
* 옵션으로 못 막습니다. UPDATE 의 WHERE 에 유니크 키가 있는지 리뷰하세요.
*/
public JdbcBatchItemWriter<Settlement> settlementUpdateWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<Settlement>()
.dataSource(dataSource)
.sql("UPDATE settlement SET net_amount = :netAmount WHERE order_id = :orderId")
.itemSqlParameterSourceProvider(s -> new MapSqlParameterSource()
.addValue("orderId", s.orderId())
.addValue("netAmount", s.netAmount()))
.assertUpdates(true)
.build();
}
// ======================================================================
// [8-5] FlatFileItemWriter — CSV
// ======================================================================
/**
* BeanWrapperFieldExtractor 는 beanMapped() 와 같은 이유로 record 에 통하지 않습니다.
* FieldExtractor 를 람다로 씁니다.
*
* headerCallback 은 open() 에서, footerCallback 은 close() 에서 실행됩니다([8-7]).
* 이 사실이 [8-8] 함정의 뿌리입니다.
*/
@Bean
public FlatFileItemWriter<Settlement> settlementFileWriter() {
return csvWriter("settlementFileWriter", "out/settlement.csv",
"# generated by settlementJob");
}
@Bean
public FlatFileItemWriter<Settlement> vipFileWriter() {
return csvWriter("vipFileWriter", "out/vip-settlement.csv",
"# vip settlements: generated by fileJob");
}
private static FlatFileItemWriter<Settlement> csvWriter(String name, String path, String footer) {
DelimitedLineAggregator<Settlement> aggregator = new DelimitedLineAggregator<>();
aggregator.setDelimiter(",");
aggregator.setFieldExtractor(s -> new Object[]{
s.orderId(), s.customerId(), s.settleDate(),
s.grossAmount(), s.feeRate(), s.feeAmount(), s.netAmount()
});
return new FlatFileItemWriterBuilder<Settlement>()
.name(name) // ExecutionContext 키의 접두사가 됩니다
.resource(new FileSystemResource(path))
.encoding("UTF-8") // 플랫폼 기본값에 맡기지 마세요
.lineSeparator("\n") // 윈도우에서 \r 이 섞이는 사고 방지
.lineAggregator(aggregator)
.headerCallback(w -> w.write(
"order_id,customer_id,settle_date,gross_amount,"
+ "fee_rate,fee_amount,net_amount"))
.footerCallback(w -> w.write(footer))
.shouldDeleteIfExists(true)
.transactional(true) // 기본값. 청크 커밋 시점에 flush
.build();
}
// ======================================================================
// [8-6] CompositeItemWriter — 같은 청크를 여러 곳에
// ======================================================================
/**
* 모든 델리게이트가 같은 청크를 받습니다.
* writeCount 는 70000 한 번만 올라갑니다. 두 곳에 썼다고 140000 이 되지 않습니다.
*
* ✅ CompositeItemWriter 는 ItemStreamWriter 를 구현하므로
* open/update/close 를 모든 델리게이트에 전파합니다. .stream() 이 필요 없습니다.
*/
@Bean
public ItemWriter<Settlement> dualWriter(JdbcBatchItemWriter<Settlement> settlementJdbcWriter,
FlatFileItemWriter<Settlement> settlementFileWriter) {
CompositeItemWriter<Settlement> writer = new CompositeItemWriter<>();
writer.setDelegates(List.of(settlementJdbcWriter, settlementFileWriter));
return writer;
}
// ======================================================================
// [8-6][8-8] ClassifierCompositeItemWriter — 아이템별 분기
// ======================================================================
/**
* VIP(fee_rate 0.0150) 15,000건은 파일로, 나머지 55,000건은 DB 로.
*
* ❌ ClassifierCompositeItemWriter 는 ItemStream 을 구현하지 않습니다.
* 안에 숨은 vipFileWriter 의 open/update/close 가 한 번도 불리지 않습니다.
* -> 파일 0줄, 또는 WriterNotOpenException.
* -> routingStep 의 .stream(vipFileWriter) 로 손수 등록해야 합니다.
*
* 이름이 비슷한 CompositeItemWriter 는 전파하는데 이쪽은 안 합니다.
* 이 비대칭이 함정의 본질입니다.
*/
@Bean
public ItemWriter<Settlement> routingWriter(JdbcBatchItemWriter<Settlement> settlementJdbcWriter,
FlatFileItemWriter<Settlement> vipFileWriter) {
ClassifierCompositeItemWriter<Settlement> writer = new ClassifierCompositeItemWriter<>();
BigDecimal vipRate = new BigDecimal("0.0150");
Classifier<Settlement, ItemWriter<? super Settlement>> classifier =
s -> vipRate.compareTo(s.feeRate()) == 0 ? vipFileWriter : settlementJdbcWriter;
writer.setClassifier(classifier);
return writer;
}
// ======================================================================
// [8-9] 커스텀 ItemStreamWriter
// ======================================================================
/**
* 외부 API 로 쓰는 writer. 실제 엔드포인트가 없으므로 그대로 실행하면
* 커넥션 예외가 납니다. 구조를 읽는 용도입니다.
*
* 규칙 셋
* (1) write() 안에서 아이템 단위 루프로 원격 호출 금지 — [8-3] 의 8배 문제와 동형입니다.
* (2) open() 에서 복원하고 update() 에서 저장 — 안 하면 재시작 시 중복 전송됩니다.
* (3) 멱등하게 — 청크가 롤백되면 같은 청크가 다시 옵니다.
* REST 호출은 롤백되지 않습니다. 이미 보낸 것은 보낸 것입니다.
* order_id 같은 멱등 키를 실어 보내고 수신 측이 중복을 무시하게 하세요.
*/
public static class SettlementApiWriter implements ItemStreamWriter<Settlement> {
private static final String CTX_KEY = "settlementApiWriter.sent";
private final RestClient restClient;
private long sent = 0;
public SettlementApiWriter(RestClient restClient) {
this.restClient = restClient;
}
@Override
public void open(ExecutionContext ctx) {
this.sent = ctx.getLong(CTX_KEY, 0L); // (2) 재시작 복원
}
@Override
public void write(Chunk<? extends Settlement> chunk) {
restClient.post().uri("/settlements/bulk") // (1) 청크 통째로 한 번에
.body(chunk.getItems())
.retrieve().toBodilessEntity();
sent += chunk.size();
}
@Override
public void update(ExecutionContext ctx) {
ctx.putLong(CTX_KEY, sent); // (2) 청크 커밋마다 저장
}
@Override
public void close() {
}
}
// ======================================================================
// Reader / Processor / Listener
// ======================================================================
@Bean
public ItemReader<Order> allOrdersReader(DataSource dataSource) {
return new JdbcCursorItemReaderBuilder<Order>()
.name("allOrdersReader")
.dataSource(dataSource)
.sql("SELECT order_id, customer_id, amount, status, ordered_at "
+ "FROM orders ORDER BY order_id")
.rowMapper((rs, rowNum) -> new Order(
rs.getLong("order_id"),
rs.getInt("customer_id"),
rs.getBigDecimal("amount"),
rs.getString("status"),
rs.getTimestamp("ordered_at").toLocalDateTime()))
.fetchSize(1000)
.build();
}
/** Step 07 의 파이프라인을 그대로 씁니다. 10만 건 읽고 3만 건 필터 -> 7만 건 출력. */
@Bean
public ItemProcessor<Order, Settlement> settlementPipeline(
org.springframework.jdbc.core.JdbcTemplate jdbcTemplate) {
var step1 = new com.example.batch.step07.Practice.CompletedOnlyProcessor();
var step2 = new com.example.batch.step07.Practice.GradeFeeSettlementProcessor(jdbcTemplate);
return order -> {
Order filtered = step1.process(order);
return filtered == null ? null : step2.process(filtered);
};
}
@Bean
public StepExecutionListener countLogger() {
return new StepExecutionListener() {
@Override
public ExitStatus afterStep(StepExecution se) {
log.info("readCount={}, writeCount={}, filterCount={}, skipCount={}",
se.getReadCount(), se.getWriteCount(),
se.getFilterCount(), se.getSkipCount());
return se.getExitStatus();
}
};
}
// ======================================================================
// Step / Job
// ======================================================================
/** [8-3] 8배 측정용. DB 에만 씁니다. */
@Bean
public Step settlementStep(JobRepository jobRepository, PlatformTransactionManager tx,
ItemReader<Order> allOrdersReader,
ItemProcessor<Order, Settlement> settlementPipeline,
JdbcBatchItemWriter<Settlement> settlementJdbcWriter,
StepExecutionListener countLogger) {
return new StepBuilder("settlementStep", jobRepository)
.<Order, Settlement>chunk(1000, tx)
.reader(allOrdersReader)
.processor(settlementPipeline)
.writer(settlementJdbcWriter)
.listener(countLogger)
.build();
}
/** [8-6] DB + 파일. CompositeItemWriter 라 .stream() 이 필요 없습니다. */
@Bean
public Step dualStep(JobRepository jobRepository, PlatformTransactionManager tx,
ItemReader<Order> allOrdersReader,
ItemProcessor<Order, Settlement> settlementPipeline,
ItemWriter<Settlement> dualWriter,
StepExecutionListener countLogger) {
return new StepBuilder("dualStep", jobRepository)
.<Order, Settlement>chunk(1000, tx)
.reader(allOrdersReader)
.processor(settlementPipeline)
.writer(dualWriter)
.listener(countLogger)
.build();
}
/**
* [8-8] 함정 재현 스텝.
*
* 순서를 지키세요.
* 1. .stream(vipFileWriter) 가 주석인 채로 실행 -> 파일 0줄 / WriterNotOpenException
* 그리고 BATCH_STEP_EXECUTION_CONTEXT 에 vipFileWriter.written 키가 없는 것 확인
* 2. 주석을 풀고 다시 실행 -> 15,002줄, 컨텍스트에 vipFileWriter.written=15000
*
* 고쳐진 코드만 보면 이 함정은 배워지지 않습니다.
*/
@Bean
public Step routingStep(JobRepository jobRepository, PlatformTransactionManager tx,
ItemReader<Order> allOrdersReader,
ItemProcessor<Order, Settlement> settlementPipeline,
ItemWriter<Settlement> routingWriter,
FlatFileItemWriter<Settlement> vipFileWriter,
StepExecutionListener countLogger) {
return new StepBuilder("routingStep", jobRepository)
.<Order, Settlement>chunk(1000, tx)
.reader(allOrdersReader)
.processor(settlementPipeline)
.writer(routingWriter)
// .stream(vipFileWriter) // ← 1단계에서는 주석, 2단계에서 해제
.listener(countLogger)
.build();
}
@Bean
public Job settlementJob(JobRepository jobRepository, Step settlementStep) {
return new JobBuilder("settlementJob", jobRepository).start(settlementStep).build();
}
@Bean
public Job fileJob(JobRepository jobRepository, Step dualStep) {
return new JobBuilder("fileJob", jobRepository).start(dualStep).build();
}
@Bean
public Job routingJob(JobRepository jobRepository, Step routingStep) {
return new JobBuilder("routingJob", jobRepository).start(routingStep).build();
}
}
Exercise.java
7문제의 문제지입니다. // 여기에 작성: 자리를 채우세요.
- 문제 2·3·4·6 은 "먼저 예측을 주석으로 적고 → 실행해서 대조"하는 형식입니다. 예측을 건너뛰면 문제의 절반이 사라집니다.
- 문제 3 은 코드가 아니라 측정 절차를 수행하는 문제입니다.
FLUSH STATUS; → Job 실행 → SHOW GLOBAL STATUS LIKE 'Com_insert'; 순서로 두 번(옵션 off/on) 재고, 결과를 표에 적습니다. Com_insert 가 70000 과 70 으로 갈리는 것을 직접 보는 것이 목적입니다.
- 문제 4 의
Q4_OverUpdate 는 WHERE 절이 빠진 UPDATE 를 일부러 담고 있습니다. 실행하면 settlement 70,000행이 전부 덮이므로, 반드시 TRUNCATE 후 재생성할 각오로 돌리세요. 이것이 assertUpdates 가 못 잡는 사고의 실체입니다.
- 문제 6 은
.stream() 이 빠진 채 주어집니다. 이 파일을 그대로 실행하면 파일이 비거나 WriterNotOpenException 이 나는 것이 정상입니다. 그 상태를 확인한 뒤 고치세요.
- 문제 7 은 실행 없이 답만 적는 서술형입니다.
open 을 뺐을 때와 update 를 뺐을 때 깨지는 것이 서로 다르다는 점이 핵심입니다.
package com.example.batch.step08;
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.item.Chunk;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemStreamWriter;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.database.JdbcBatchItemWriter;
import org.springframework.batch.item.database.builder.JdbcBatchItemWriterBuilder;
import org.springframework.batch.item.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.batch.item.file.transform.DelimitedLineAggregator;
import org.springframework.batch.item.support.ClassifierCompositeItemWriter;
import org.springframework.core.io.FileSystemResource;
import org.springframework.jdbc.core.namedparam.MapSqlParameterSource;
import org.springframework.transaction.PlatformTransactionManager;
import javax.sql.DataSource;
import java.util.List;
/**
* Step 08 — 연습문제 7문항.
*
* 준비
* mkdir -p ./out
* mysql -h127.0.0.1 -P3308 -uroot -proot1234 batchdb -e "TRUNCATE TABLE settlement;"
*
* 규칙
* - "여기에 작성:" 자리를 채우세요.
* - 문제 2·3·4·6 은 반드시 "예측 먼저, 실행 나중" 입니다.
*
* 정답은 Solution.java.
*/
public class Exercise {
// ======================================================================
// 문제 1 — 4.x writer 를 5.x 시그니처로 마이그레이션
// ======================================================================
/**
* 아래는 Spring Batch 4.x 스타일 writer 입니다. 그대로 두면 컴파일되지 않습니다.
* doWrite(List) 본문은 건드리지 말고, 시그니처만 5.x 로 바꾸세요.
*
* 컴파일 에러 메시지를 먼저 예측해 보세요.
* 여기에 작성(예상 에러):
*/
public static class Q1_LegacyWriter implements ItemWriter<Settlement> {
// 여기에 작성: 5.x 시그니처로 write() 를 구현하세요.
/** 4.x 시절 본문. 수정하지 마세요. */
private void doWrite(List<? extends Settlement> items) {
System.out.println("wrote " + items.size() + " items");
}
}
// ======================================================================
// 문제 2 — record 에 맞는 JdbcBatchItemWriter
// ======================================================================
/**
* (a) 아래 writer 에서 .beanMapped() 를 쓰면 어떤 예외가 납니까?
* 예외 클래스 : ______________________
* 메시지 요지 : ______________________
* 그 이유 : 여기에 작성:
*
* (b) settlement 테이블의 컬럼이 전부 nullable 이었다면 무슨 일이 벌어집니까?
* writeCount = ______, Job STATUS = ______, 데이터 = ______
* 여기에 작성(설명):
*
* (c) record 에 맞는 writer 를 완성하세요.
*/
public static class Q2_RecordWriter {
public JdbcBatchItemWriter<Settlement> build(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 를 람다로 채우세요.
.assertUpdates(true)
.build();
}
}
// ======================================================================
// 문제 3 — rewriteBatchedStatements 실측 (측정 절차 문제)
// ======================================================================
/**
* 코드가 아니라 절차를 수행하는 문제입니다. 두 번 측정하세요.
*
* 각 측정마다:
* 1) application.yml 의 datasource.url 을 (A) 또는 (B) 로 바꾼다
* (A) ...serverTimezone=Asia/Seoul
* (B) ...serverTimezone=Asia/Seoul&rewriteBatchedStatements=true
* 2) mysql ... -e "TRUNCATE TABLE settlement;"
* 3) mysql ... -uroot -proot1234 -e "FLUSH STATUS;"
* 4) ./gradlew bootRun --args='--spring.batch.job.name=settlementJob'
* 5) mysql ... -uroot -proot1234 -e "SHOW GLOBAL STATUS LIKE 'Com_insert';"
*
* 결과를 채우세요.
* ┌───────────────┬──────────────┬─────────────┬────────────┐
* │ │ Step 실행시간 │ Com_insert │ writeCount │
* ├───────────────┼──────────────┼─────────────┼────────────┤
* │ (A) 옵션 없음 │ │ │ │
* │ (B) 옵션 있음 │ │ │ │
* └───────────────┴──────────────┴─────────────┴────────────┘
* 배수: 약 ______ 배
*
* (a) 두 실행의 writeCount / Job STATUS / 결과 행 수가 같습니까? 답:
* 여기에 작성:
* (b) 그렇다면 이 문제를 "모니터링만으로" 발견할 수 있습니까? 답과 대안:
* 여기에 작성:
*/
// ======================================================================
// 문제 4 — assertUpdates 가 못 잡는 두 상황
// ======================================================================
/**
* ⚠️ 이 문제는 settlement 데이터를 파괴합니다. 실행 후 TRUNCATE 하고 재생성하세요.
*
* (a) 아래 writer 는 WHERE 절이 없습니다.
* 70,000건 청크를 흘리면 최종적으로 settlement 는 어떤 상태가 됩니까?
* 여기에 작성(예측):
* assertUpdates(true) 가 이것을 잡습니까? 답과 이유:
* 여기에 작성:
*
* (b) rewriteBatchedStatements=true 인 상태에서 UPDATE 배치가 다중 문장으로
* 재작성되면 드라이버는 각 항목에 어떤 값을 돌려줍니까?
* 값 = ______ (상수 이름과 숫자)
* JdbcBatchItemWriter 의 검사는 `updateCounts[i] == 0` 입니다.
* 이 값이 그 검사를 통과합니까? 답과 그 의미:
* 여기에 작성:
*
* (c) 두 상황 각각의 대안 방어책을 적으세요.
* (a) 의 방어책: 여기에 작성:
* (b) 의 방어책: 여기에 작성:
*/
public static class Q4_OverUpdate {
public JdbcBatchItemWriter<Settlement> dangerous(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<Settlement>()
.dataSource(dataSource)
.sql("UPDATE settlement SET fee_rate = :feeRate") // WHERE 가 없다
.itemSqlParameterSourceProvider(s -> new MapSqlParameterSource()
.addValue("feeRate", s.feeRate()))
.assertUpdates(true)
.build();
}
}
// ======================================================================
// 문제 5 — FlatFileItemWriter 와 header / footer
// ======================================================================
/**
* out/exercise-settlement.csv 를 만드세요.
* - 구분자 ","
* - 헤더: order_id,customer_id,settle_date,gross_amount,fee_rate,fee_amount,net_amount
* - 푸터: # total lines above = 70000
* - 인코딩 UTF-8, 줄 구분자 "\n", 기존 파일은 삭제 후 새로 쓰기
*
* 실행 후 `wc -l out/exercise-settlement.csv` 의 결과를 예측하세요.
* 예측 = ______줄
* 그 이유: 여기에 작성:
*
* 힌트: FieldExtractor 로 BeanWrapperFieldExtractor 를 쓰면 안 되는 이유는
* 문제 2 의 beanMapped() 와 완전히 같습니다.
*/
public static class Q5_FileWriter {
public FlatFileItemWriter<Settlement> build() {
DelimitedLineAggregator<Settlement> aggregator = new DelimitedLineAggregator<>();
aggregator.setDelimiter(",");
// 여기에 작성: setFieldExtractor
return new FlatFileItemWriterBuilder<Settlement>()
.name("exerciseFileWriter")
.resource(new FileSystemResource("out/exercise-settlement.csv"))
// 여기에 작성: encoding / lineSeparator / lineAggregator /
// headerCallback / footerCallback / shouldDeleteIfExists
.build();
}
}
// ======================================================================
// 문제 6 — ClassifierCompositeItemWriter 와 .stream()
// ======================================================================
/**
* VIP(fee_rate 0.0150) 15,000건은 파일로, 나머지 55,000건은 DB 로 보냅니다.
*
* 아래 Step 에는 .stream() 이 빠져 있습니다.
*
* 1단계: 이 상태로 그대로 실행하세요. (실패가 정상입니다)
* Job STATUS = ______
* DB 행 수 = ______
* 파일 줄 수 = ______
* 예외가 났다면 클래스 = ______________________
* BATCH_STEP_EXECUTION_CONTEXT 에 vipFileWriter.written 키가 있습니까? ______
*
* 2단계: 왜 이렇게 되는지 적으세요.
* 여기에 작성:
* (힌트: ClassifierCompositeItemWriter 는 무엇을 구현하지 "않는가"?
* CompositeItemWriter 와 비교하세요.)
*
* 3단계: 고치세요.
*/
public static class Q6_RoutingStep {
public Step build(JobRepository jobRepository, PlatformTransactionManager tx,
ItemReader<Order> reader,
ItemProcessor<Order, Settlement> processor,
ClassifierCompositeItemWriter<Settlement> routingWriter,
FlatFileItemWriter<Settlement> vipFileWriter) {
return new StepBuilder("q6RoutingStep", jobRepository)
.<Order, Settlement>chunk(1000, tx)
.reader(reader)
.processor(processor)
.writer(routingWriter)
// 여기에 작성(3단계):
.build();
}
}
// ======================================================================
// 문제 7 — 커스텀 ItemStreamWriter 의 생명주기 (서술형)
// ======================================================================
/**
* 아래 writer 를 완성하고, 다음 질문에 답하세요.
*
* (a) open() 을 비워 두면 무엇이 깨집니까?
* 여기에 작성:
* (b) update() 를 비워 두면 무엇이 깨집니까?
* 여기에 작성:
* (c) (a) 와 (b) 중 어느 쪽이 더 발견하기 어렵습니까? 왜입니까?
* 여기에 작성:
* (d) 청크가 롤백되면 이미 전송한 아이템은 어떻게 됩니까? 방어책은?
* 여기에 작성:
*/
public static class Q7_CustomWriter implements ItemStreamWriter<Settlement> {
private static final String CTX_KEY = "q7Writer.sent";
private long sent = 0;
@Override
public void open(ExecutionContext ctx) {
// 여기에 작성:
}
@Override
public void write(Chunk<? extends Settlement> chunk) {
// 실제 전송 대신 카운트만 합니다.
sent += chunk.size();
}
@Override
public void update(ExecutionContext ctx) {
// 여기에 작성:
}
@Override
public void close() {
}
}
}
Solution.java
7문제의 정답과 "왜 그 답인가"를 설명하는 긴 주석입니다. 풀어 본 뒤에 여세요.
- 정답 2 는
beanMapped() 가 실패하는 이유를 record 접근자 규약(orderId() vs getOrderId())까지 내려가 설명하고, **가장 위험한 시나리오가 "예외가 나는 경우"가 아니라 "nullable 테이블에서 조용히 성공하는 경우"**라는 점을 강조합니다. settlement 를 nullable 로 바꿔 재현해 보는 방법도 주석에 적어 두었습니다.
- 정답 3 은 61.4초 / 7.7초와
Com_insert 70000 / 70 을 표로 제시하고, "왜 8배인가" 를 분해합니다. 네트워크 왕복 1회당 약 0.77ms 라고 보면 69,930회 절약 × 0.77ms ≈ 54초 — 실측 차이 53.7초와 맞아떨어집니다. 병목이 DB 연산이 아니라 왕복이라는 결론이 여기서 나옵니다.
- 정답 4 는 두 구멍을 각각 다룹니다. (a)
SUCCESS_NO_INFO(-2) 는 == 0 검사를 통과한다는 것 — 즉 성능 옵션을 켜면 안전장치가 조용히 꺼진다는 것, (b) 과다 갱신은 애초에 갱신 건수로 판별할 수 없다는 것. 대안으로 Step 뒤에 검증 Tasklet 을 붙이는 코드가 함께 들어 있습니다.
- 정답 6 은
.stream() 이 없을 때 벌어지는 일을 세 단계로 나눕니다: open() 미호출 → 헤더 없음 + 파일 미개방, update() 미호출 → ExecutionContext 에 written 키 없음(재시작 파괴), close() 미호출 → 푸터 없음 + 마지막 버퍼 유실. 그리고 CompositeItemWriter 는 전파하고 ClassifierCompositeItemWriter 는 안 한다는 비대칭이 이 함정의 진짜 원인이라고 결론짓습니다.
- 정답 7 은
open 과 update 를 뺐을 때의 증상이 다르다는 점을 표로 정리합니다. update 만 빠지면 처음 실행은 완벽하게 성공하고 재시작할 때만 중복 전송이 일어납니다 — 즉 테스트 환경에서는 절대 발견되지 않는 종류의 버그입니다.
package com.example.batch.step08;
import com.example.batch.domain.Order;
import com.example.batch.domain.Settlement;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepContribution;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.scope.context.ChunkContext;
import org.springframework.batch.core.step.builder.StepBuilder;
import org.springframework.batch.core.step.tasklet.Tasklet;
import org.springframework.batch.repeat.RepeatStatus;
import org.springframework.batch.item.Chunk;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.batch.item.ItemProcessor;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.ItemStreamWriter;
import org.springframework.batch.item.ItemWriter;
import org.springframework.batch.item.database.JdbcBatchItemWriter;
import org.springframework.batch.item.database.builder.JdbcBatchItemWriterBuilder;
import org.springframework.batch.item.file.FlatFileItemWriter;
import org.springframework.batch.item.file.builder.FlatFileItemWriterBuilder;
import org.springframework.batch.item.file.transform.DelimitedLineAggregator;
import org.springframework.batch.item.support.ClassifierCompositeItemWriter;
import org.springframework.core.io.FileSystemResource;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.namedparam.MapSqlParameterSource;
import org.springframework.transaction.PlatformTransactionManager;
import javax.sql.DataSource;
import java.util.List;
/**
* Step 08 — 연습문제 정답과 해설.
*
* 문제를 직접 풀어 본 뒤에 여세요.
*/
public class Solution {
// ======================================================================
// 정답 1 — 5.x 시그니처로 마이그레이션
// ======================================================================
/*
* 예상 컴파일 에러:
* error: Q1_LegacyWriter is not abstract and does not override abstract method
* write(Chunk<? extends Settlement>) in ItemWriter
* error: method does not override or implement a method from a supertype
*
* 4.x 의 write(List<? extends T>) 는 5.0 에서 write(Chunk<? extends T>) 로 바뀌었습니다.
* 소스 호환성이 깨지는 변경이고, 이것은 "좋은 실패" 입니다 — 컴파일러가 잡아 주니까요.
* 이 코스에서 다루는 대부분의 함정은 컴파일러가 못 잡는 것들입니다.
*
* 마이그레이션은 어댑터 한 줄이면 끝납니다. 본문(doWrite)은 손대지 않습니다.
* Chunk 로 바뀐 실질적 이득은 getSkips() / getErrors() 처럼
* "아이템 말고 다른 정보"를 함께 실을 수 있다는 점입니다. List 로는 불가능했습니다.
*/
public static class A1_MigratedWriter implements ItemWriter<Settlement> {
@Override
public void write(Chunk<? extends Settlement> chunk) {
doWrite(chunk.getItems());
}
private void doWrite(List<? extends Settlement> items) {
System.out.println("wrote " + items.size() + " items");
}
}
// ======================================================================
// 정답 2 — record 와 beanMapped()
// ======================================================================
/*
* (a)
* 예외 클래스: org.springframework.dao.InvalidDataAccessApiUsageException
* 메시지 요지: No value supplied for the SQL parameter 'orderId':
* Invalid property 'orderId' of bean class
* [com.example.batch.domain.Settlement]
*
* 이유: beanMapped() 는 BeanPropertyItemSqlParameterSourceProvider 를 쓰고,
* 이 클래스는 BeanWrapper 를 통해 "자바빈 규약"으로 값을 읽습니다.
* 자바빈 규약의 읽기 접근자는 getOrderId() 입니다.
* record 의 접근자는 orderId() — get 접두사가 없습니다.
* 그래서 BeanWrapper 는 orderId 라는 읽을 수 있는 프로퍼티가 없다고 판단합니다.
*
* (b) settlement 가 nullable 이었다면 — 이쪽이 진짜 무서운 경우입니다.
* writeCount = 70000
* Job STATUS = COMPLETED
* 데이터 = 70,000행이 들어가 있고 order_id 부터 net_amount 까지 전부 NULL
*
* 로그도 정상, 카운터도 정상, 행 수도 정상입니다.
* "정산이 안 맞는다"는 문의가 며칠 뒤에 옵니다.
*
* 실제로 재현해 보려면:
* CREATE TABLE settlement_nullable LIKE settlement;
* ALTER TABLE settlement_nullable
* MODIFY gross_amount DECIMAL(12,2) NULL,
* MODIFY net_amount DECIMAL(12,2) NULL,
* DROP INDEX uk_settlement_order;
* 그리고 이 테이블로 beanMapped() writer 를 돌려 보세요.
*
* 교훈: 결과 테이블의 컬럼은 되도록 NOT NULL 로 만드세요.
* 제약은 성능 장치이기 이전에 "조용한 실패를 시끄럽게 만드는 장치" 입니다.
* project 스펙의 settlement 가 전 컬럼 NOT NULL + uk_settlement_order 인 것은
* 우연이 아닙니다.
*
* (c) 정답 코드는 아래.
* 장황해 보이지만 record 의 컴포넌트 이름을 바꾸면 여기서 컴파일 에러가 납니다.
* beanMapped() 는 리팩터링을 따라오지 못하고, 런타임에야 알려 줍니다.
*
* 읽기 쪽도 대칭입니다: BeanPropertyRowMapper 는 record 를 못 채우고,
* DataClassRowMapper 가 생성자 기반으로 채웁니다(Step 06).
*/
public static class A2_RecordWriter {
public JdbcBatchItemWriter<Settlement> build(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()))
.assertUpdates(true)
.build();
}
}
// ======================================================================
// 정답 3 — rewriteBatchedStatements 실측
// ======================================================================
/*
* ┌───────────────┬──────────────┬─────────────┬────────────┐
* │ │ Step 실행시간 │ Com_insert │ writeCount │
* ├───────────────┼──────────────┼─────────────┼────────────┤
* │ (A) 옵션 없음 │ 1m 1s 412ms │ 70000 │ 70000 │
* │ (B) 옵션 있음 │ 7s 702ms │ 70 │ 70000 │
* └───────────────┴──────────────┴─────────────┴────────────┘
* 61.4초 -> 7.7초. 약 8배.
*
* 왜 8배인가 — 분해해 봅니다.
* 줄어든 왕복 = 70,000 - 70 = 69,930회
* 줄어든 시간 = 61.4 - 7.7 = 53.7초
* 왕복 1회당 = 53.7s / 69,930 ≈ 0.77ms
* 로컬 도커 MySQL 에 대한 TCP 왕복 + 파싱 비용으로 타당한 값입니다.
*
* 즉 사라진 것은 DB 가 한 일이 아니라 "왕복" 입니다.
* InnoDB 가 실제로 삽입한 행 수는 양쪽 다 70,000 으로 같습니다.
* 그래서 서버를 키워도, 인덱스를 지워도 8배가 안 나옵니다.
* 병목이 DB 가 아니기 때문입니다.
*
* (a) writeCount, Job STATUS, 결과 행 수가 모두 같습니까?
* 네, 완전히 같습니다. 70000 / COMPLETED / 70000.
*
* (b) 모니터링만으로 발견 가능한가?
* 아니오. 어떤 에러도 경고도 남지 않습니다.
* 그래서 "우리 배치는 원래 한 시간 걸려요" 가 굳어집니다.
*
* 대안 진단법: 서버 쪽 Com_insert 를 재세요.
* FLUSH STATUS; -> Job 실행 -> SHOW GLOBAL STATUS LIKE 'Com_insert';
* 이 값이 "처리 건수와 같으면" 배치 쓰기가 배치가 아닌 것입니다.
* 정상이라면 처리건수 / 청크크기 정도의 값이 나와야 합니다(70,000/1,000 = 70).
*
* 참고: PostgreSQL·Oracle 드라이버에는 이 옵션이 없습니다.
* 기본적으로 배치를 제대로 묶어 보냅니다. MySQL 만의 기본값 문제입니다.
*/
// ======================================================================
// 정답 4 — assertUpdates 의 두 구멍
// ======================================================================
/*
* (a) WHERE 절 없는 UPDATE
* 각 아이템마다 settlement 70,000행이 전부 갱신됩니다.
* 청크의 1,000개 아이템이 순서대로 전체를 덮으므로,
* 최종적으로 남는 값은 "마지막에 처리된 아이템의 fee_rate" 하나입니다.
* 70,000행 전체가 같은 fee_rate 를 갖게 됩니다.
*
* assertUpdates(true) 가 잡습니까? 아니오.
* 검사는 `updateCounts[i] == 0` 뿐입니다. 여기서는 각 항목이 70000 입니다.
* 0 이 아니므로 통과합니다. writeCount 도 70000 으로 정상입니다.
*
* 이것이 assertUpdates 의 본질적 한계입니다.
* "너무 적게 갱신됨"만 잡고 "너무 많이 갱신됨"은 개념적으로 잡을 수 없습니다.
*
* (b) SUCCESS_NO_INFO
* 값 = java.sql.Statement.SUCCESS_NO_INFO, 숫자로는 -2.
*
* rewriteBatchedStatements=true 일 때 UPDATE/DELETE 배치는 멀티 VALUES 로
* 합칠 수 없어서 여러 문장을 이어 보내는 형태로 재작성됩니다.
* 이때 드라이버는 개별 문장의 갱신 건수를 분리해 낼 수 없어 -2 를 채웁니다.
*
* `updateCounts[i] == 0` 검사를 통과합니까? 네, 통과합니다. -2 != 0 이니까요.
*
* 의미: 성능 옵션을 켜는 순간 UPDATE writer 의 안전장치가 조용히 꺼집니다.
* 8-3 에서 "가장 가성비 좋은 한 줄"이라고 한 옵션이,
* 8-4 의 안전장치를 무력화합니다. 두 절을 따로 읽으면 이 상호작용을 놓칩니다.
*
* (c) 대안 방어책
* (a) 의 방어책 — 코드 리뷰 규칙으로 못 박습니다.
* "ItemWriter 의 UPDATE/DELETE 문에는 반드시 유니크 키가 WHERE 에 있어야 한다."
* settlement 라면 order_id (uk_settlement_order) 입니다.
* 자동화하려면 SQL 문자열을 정적 검사하는 테스트를 붙이세요.
*
* (b) 의 방어책 — 갱신 건수를 믿지 말고 Step 뒤에서 결과를 직접 검증합니다.
* 아래 A4_VerifyTasklet 처럼 Tasklet Step 을 붙이고,
* 기대 건수와 다르면 예외를 던져 Job 을 FAILED 로 만듭니다.
* Step 10 의 흐름 제어를 쓰면 검증 실패 시 보정 Step 으로 분기시킬 수도 있습니다.
*/
public static class A4_VerifyTasklet implements Tasklet {
private final JdbcTemplate jdbcTemplate;
private final long expected;
public A4_VerifyTasklet(JdbcTemplate jdbcTemplate, long expected) {
this.jdbcTemplate = jdbcTemplate;
this.expected = expected;
}
@Override
public RepeatStatus execute(StepContribution contribution, ChunkContext ctx) {
Long actual = jdbcTemplate.queryForObject(
"SELECT COUNT(*) FROM settlement", Long.class);
if (actual == null || actual != expected) {
// 조용히 틀리느니 시끄럽게 실패하는 편이 낫습니다.
throw new IllegalStateException(
"정산 건수 불일치: expected=" + expected + ", actual=" + actual);
}
return RepeatStatus.FINISHED;
}
}
// ======================================================================
// 정답 5 — FlatFileItemWriter
// ======================================================================
/*
* wc -l 예측 = 70002 줄.
* 데이터 70,000행 + headerCallback 1줄 + footerCallback 1줄.
*
* FieldExtractor 를 람다로 쓰는 이유는 문제 2 와 완전히 같습니다.
* BeanWrapperFieldExtractor 는 이름 그대로 BeanWrapper 를 쓰므로
* record 의 orderId() 를 읽지 못합니다.
*
* 옵션 중 실무에서 실제로 사고를 내는 것 둘:
* - encoding 을 지정하지 않으면 플랫폼 기본 인코딩이 쓰입니다.
* 개발자 맥에서는 UTF-8, 운영 리눅스 컨테이너에서 POSIX 로케일이면 US-ASCII 가
* 되어 한글이 '?' 로 깨집니다. 로컬에서는 절대 재현되지 않습니다.
* - lineSeparator 를 지정하지 않으면 System.lineSeparator() 입니다.
* 윈도우에서 만든 파일에 \r\n 이 섞여 수신 시스템의 파서가 마지막 필드에
* \r 을 붙여 읽습니다. 이것도 조용한 실패입니다.
* 둘 다 "명시하면 끝나는" 문제라 명시하는 것이 정답입니다.
*/
public static class A5_FileWriter {
public FlatFileItemWriter<Settlement> build() {
DelimitedLineAggregator<Settlement> aggregator = new DelimitedLineAggregator<>();
aggregator.setDelimiter(",");
aggregator.setFieldExtractor(s -> new Object[]{
s.orderId(), s.customerId(), s.settleDate(),
s.grossAmount(), s.feeRate(), s.feeAmount(), s.netAmount()
});
return new FlatFileItemWriterBuilder<Settlement>()
.name("exerciseFileWriter")
.resource(new FileSystemResource("out/exercise-settlement.csv"))
.encoding("UTF-8")
.lineSeparator("\n")
.lineAggregator(aggregator)
.headerCallback(w -> w.write(
"order_id,customer_id,settle_date,gross_amount,"
+ "fee_rate,fee_amount,net_amount"))
.footerCallback(w -> w.write("# total lines above = 70000"))
.shouldDeleteIfExists(true)
.build();
}
}
// ======================================================================
// 정답 6 — .stream() 수동 등록
// ======================================================================
/*
* 1단계 관찰 결과
* Job STATUS = FAILED (WriterNotOpenException) 또는
* COMPLETED (파일은 0줄) — vipFileWriter 가 @Bean 이라
* afterPropertiesSet() 이 불려 빈 파일만 만들어진 경우
* DB 행 수 = 55000 (JdbcBatchItemWriter 는 ItemStream 이 아니라 멀쩡합니다)
* 파일 줄 수 = 0
* 예외 클래스 = org.springframework.batch.item.WriterNotOpenException
* "Writer must be open before it can be written to"
* vipFileWriter.written 키 = 없음
*
* 2단계 — 왜 이렇게 되는가
* StepBuilder 는 .reader()/.processor()/.writer() 로 "직접 넘긴" 객체가
* ItemStream 이면 자동으로 stream() 에 등록합니다.
* 여기서 직접 넘긴 것은 routingWriter 하나뿐이고,
* ClassifierCompositeItemWriter 는 ItemStream 을 구현하지 않습니다.
*
* public class ClassifierCompositeItemWriter<T> implements ItemWriter<T>
* ^^^^^^^^^^^^ 이것뿐
*
* 그래서 안에 숨은 vipFileWriter 의 콜백이 하나도 불리지 않습니다.
* 빠진 것은 셋이고 각각 다른 피해를 냅니다.
*
* open() 미호출 -> 파일이 열리지 않음 + headerCallback 미실행
* update() 미호출 -> ExecutionContext 에 written 위치가 안 남음
* => 재시작하면 파일을 처음부터 다시 씀.
* 리더는 중단 지점부터 이어가는데 파일은 처음부터라
* 파일에는 뒤쪽 데이터만 남습니다. DB 는 맞고 파일은 틀립니다.
* close() 미호출 -> footerCallback 미실행 + 마지막 버퍼 유실
* (transactional=true 라 출력이 버퍼에 있습니다)
*
* 진짜 원인은 "비대칭" 입니다.
* CompositeItemWriter : ItemStreamWriter 구현 -> 델리게이트에 전파함
* ClassifierCompositeItemWriter : ItemWriter 만 구현 -> 전파 안 함
* 이름이 한 단어 다를 뿐인데 정반대로 동작합니다.
* "CompositeItemWriter 로 잘 되던 코드"를 분기가 필요해져 바꾸는 순간
* 파일이 조용히 비기 시작합니다. 코드 리뷰에서는 클래스 이름만 바뀐 것으로 보입니다.
*
* 3단계 — 고치기: .stream(vipFileWriter) 한 줄.
* JdbcBatchItemWriter 는 ItemStream 이 아니므로 등록할 필요가 없습니다
* (파일 핸들 같은 상태가 없습니다).
*
* 고친 뒤 검증:
* DB 55000 + 파일 15002줄 (= 15000 + 헤더 + 푸터) = 70000
* SELECT SHORT_CONTEXT FROM BATCH_STEP_EXECUTION_CONTEXT
* ORDER BY STEP_EXECUTION_ID DESC LIMIT 1;
* -> {"vipFileWriter.written":15000, "allOrdersReader.read.count":100000}
*
* 습관으로 만들 점검 3줄:
* 1. ItemStream 구현체가 .reader()/.writer() 에 직접 들어갔는가?
* 2. 아니라면 .stream(...) 으로 손수 등록했는가?
* 3. 한 번 돌리고 컨텍스트에 <name>.written 키가 있는가?
* 3번은 10초면 됩니다. 안 하면 재시작이 필요해지는 그날까지 아무도 모릅니다.
*/
public static class A6_RoutingStep {
public Step build(JobRepository jobRepository, PlatformTransactionManager tx,
ItemReader<Order> reader,
ItemProcessor<Order, Settlement> processor,
ClassifierCompositeItemWriter<Settlement> routingWriter,
FlatFileItemWriter<Settlement> vipFileWriter) {
return new StepBuilder("q6RoutingStep", jobRepository)
.<Order, Settlement>chunk(1000, tx)
.reader(reader)
.processor(processor)
.writer(routingWriter)
.stream(vipFileWriter) // ← 숨어 있는 ItemStream 을 수동 등록
.build();
}
}
// ======================================================================
// 정답 7 — 커스텀 ItemStreamWriter 의 생명주기
// ======================================================================
/*
* (a) open() 을 비워 두면
* 재시작 시 이전에 전송한 건수를 복원하지 못해 sent 가 0 부터 시작합니다.
* 이 예제처럼 sent 를 "몇 건 보냈나" 리포트에만 쓴다면 리포트만 틀립니다.
* 그러나 sent 를 "여기부터 보내면 된다"는 위치로 쓴다면 중복 전송이 일어납니다.
* 외부 리소스를 여는 writer(파일·소켓·커넥션)라면 아예 열리지 않아
* 첫 write() 에서 즉시 터집니다.
*
* (b) update() 를 비워 두면
* ExecutionContext 에 위치가 저장되지 않습니다.
* open() 이 복원할 것이 애초에 없으므로, (a) 와 같은 결과가 되지만
* 발현 시점이 다릅니다.
*
* (c) 어느 쪽이 더 발견하기 어려운가 — update() 쪽입니다.
*
* ┌───────────┬────────────────────┬──────────────────────────┐
* │ │ 첫 실행 │ 재시작 │
* ├───────────┼────────────────────┼──────────────────────────┤
* │ open 누락 │ 즉시 실패하기도 함 │ 위치 복원 안 됨 │
* │ update 누락│ 완벽하게 성공 │ 처음부터 다시 전송(중복) │
* └───────────┴────────────────────┴──────────────────────────┘
*
* update() 누락은 첫 실행이 완벽하게 성공합니다.
* 단위 테스트도 통과하고, 스테이징에서도 통과하고, 운영 첫날도 통과합니다.
* 재시작이 필요해지는 날 — 즉 이미 장애 대응 중인 순간에 — 두 번째 장애로
* 모습을 드러냅니다. 최악의 타이밍입니다.
*
* 이 코스가 반복해 말하는 종류의 버그입니다:
* 문법 에러는 금방 고치지만, 에러 없이 조용히 틀리는 코드가 진짜 위험합니다.
*
* (d) 청크 롤백 시
* DB writer 는 함께 롤백됩니다. 외부 API 호출은 롤백되지 않습니다.
* 청크의 900번째에서 예외가 나면 DB 는 깨끗이 되돌아가지만
* API 쪽에는 이미 1,000건이 들어가 있습니다. 재시도하면 2,000건이 됩니다.
*
* 방어책은 멱등 키입니다.
* order_id 처럼 아이템마다 고유한 값을 요청에 실어 보내고,
* 수신 측이 같은 키를 두 번 받으면 무시하게 만듭니다.
* 이 코스의 settlement.order_id 에 걸린 UNIQUE 제약이 바로 그 DB 버전입니다.
* "한 주문에 정산 한 번"을 스키마가 보장하므로, 배치를 두 번 돌려도
* 정산이 두 배가 되는 대신 DuplicateKeyException 으로 시끄럽게 실패합니다.
*/
public static class A7_CustomWriter implements ItemStreamWriter<Settlement> {
private static final String CTX_KEY = "q7Writer.sent";
private long sent = 0;
@Override
public void open(ExecutionContext ctx) {
this.sent = ctx.getLong(CTX_KEY, 0L); // 재시작 복원
}
@Override
public void write(Chunk<? extends Settlement> chunk) {
// 실제 구현이라면 여기서 chunk.getItems() 를 통째로 한 번에 보냅니다.
// 아이템 단위 루프로 원격 호출하면 8-3 의 8배 문제가 그대로 재현됩니다.
sent += chunk.size();
}
@Override
public void update(ExecutionContext ctx) {
ctx.putLong(CTX_KEY, sent); // 청크 커밋마다 위치 저장
}
@Override
public void close() {
}
}
}