데이터 파이프라인 · 배치와 스트리밍 · 실습
배치 적재와 워터마크
목표
추출, 적재, 워터마크 기록, 증분 처리, 재실행 검증까지 배치 파이프라인의 한 사이클을 손으로 돌려 봅니다.
왜 중요한가
배치 파이프라인에서 사고가 나는 지점은 거의 정해져 있습니다. 경계 조건과 재실행입니다.
경계 조건은 > 와 >= 하나 차이로 매 실행마다 한 건이 중복되거나 누락됩니다. 하루에 한 번 도는 배치라면 1년에 365건이 틀어지고, 그동안 아무도 모릅니다. 그래서 추출 직후 건수를 원본과 대조하는 검증 단계를 파이프라인 안에 넣는 것이 좋습니다.
재실행은 더 중요합니다. 파이프라인은 반드시 실패하고, 실패하면 다시 돌립니다. 이때 결과가 달라지면 데이터가 오염됩니다. 그래서 적재 대상 표에 유일 제약을 걸고 충돌 시 무시하거나 갱신하도록 설계해야 합니다. 이 실습에서는 같은 증분 작업을 두 번 돌려 건수가 변하지 않는지 직접 확인합니다.
단계
작업 디렉터리는 /root/etl 입니다.
1. ordered_at 이 2025년 1월인 주문의 id, ordered_at, status, total_amount 를 /root/etl/orders_2025_01.csv 로 추출합니다. 헤더를 포함합니다. 기준 시각은 TIMESTAMPTZ '2025-01-01 00:00:00+09' 이상 TIMESTAMPTZ '2025-02-01 00:00:00+09' 미만입니다.
2. 파일의 줄 수가 해당 기간 주문 건수에 헤더 한 줄을 더한 값과 같은지 확인합니다.
3. orders_archive 표를 만듭니다. 컬럼은 id, ordered_at, status, total_amount 순서이며 id 에 기본 키나 유일 제약이 있어야 합니다.
4. 1번에서 만든 CSV 를 orders_archive 에 적재합니다.
5. etl_watermark 표를 만듭니다. 컬럼은 job_name, last_ordered_at 이며, job_name 이 orders_archive 인 행에 적재한 데이터의 최대 ordered_at 을 기록합니다.
6. 워터마크 이후부터 TIMESTAMPTZ '2025-03-01 00:00:00+09' 미만까지의 주문을 증분 적재하고 워터마크를 갱신합니다. 적재 후 orders_archive 는 1월과 2월 주문을 모두 담고 있어야 합니다.
7. 같은 증분 작업을 한 번 더 실행합니다. 실행 전후의 orders_archive 건수를 /root/etl/rerun.txt 에 한 줄씩 두 줄로 기록합니다. 두 값이 같아야 하고 중복 행이 없어야 합니다.
8. orders_archive 의 상태별 건수를 /root/etl/report.tsv 에 저장합니다. 상태와 건수를 탭으로 구분하고 헤더는 넣지 않습니다.
참고
\copy (SELECT ...) TO '/root/etl/x.csv' WITH (FORMAT csv, HEADER true)\copy 표이름 (컬럼목록) FROM '/root/etl/x.csv' WITH (FORMAT csv, HEADER true)- 충돌 무시:
INSERT ... ON CONFLICT (id) DO NOTHING - 흔한 실수 1: 경계값을 양쪽 다 포함하면 월말 주문이 두 배로 들어갑니다.
- 흔한 실수 2: 증분 적재 후 워터마크 갱신을 빠뜨리면 다음 실행이 같은 구간을 또 읽습니다.
단계 8개
- 1월 주문을 CSV 로 추출하기
- 추출 건수 검증하기
- 적재 대상 표 만들기
- CSV 를 표에 적재하기
- 워터마크 기록하기
- 2월분 증분 적재하기
- 두 번 돌려도 같은지 확인하기
- 적재 결과 보고서 만들기