バッチ取り込みとウォーターマーク
한국어 원문으로 표시합니다.
목표
추출, 적재, 워터마크 기록, 증분 처리, 재실행 검증까지 배치 파이프라인의 한 사이클을 손으로 돌려 봅니다.
왜 중요한가
배치 파이프라인에서 사고가 나는 지점은 거의 정해져 있습니다. 경계 조건과 재실행입니다.
경계 조건은 > 와 >= 하나 차이로 매 실행마다 한 건이 중복되거나 누락됩니다. 하루에 한 번 도는 배치라면 1년에 365건이 틀어지고, 그동안 아무도 모릅니다. 그래서 추출 직후 건수를 원본과 대조하는 검증 단계를 파이프라인 안에 넣는 것이 좋습니다.
재실행은 더 중요합니다. 파이프라인은 반드시 실패하고, 실패하면 다시 돌립니다. 이때 결과가 달라지면 데이터가 오염됩니다. 그래서 적재 대상 표에 유일 제약을 걸고 충돌 시 무시하거나 갱신하도록 설계해야 합니다. 이 실습에서는 같은 증분 작업을 두 번 돌려 건수가 변하지 않는지 직접 확인합니다.
단계
작업 디렉터리는 /root/etl 입니다.
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'미만입니다.- 파일의 줄 수가 해당 기간 주문 건수에 헤더 한 줄을 더한 값과 같은지 확인합니다.
orders_archive표를 만듭니다. 컬럼은id,ordered_at,status,total_amount순서이며id에 기본 키나 유일 제약이 있어야 합니다.- 1번에서 만든 CSV 를
orders_archive에 적재합니다. etl_watermark표를 만듭니다. 컬럼은job_name,last_ordered_at이며,job_name이orders_archive인 행에 적재한 데이터의 최대ordered_at을 기록합니다.- 워터마크 이후부터
TIMESTAMPTZ '2025-03-01 00:00:00+09'미만까지의 주문을 증분 적재하고 워터마크를 갱신합니다. 적재 후orders_archive는 1월과 2월 주문을 모두 담고 있어야 합니다. - 같은 증분 작업을 한 번 더 실행합니다. 실행 전후의
orders_archive건수를/root/etl/rerun.txt에 한 줄씩 두 줄로 기록합니다. 두 값이 같아야 하고 중복 행이 없어야 합니다. 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: 증분 적재 후 워터마크 갱신을 빠뜨리면 다음 실행이 같은 구간을 또 읽습니다.
1월 주문을 CSV 로 추출하기
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' 미만입니다.
psql 의 \copy 는 클라이언트 쪽 파일로 씁니다. 헤더를 포함하는 옵션이 있습니다.
추출 건수 검증하기
파일의 줄 수가 해당 기간 주문 건수에 헤더 한 줄을 더한 값과 같은지 확인합니다.
파일 줄 수는 데이터 행 수에 헤더 한 줄을 더한 값이어야 합니다. 경계 조건을 한쪽만 포함했는지 확인하세요.
적재 대상 표 만들기
orders_archive 표를 만듭니다. 컬럼은 id, ordered_at, status, total_amount 순서이며 id 에 기본 키나 유일 제약이 있어야 합니다.
재실행 안전성은 여기서 결정됩니다. 같은 주문이 두 번 들어올 수 없게 만드는 제약이 필요합니다.
CSV 를 표에 적재하기
1번에서 만든 CSV 를 orders_archive 에 적재합니다.
\copy 는 반대 방향으로도 씁니다. 헤더 줄을 건너뛰는 옵션을 잊지 마세요.
워터마크 기록하기
etl_watermark 표를 만듭니다. 컬럼은 job_name, last_ordered_at 이며, job_name 이 orders_archive 인 행에 적재한 데이터의 최대 ordered_at 을 기록합니다.
어디까지 처리했는지를 표에 남깁니다. 값은 적재한 데이터의 마지막 시각이어야 합니다.
2월분 증분 적재하기
워터마크 이후부터 TIMESTAMPTZ '2025-03-01 00:00:00+09' 미만까지의 주문을 증분 적재하고 워터마크를 갱신합니다. 적재 후 orders_archive 는 1월과 2월 주문을 모두 담고 있어야 합니다.
워터마크 이후의 데이터만 골라 넣고, 끝나면 워터마크를 갱신합니다.
두 번 돌려도 같은지 확인하기
같은 증분 작업을 한 번 더 실행합니다. 실행 전후의 orders_archive 건수를 /root/etl/rerun.txt 에 한 줄씩 두 줄로 기록합니다. 두 값이 같아야 하고 중복 행이 없어야 합니다.
같은 증분 작업을 한 번 더 실행하고 전후 건수를 파일에 두 줄로 남깁니다. 충돌 시 무시하는 절이 필요합니다.
적재 결과 보고서 만들기
orders_archive 의 상태별 건수를 /root/etl/report.tsv 에 저장합니다. 상태와 건수를 탭으로 구분하고 헤더는 넣지 않습니다.
상태별 건수를 탭으로 구분해 저장합니다. 헤더 없이 값만 남기세요.