정상 실행만 생각하면 배치는 읽고, 변환하고, 쓰는 반복문입니다. 하지만 10만 건의 주문 중 6만 건을 처리한 뒤 프로세스가 멈추면 어디까지 성공했고, 어디서부터 다시 시작하며, 이미 쓴 데이터는 어떻게 판단해야 할까요?
이 질문에 답하려면 처리량보다 먼저 실패했을 때 어디까지 commit됐는지를 정의해야 합니다. Flowpath의 초기 Core를 설계하면서 Spring Batch를 참고한 이유도 여기에 있습니다. Spring Batch의 Java 객체 구조를 Rust로 옮기려는 것이 아니라, 오랫동안 다듬어 온 Chunk, checkpoint, transaction, restart semantics를 Rust의 async Stream과 ownership, type-state로 다시 표현해 보고 싶었습니다.
이 글에서는 왜 Source와 Cursor, Pipeline과 Runtime, Writer와 Committer를 나누게 됐는지 설명합니다.
먼저 상황을 가정해보겠습니다.
- 10만 건의 주문을 읽는다.
- 1,000건씩 변환해 저장하고 처리 위치를 기록한다.
- 6만 건까지는 저장과 위치 기록이 끝났다.
- 다음 1,000건을 처리하던 중 업무 데이터 저장과 위치 기록 사이에서 프로세스가 종료됐다.
이 상황에서는 최소한 네 가지를 구분해야 합니다.
- Source가 실제로 읽은 위치
- 다시 처리하지 않아도 되는, checkpoint에 commit된 위치
- 업무 데이터가 저장된 상태
- 처리 건수와 실행 상태가 기록된 상태
업무 데이터는 저장됐지만 위치 기록에 실패했다면 성공일까요? 다음 실행이 이전 위치에서 다시 시작하면 같은 데이터를 중복으로 쓸 수 있습니다. 반대로 위치만 먼저 저장하면 아직 쓰지 못한 데이터를 건너뛸 수 있습니다. 모든 입력이 필터링되어 쓸 데이터가 없더라도 마지막으로 처리한 Source 위치는 기록해야 합니다.
결국 배치 프레임워크의 핵심은 read-transform-write 반복문 자체가 아닙니다. 이 상태들이 어느 transaction boundary에서 함께 commit되는지를 정하는 execution contract가 더 중요합니다.
Spring Batch의 일반적인 Chunk 지향 처리는 Item을 하나씩 읽고 필요하면 가공한 뒤, 설정된 commit interval만큼 모아 ItemWriter에 전달하고 transaction을 commit합니다. 여기서 Chunk는 여러 Item을 하나의 transaction으로 처리하는 단위입니다. 1,000이라는 수치는 예시일 뿐입니다. 중요한 것은 개수가 아니라 transaction과 checkpoint가 commit되는 시점입니다. Spring Batch 공식 문서도 Chunk를 transaction 범위에서 처리되는 Item으로 정의합니다.
Spring Batch의 실행 모델에는 단순한 Reader와 Writer보다 더 많은 개념이 있습니다.
JobInstance는 Job과 식별용 JobParameters로 구분되는 논리적인 작업입니다.JobExecution은 그 작업을 실제로 실행한 한 번의 시도입니다.ExecutionContext는 실패 후 재시작에 필요한 상태를 보관합니다.JobRepository는 JobExecution, StepExecution 같은 batch metadata를 저장합니다.
특히 Step에 속한 ExecutionContext는 commit 지점마다 저장되며, ItemStream.update는 commit 전에 현재 상태를 ExecutionContext에 반영합니다. 그래서 실패한 같은 JobInstance를 다시 실행할 때 Reader가 이전 위치에서 이어갈 수 있습니다. 이 동작은 Domain Language와 ItemStream 문서에 정리되어 있습니다.
다만 여기서 Spring Batch는 업무 데이터와 metadata를 언제나 하나의 transaction으로 묶는다고 일반화하면 안 됩니다. Step의 transaction manager와 JobRepository가 사용하는 transaction manager가 다르면 처리 DB와 repository는 같은 transaction 범위가 아닙니다. Spring Batch 문서도 이 경우 처리 완료 후 repository 갱신 전에 실패하면 Step이 재실행되어 중복 처리가 발생할 수 있으며, idempotent processing이나 외부 transaction 관리가 필요하다고 설명합니다. 자세한 조건은 Configuring a Step에서 확인할 수 있습니다.
Flowpath가 가져오고 싶었던 것은 이런 execution semantics 자체였습니다. 그리고 다음과 같이 반영하였습니다.
| Spring Batch의 개념 | Flowpath의 API |
|---|---|
| ItemReader의 읽기 상태와 ExecutionContext | Source, typed Cursor, CursorCodec |
| Job/Step 구성과 실행 | immutable Pipeline, stateful Runtime |
| commit interval | 소비한 Source Item 기준 Chunk |
| ItemWriter와 transaction manager | Writer, Committer, TransactionManager |
| 업무 resource와 repository의 transaction 조건 | AtomicGuarantee, CoordinatedGuarantee |
가장 먼저 정한 것은 의존 방향이었습니다.
Core에는 Pipeline과 execution contract, adapter가 구현해야 할 capability trait만 둡니다. Runtime, 데이터베이스, 파일 시스템, observability SDK는 Core가 알지 못합니다. 테스트 지원 crate도 실제 adapter와 같은 public contract를 검증하지만 제품 Runtime의 의존성은 아닙니다.
이 분리가 중요한 이유는 Committer에 있습니다. Core가 특정 데이터베이스의 transaction 타입을 직접 알기 시작하면 무엇을 함께 commit해야 하는가라는 본래 문제가 SQL client와 lifetime 같은 구현에 묻힙니다. 반대로 저장 기술을 Core 밖에 두면 Committer는 구체 query를 노출하지 않고 guarantee만 API에 드러낼 수 있습니다.
현재 flowpath-runtime-local은 의존 방향을 고정한 골격만 있습니다. execution contract는 deterministic in-memory Runtime으로 테스트하고 있습니다.
Spring Batch의 stateful ItemReader와 ExecutionContext가 함께 다루는 문제를 Flowpath에서는 데이터를 만드는 Source와 진행 위치를 나타내는 Cursor로 나눴습니다.
pub trait CursorStream:
Stream<Item = Result<Self::SourceItem, Self::ReadError>>
{
type SourceItem;
type Cursor;
type ReadError;
fn safe_cursor(&self) -> Option<&Self::Cursor>;
}
pub trait CursorSource: Source {
type Cursor;
type CursorStream<'source>: CursorStream<
SourceItem = Self::Item,
Cursor = Self::Cursor,
ReadError = Self::ReadError,
> + 'source
where
Self: 'source;
fn stream_from(
&self,
start: SourceStart<'_, Self::Cursor>,
) -> impl Future<Output = Result<Self::CursorStream<'_>, Self::OpenError>>;
}
우선 Stream의 public Item에 Cursor를 넣지 않습니다. 사용자 transform은 Cursor나 내부 sequence가 아니라 순수한 Item만 받습니다. source progress와 업무 변환이 섞이지 않으므로 같은 transform을 cursorless Source에도 적용할 수 있습니다.
다음으로 safe_cursor()는 adapter가 page를 미리 가져온 위치가 아니라 consumer에게 마지막으로 전달한 Item의 위치를 반환합니다. 내부 buffer에 100개를 읽어 놓고 그중 30개만 전달했다면 30번째 Item의 Cursor를 반환해야 합니다. read error가 발생한 Item의 Cursor는 반영하지 않습니다.
Source마다 별도의 Cursor type을 정의합니다. 데이터베이스 keyset이라면 (created_at, id), 파일이라면 (byte_offset, line_number), partition stream이라면 (partition, offset)처럼 표현합니다. CursorCodec은 Cursor type identifier와 serialization version을 함께 저장하고, Source를 열기 전에 잘못된 타입이나 지원하지 않는 version을 거부합니다.
이 설계는 모든 Item이 filter되는 경우에도 의미가 있습니다.
읽은 Item: 1, 3, 5
filter 결과: []
업무 output: 없음
Source progress: consumer가 전달받은 마지막 Item 5의 위치를 기록
output 개수로 source progress를 계산하면 이 처리 단위의 위치를 잃습니다. 그래서 Chunk boundary는 변환 결과 수가 아니라 소비한 Source Item 수로 계산하고, output이 비어 있어도 safe Cursor를 Committer에 전달합니다.
구현에는 GAT와 trait의 impl Future 반환을 사용했습니다. Source를 빌리는 구체 Stream을 반환해 기본 경로에서 boxing과 dynamic dispatch를 강제하지 않기 위해서입니다. 그 대가로 현재 Source trait은 object-safe하지 않습니다. 나중에 서로 다른 adapter를 하나의 registry에 보관해야 할 요구가 생기면 별도의 type-erasure layer를 추가할 수 있지만, 아직 필요하지 않은 할당과 virtual call을 모든 adapter의 기본 비용으로 만들지는 않았습니다.
사용자는 Pipeline을 다음처럼 구성합니다.
let pipeline = Pipeline::from(source)
.map(validate)
.then(enrich)
.filter(accepted)
.chunks(1_000)?
.commit_with(committer);
let report = runtime.run(pipeline).await?;
Pipeline은 immutable execution plan입니다. 각 combinator는 기존 Pipeline을 소비하고 새 Pipeline을 반환합니다. Pipeline 자체에는 run()이 없으며, Runtime이 완성된 Pipeline을 넘겨받아 Source stream, transform state, output buffer, Committer와 report를 소유합니다.
Pipeline의 구성 단계는 type-state로 구분했습니다.
Pipeline
└─ chunks(non-zero) → ChunkedPipeline
└─ commit_with → ExecutablePipeline
Chunk 크기가 없거나 Committer가 연결되지 않은 계획은 Runtime에 넘길 수 없습니다. .chunks(0)은 실행 중 오류가 아니라 구성 시점의 ChunkSizeError가 됩니다. transform의 associated Output이 다음 closure의 입력과 맞지 않으면 compile time에 거부됩니다. map과 then은 실패하지 않는 transform에, try_map과 try_then은 실패할 수 있는 transform에 사용합니다. 후자는 원래 error type을 그대로 보존합니다.
여기서 .chunks(1_000)은 Source의 fetch size가 아닙니다. Source adapter는 I/O 효율을 위해 한 번에 더 많거나 적은 데이터를 읽을 수 있지만, Runtime은 소비한 Source Item 1,000개를 기준으로 commit을 시도합니다. fetch buffer와 commit 단위를 분리해야 adapter의 성능 최적화가 restart semantics를 바꾸지 않습니다.
한 처리 단위의 성공을 Writer 호출로 판정하면 업무 output, checkpoint, execution stats가 서로 다른 상태로 남을 수 있습니다. 그래서 Flowpath에서는 이 셋의 commit을 Committer라는 하나의 API로 묶었습니다.
TransactionManager.beginWriter.writeCheckpointStore.saveExecutionStore.updateTransaction.commitCommitReceipt
Committer 내부 역할도 trait별로 나눴습니다.
TransactionManager와Transaction: begin, commit, rollback lifecycleWriter: 업무 output 기록CheckpointStore: Source progress 기록ExecutionStore: committed stats 기록Committer: 위 capability를 하나의 성공 또는 실패로 조합
public API와 transaction ownership은 다음 두 trait으로 구현하였습니다.
pub trait Committer<Item, Cursor> {
type Guarantee;
type Error;
fn commit(
&mut self,
batch: CommitBatch<Item>,
context: CommitContext<Cursor>,
) -> impl Future<
Output = Result<CommitReceipt<Self::Guarantee>, Self::Error>,
>;
}
pub trait Transaction {
type Error;
fn commit(self) -> impl Future<Output = Result<(), Self::Error>>;
fn rollback(self) -> impl Future<Output = Result<(), Self::Error>>;
}
Committer는 output을 담은 CommitBatch와 operation ID, Cursor, stats를 담은 CommitContext를 한 호출에서 받습니다. Transaction::commit과 rollback은 self를 소비합니다. 한 번 commit하거나 rollback한 transaction을 다시 사용하는 경로를 ownership 수준에서 막기 위한 선택입니다.
다음은 AtomicCommitter::commit의 실제 플로우입니다.
async fn commit(
&mut self,
batch: CommitBatch<Item>,
context: CommitContext<Cursor>,
) -> Result<CommitReceipt<Self::Guarantee>, Self::Error> {
let mut transaction = self.manager.begin().await.map_err(|error| {
AtomicCommitError {
failure: AtomicCommitFailure::Begin(error),
rollback: None,
}
})?;
let (operation_id, cursor, stats) = context.into_parts();
if let Err(error) = self.destination
.write(&mut transaction, batch.into_items(), &operation_id)
.await
{
let rollback = transaction.rollback().await.err();
return Err(AtomicCommitError {
failure: AtomicCommitFailure::Writer(error),
rollback,
});
}
if let Err(error) = self.checkpoints
.save(&mut transaction, cursor, &operation_id)
.await
{
let rollback = transaction.rollback().await.err();
return Err(AtomicCommitError {
failure: AtomicCommitFailure::Checkpoint(error),
rollback,
});
}
if let Err(error) = self.executions
.update(&mut transaction, stats, &operation_id)
.await
{
let rollback = transaction.rollback().await.err();
return Err(AtomicCommitError {
failure: AtomicCommitFailure::Execution(error),
rollback,
});
}
transaction.commit().await.map_err(|error| AtomicCommitError {
failure: AtomicCommitFailure::Commit(error),
rollback: None,
})?;
Ok(CommitReceipt::committed(operation_id))
}
begin에 실패하면 rollback할 transaction이 없습니다. 반면 Writer, CheckpointStore, ExecutionStore에서 실패하면 비동기 rollback을 명시적으로 기다립니다. rollback까지 실패하더라도 원래 stage 오류를 덮어쓰지 않고 두 오류를 함께 보존합니다. 마지막 transaction commit이 성공한 뒤에만 CommitReceipt를 반환합니다.
실패 지점마다 결과는 다음과 같습니다.
| 실패 지점 | 결과 |
|---|---|
| begin | transaction이 없으므로 변경과 rollback 없음 |
| writer | 명시적 rollback, Cursor 불변 |
| checkpoint | 앞서 수행한 업무 write rollback |
| execution update | write와 checkpoint rollback |
| commit | 성공 receipt를 반환하지 않음 |
| rollback | 원래 stage 오류와 rollback 오류를 모두 보존 |
이 동작은 contract test로 확인합니다. 다음 테스트는 checkpoint와 rollback 오류를 함께 주입했을 때 두 오류가 모두 남는지 검증하는 실제 테스트입니다.
#[test]
fn rollback_failure_preserves_both_errors() {
block_on(async {
let store = InMemoryCommitStore::new();
let mut committer = in_memory_atomic_committer(
store,
[CommitFault::Checkpoint, CommitFault::Rollback],
);
let error = committer
.commit(CommitBatch::new(vec![1]), context("run-1/chunk-0"))
.await
.expect_err("checkpoint and rollback fail");
assert!(matches!(
error.failure,
AtomicCommitFailure::Checkpoint(_)
));
assert_eq!(
error.rollback.expect("rollback error retained").fault,
CommitFault::Rollback
);
});
}
AtomicGuarantee가 의미하는 범위는 같은 adapter transaction에 참여한 write, checkpoint, stats가 함께 보이거나 모두 보이지 않는 것입니다. 이것만으로 여러 Runtime 인스턴스가 같은 작업을 동시에 실행할 때의 중복을 막을 수는 없습니다.
예를 들어 두 worker가 checkpoint 10을 동시에 읽으면 둘 다 다음 처리 단위를 계산할 수 있습니다.
Worker A: checkpoint 10 읽기 → 다음 단위 계산 → begin
Worker B: checkpoint 10 읽기 → 다음 단위 계산 → begin
Worker A: output + checkpoint 11 commit
Worker B: output + checkpoint 11 commit
각 transaction 내부는 원자적이어도 업무 output은 두 번 반영될 수 있습니다. in-memory store의 Mutex는 분산 lock이 아니며, 이 경쟁을 검증하는 PostgreSQL integration test도 현재 범위에는 없습니다.
멀티 인스턴스 정합성을 다루려면 같은 논리 작업과 처리 단위를 식별하는 stable identity, operation ID unique constraint, checkpoint compare-and-set, lease, heartbeat, fencing token이 추가로 필요합니다. 외부 API나 다른 데이터베이스처럼 하나의 local transaction에 넣을 수 없는 side effect에는 idempotency key와 reconciliation도 필요합니다.
그래서 guarantee를 Atomic과 Coordinated로 구분했습니다.
AtomicGuarantee
→ 같은 adapter transaction에 참여한 변경을 함께 commit
CoordinatedGuarantee
→ partial state를 드러내고 operation ID를 기반으로 reconciliation 요구
CoordinatedGuarantee는 atomicity의 다른 이름이 아닙니다. partial success를 숨기지 않고, 어디까지 진행됐으며 reconciliation에 어떤 정보가 필요한지를 error에 남깁니다.
두 설계는 같은 실패 문제를 다루지만 범위와 표현 방식이 다릅니다.
| 관점 | Spring Batch | 현재 Flowpath Core |
|---|---|---|
| 성숙도 | production framework와 풍부한 생태계 | 핵심 execution contract를 검증하는 초기 구현 |
| 구성 API | Job과 Step 구성 | immutable functional Pipeline |
| 읽기 | ItemReader, ItemStream, ExecutionContext | async Source, typed Cursor, CursorCodec |
| 실행 주체 | framework가 Job/Step lifecycle 실행 | Runtime이 Pipeline plan을 소비 |
| Chunk | commit interval과 transaction boundary | 소비한 Source Item 기준 commit boundary |
| 저장 | ItemWriter와 transaction manager | Writer와 TransactionManager capability 분리 |
| source progress | JobRepository와 ExecutionContext | CheckpointStore contract와 in-memory test double |
| consistency guarantee | 업무 resource와 repository의 transaction 구성에 따라 달라짐 | Atomic/Coordinated guarantee를 타입으로 구분 |
| 재시작과 오류 정책 | restart, retry, skip 제공 | persistent restart와 retry는 아직 미구현 |
| adapter | DB, file, messaging 통합 | production adapter는 아직 없음 |
Spring Batch가 완성된 실행 lifecycle과 adapter 생태계를 제공한다면, 현재 Flowpath는 같은 실패 문제를 Rust에서 어떤 Core contract로 표현할지 먼저 검증하고 있습니다.
이제 같은 실패를 새 API 기준으로 다시 따라가 보겠습니다.
Source가 아직 commit되지 않은 현재 처리 단위의 Item을 전달한다.
→ safe Cursor는 메모리에만 존재한다.
→ Pipeline이 Item을 변환하고 output을 모은다.
→ Committer가 transaction을 시작한다.
→ Writer가 업무 output을 기록한다.
→ CheckpointStore 저장에 실패한다.
→ Committer가 명시적으로 rollback한다.
→ committed output과 checkpoint는 이전 처리 단위 상태를 유지한다.
여기까지가 현재 구현이 보장하는 범위입니다. Core에는 async Source와 typed Cursor, functional Pipeline, Committer contract가 있고, in-memory Runtime과 failure matrix test로 주요 규칙을 검증했습니다. 하지만 persistent metadata에서 checkpoint를 불러와 새 실행을 만드는 restart command는 아직 없습니다. RDB/file adapter, retry, skip, replay, parallel/distributed Runtime, lease와 fencing도 구현 범위 밖입니다.
따라서 지금의 구현은 재시작을 지원한다거나 멀티 인스턴스 exactly-once를 보장한다기 보다는 아직 commit되지 않은 처리 단위가 persistent state를 변경하지 않도록 Core contract를 만들었다는 것입니다.
Spring Batch에서 가장 인상적이었던 것은 Job, Step, ItemReader라는 이름보다 실패와 재시작을 일관되게 설명하는 실행 모델이었습니다. Flowpath는 그 객체 구조를 그대로 번역하는 대신 읽은 위치와 commit된 위치, Pipeline과 Runtime, write와 commit을 Rust의 타입으로 구분했습니다.
그 결과 정상 경로보다 먼저 실패 경로를 설명할 수 있게 됐습니다. source progress는 output 수와 별개로 계산하고, 실행 가능한 Pipeline은 필요한 설정을 타입으로 갖추며, Committer는 업무 데이터와 checkpoint, stats를 하나의 commit 결과로 다룹니다.
이제 다음 단계는 이 contract 위에 persistent execution identity와 실제 restart를 올리고, RDB의 unique constraint와 checkpoint CAS, lease와 fencing으로 멀티 인스턴스 경쟁을 검증하는 일입니다. 그 단계에서도 기준은 같습니다. 기능 목록보다 먼저, 실패했을 때 어떤 상태가 남는지를 설명할 수 있어야 합니다.