대규모 데이터 파이프라인을 운영하다 보면 필연적으로 원천 로그 재수집이나 과거 데이터 재적재(Backfill) 상황을 마주하게 됩니다. 이때 가장 치명적인 문제는 재실행 시 중복 적재(Duplicate Insert)나 데이터 유실이 발생하는 것입니다.
“Airflow DAG? 그냥 실패하면 웹 UI에서 Clear 눌러서 다시 돌리면 되는 거 아닌가?”
실제 운영 환경에서 단순 INSERT INTO로 작성된 배치를 재실행했다가 집계 수치가 2배로 부풀려지는 대형 장애를 겪은 후, 데이터 엔지니어링의 핵심 원칙인 멱등성(Idempotency)을 확보하기 위해 아키텍처와 쿼리 구조를 완전히 재설계했습니다.
1. 멱등성을 보장하는 2대 핵심 적재 패턴
Airflow 태스크는 동일한 실행 파라미터(예: {{ ds }})로 N번을 반복 실행해도 데이터베이스의 최종 상태가 정확히 1번 실행된 것과 완전히 동일해야 합니다.
패턴 A: Partition Delete-Write (트랜잭션 원자적 교체)
해당 일자 파티션 데이터를 단일 트랜잭션 내에서 제거 후 재적재하는 방식으로, 대용량 일 단위 파티션 적재에 가장 확실하고 안전합니다.
BEGIN;
-- 1. 대상 실행일자(ds) 파티션 데이터 원자적 삭제
DELETE FROM daily_summary WHERE stat_date = '{{ ds }}';
-- 2. 해당 일자 RAW 데이터 정제 및 재적재
INSERT INTO daily_summary (stat_date, user_id, amount)
SELECT '{{ ds }}', user_id, SUM(amount)
FROM raw_logs WHERE log_date = '{{ ds }}'
GROUP BY user_id;
COMMIT;
패턴 B: UPSERT (ON CONFLICT DO UPDATE)
고유 식별 키(Unique Key)가 존재하는 이벤트 로그 테이블의 경우 충돌 발생 시 기존 데이터를 갱신(Update)하거나 무시(Do Nothing)하도록 설계합니다.
INSERT INTO raw_event_logs (source_file, event_id, occur_time, payload)
VALUES ('log_20260827_01.csv', 'EVT-9821', '2026-08-27 07:00:00', '{"status":"OK"}')
ON CONFLICT (source_file, event_id)
DO UPDATE SET
occur_time = EXCLUDED.occur_time,
payload = EXCLUDED.payload;
2. 1,973만 건 대용량 RAW 로그 스트리밍 최적화
단순 쿼리 멱등성 외에도 메모리 오버헤드(OOM) 방지를 위해 Server-side Cursor(chunk_size=30000)와 Range Partitioning 기반 COPY 적재를 도입했습니다. 이를 통해 전체 메모리 점유율을 90% 이상 절감하고 배치 실행 시간을 4,175초에서 2,896초로 약 30% 단축시켰습니다.
flowchart TD
A["FTP 원천 RAW 로그 스캔"] --> B["Airflow Worker Task 시작"]
B --> C["Server-side Cursor 스트리밍 (Chunk 30,000)"]
C --> D["EUV / ER Dose 로그 파싱 & 정제"]
D --> E["PostgreSQL Range Partition COPY 적재"]
E --> F["DB Commit 완료 후 원천 FTP 파일 정리"]
3. 연관 아티클 및 실무 링크
- PostgreSQL MVCC와 Vacuum: 데이터 지우는데 용량이 왜 더 늘어날까?
- PostgreSQL FOR UPDATE SKIP LOCKED: DB 락 경합 없이 고성능 작업 큐 구축하기
- Redis Cache Stampede 현상 극복기
📚 공식 레퍼런스 (References)
- Apache Airflow Official Best Practices — Idempotency
- PostgreSQL Official Documentation — ON CONFLICT Clause
- PostgreSQL Table Partitioning Guide
함께 읽으면 좋은 글
- PostgreSQL MVCC와 Vacuum: 데이터 지우는데 용량이 왜 더 늘어날까?
- Redis Cache Stampede 현상 극복기: 캐시 적용하면 DB 부하 끝 아닌가?
- JPA N+1 문제와 Fetch Join 실무 고찰: lazy loading이면 다 해결되나?
공식 문서와 확인 자료
아래 1차 자료를 기준으로 명령 동작과 적용 조건을 다시 확인할 수 있습니다.