Data Ingestion & Transformation Engineering
원천 시스템에서 데이터를 수집하는 수송 과정과 분석 목적에 맞게 재가공하는 ETL/ELT 파이프라인의 물리적 구축을 다루는 학습 노드입니다.
목차 보기22
1. Overview
데이터 수집 및 변환 엔지니어링(Data Ingestion & Transformation Engineering, ITE)은 파편화된 수많은 원천 시스템(DB, API, Log)에서 발생하는 거대한 데이터의 격류를 안전하게 빨아들여, 분석가가 씹어먹기 좋은 형태의 정제된 테이블로 가공해 내는 '데이터의 혈맥과 제련소'를 다룹니다.
데이터는 생성되는 순간과 분석되는 순간의 물리적 형태가 완전히 다릅니다. 학습자는 백엔드 서버에 부하를 주지 않고 변경분만 조용히 빼내는 CDC(Change Data Capture) 기술과, 데이터 폭주로 인한 파이프라인 붕괴를 막는 Kafka 같은 분산 메시지 버퍼링 역학을 배웁니다. 나아가 수백 개의 데이터 변환 작업이 얽힌 복잡한 의존성 그래프(DAG)를 Airflow 같은 오케스트레이터로 제어하며, 중간에 서버가 터져도 처음부터 다시 돌리면 데이터가 이중으로 쌓이지 않는 '멱등성(Idempotency)'의 물리 법칙을 훈련합니다.
2. Scope & Boundaries
In-Scope
- 수집 아키텍처 (Ingestion Patterns): Push vs Pull 모델, 배치(Batch) vs 스트리밍(Streaming) 윈도우링(Windowing), CDC(Change Data Capture), 로그 테일링(Log Tailing).
- 파이프라인 구조 (Pipeline Models): ETL(Extract-Transform-Load) vs ELT(Extract-Load-Transform), 스키마-온-라이트(Schema-on-write) vs 스키마-온-리드(Schema-on-read).
- 데이터 가공 로직 (Transformation Logic): 결측치(Null) 처리, 역정규화(Denormalization), 데이터 평탄화(Flattening), 지연 도착 데이터(Late Data)의 워터마크(Watermark) 보정.
- 오케스트레이션 및 신뢰성 (Orchestration): DAG(Directed Acyclic Graph) 기반 스케줄링, 멱등성(Idempotency), 자가 치유(Self-healing) 재시도 로직, 데드 레터 큐(Dead Letter Queue).
Out-of-Scope
- 데이터 레이크하우스(Lakehouse) 자체의 아키텍처 구조: 데이터가 최종적으로 쌓여 저장되는 Snowflake나 Iceberg의 스토리지 포맷 설계 → 06-07. Lakehouse Arch 영역으로 위임.
- 스트리밍 클러스터의 내부 인프라 구성: Kafka 브로커의 컨트롤러 리더 선출이나 주키퍼(Zookeeper) 노드 관리 등 순수 인프라/옵스 → 07-06. Message Queues 영역으로 위임.
Boundaries
- ITE vs. Lakehouse Arch (06-07): Lakehouse(06-07)가 '적재된 데이터가 어떻게 효율적으로 저장되고 쿼리되는가'라는 창고(Warehouse)의 구조에 집중한다면, ITE는 **'그 창고까지 어떻게 물건을 파손 없이, 제시간에, 예쁘게 포장해서 배달할 것인가'**라는 물류 파이프라인(Logistics)에 초점을 맞춥니다.
3. Counterexample
- API 폴링에 의존하는 원천지 폭파 (Ingestion Fallacy): 운영 중인 메인 DB에서 최신 주문 데이터를 가져오기 위해, 데이터 엔지니어가 1분마다
SELECT * FROM orders WHERE updated_at > ?쿼리를 날리는 배치 스크립트를 짜는 행위. 인덱스가 제대로 안 걸려있다면 운영 DB에 풀 스캔 부하를 걸어 서비스 전체를 다운시키는 테러가 됩니다. 원천 DB의 트랜잭션 로그(WAL/Binlog)를 백그라운드에서 읽어내는 **CDC(Change Data Capture)**의 물리적 오버헤드 회피 개념을 모르는 안티패턴입니다. - 멱등성(Idempotency)이 없는 파이프라인 재시도 (Reliability Fallacy): 어제 날짜의 데이터를 집계하는
INSERT INTO daily_stats ...파이프라인이 중간에 네트워크 오류로 실패했을 때, 엔지니어가 그냥 재실행(Retry) 버튼을 눌렀다가 매출액이 2배로 뻥튀기되는 대참사. "몇 번을 재실행해도 결과는 1번 실행한 것과 동일해야 한다"는 멱등성 원칙을 무시하고,UPSERT나 파티션 덮어쓰기(Overwrite) 물리를 파이프라인 코드에 설계하지 않은 것은 재앙입니다.
4. Prerequisites
- 관계형 시스템 (Basic): 소스 DB의 스키마와 데이터 타입을 분석 DW로 변환(Casting)하려면 RDBMS의 데이터 모델에 대한 지식이 필요합니다. (06-01. RS)
- 분산 로직 및 저장 물리 (Recommended): 스트리밍 시스템에서 "정확히 한 번(Exactly-once)" 처리가 얼마나 물리적으로 달성하기 어려운 문제인지 이해하려면 분산 합의 지식이 권장됩니다. (06-03. DLP)
5. Learning Map
6. Learning Topics
Basic
Core Topic 01: 데이터 수집 패턴과 변경 데이터 캡처 (Ingestion & CDC)
- Why to Learn: 살아 움직이는 라이브 서비스의 성능에 1%의 타격도 주지 않고, 발생한 모든 데이터 변경(Insert, Update, Delete) 이력을 훔쳐(?) 오기 위함입니다.
- What to Learn:
- Concepts: Full Load(전체 적재) vs Incremental Load(증분 적재), 스냅샷(Snapshot), CDC(Change Data Capture).
- Skills: DB 트랜잭션 로그(MySQL Binlog, PostgreSQL WAL) 읽기 원리 파악,
updated_at워터마크 기반의 논리적 캡처 vs 로그 기반의 물리적 캡처 비교. - Tools: Debezium, AWS DMS, Kafka Connect.
- Trade-offs: 배치 쿼리로 당겨오는 방식의 쉬운 구현성 vs 소스 DB 부하 및 "A → B → C"로 변한 데이터 중 중간 B가 누락되는 로직 캡처의 한계. 반면 로그 기반 CDC는 모든 이력을 완벽히 가져오지만 DBA의 권한 협조와 인프라 세팅 비용이 큽니다.
- How to Learn:
- 1단계: RDBMS의
users테이블에 쿼리 수집으로 매일 증분 데이터를 가져올 때, 회원이 회원탈퇴(Hard Delete)를 해버리면 타겟 DW에는 여전히 그 회원이 살아있는 끔찍한(데이터 불일치) 시나리오를 증명합니다. - 2단계: 이를 해결하기 위해 Debezium(CDC)을 붙여 트랜잭션 로그에서
{"op": "d", "before": {"id": 123}}라는 순수 삭제 이벤트를 스트리밍 받아 타겟 DW에 완벽히 동기화하는 파이프라인을 스케치합니다.
- 1단계: RDBMS의
- Implement: 특정 로컬 DB의 테이블이 업데이트될 때마다, 트리거(Trigger)나 임시 스크립트를 통해 변경된 Before/After JSON 페이로드를 로그 파일로 남기는 미니 CDC 모형.
Recommended
Core Topic 02: 배치와 스트림 처리 엔진의 역학 (Batch & Stream Processing)
- Why to Learn: "어제 매출액"을 계산하는 덤프트럭(배치)과 "현재 접속자 수"를 계산하는 F1 머신(스트리밍)의 물리적 엔진 구조가 완전히 다름을 이해하고 통제하기 위해서입니다.
- What to Learn:
- Concepts: 배치 처리(Bounded Data), 스트림 처리(Unbounded Data), 마이크로 배치(Micro-batch).
- Skills: 이벤트 시간(Event Time: 데이터 생성 시각) vs 처리 시간(Processing Time: 서버 도착 시각) 불일치 해결, 텀블링 윈도우(Tumbling), 호핑 윈도우(Hopping), 워터마크(Watermark)를 통한 지연 데이터(Late Data) 폐기/수용 로직.
- Tools: Apache Spark(배치/마이크로 배치), Apache Flink(순수 스트리밍).
- Trade-offs: 시스템 복잡도는 높으나 수 밀리초의 지연(Latency)을 보장하는 스트리밍 아키텍처 vs 구현은 직관적이나 지표 확인까지 하루(24시간)를 기다려야 하는 배치 아키텍처.
- How to Learn:
- 1단계: 유저가 12<05에>05에> 결제했지만 오프라인 상태라 데이터가 12<15에>15에> 서버로 들어왔을 때, 이를 처리 시간 기준으로 집계하면 12<15분대>15분대> 매출이 폭등하는 통계 왜곡을 인지합니다.
- 2단계: 이벤트 시간(Event Time)을 기준으로 윈도우(12<00>00>~12<10>10>)를 열어두고, 12<15에>15에> 도착한 데이터를 워터마크(허용 지연 시간) 규칙에 따라 뒤늦게 12<05분>05분> 윈도우로 꽂아 넣어 통계를 재계산하는 스트리밍 엔진의 상태 보정(State Correction) 물리를 추적합니다.
- Implement: 끝없이 들어오는 실시간 JSON 로그 스트림(예: 초당 10건)을 읽어들여, 10초 단위의 텀블링 윈도우(Tumbling Window)로 묶어 에러 발생 횟수를 집계해 터미널에 뿌려주는 로컬 스크립트.
Practical
Core Topic 03: ELT 전환과 데이터 연금술 (Modern ELT & Transformation)
- Why to Learn: 과거 비싼 DB에 넣기 전 별도 서버에서 데이터를 깎던 시대(ETL)에서 벗어나, 데이터 웨어하우스 자체의 압도적 분산 컴퓨팅 파워를 활용해 SQL만으로 수십억 건의 데이터를 가공하는(ELT) 현대적 연금술을 익히기 위함입니다.
- What to Learn:
- Concepts: ETL vs ELT, 데이터 웨어하우스 컴퓨팅, 뷰(View)와 구체화된 뷰(Materialized View).
- Skills: JSON/Array 등 반정형(Semi-structured) 컬럼 평탄화(Flattening), 결측치/이상치 대체 함수(Coalesce 등) 활용, 비즈니스 도메인에 맞춘 차원(Dimension) 및 팩트(Fact) 테이블 생성 SQL.
- Tools: dbt (data build tool), Snowflake / BigQuery SQL Dialect.
- Trade-offs: 파이썬/스칼라 코드로 복잡한 객체지향적 변환 로직을 짤 수 있는 ETL의 유연성 vs SQL 하나만 알면 모든 데이터 분석가가 직접 파이프라인에 기여할 수 있는 압도적 생산성(ELT/dbt).
- How to Learn:
- 1단계: 원천 DB의
{"user": "A", "items": ["apple", "banana"]}JSON 도큐먼트를 타겟 DW의 Raw 스키마에 통째로 적재(Load)합니다. - 2단계: 타겟 DW 안에서 배열을 행(Row)으로 풀어내는(Flatten/Unnest) SQL 쿼리를 작성하여,
User-Item쌍의 정규화된 마트(Mart) 테이블을 만들어내는 ELT 물리적 흐름을 체험합니다.
- 1단계: 원천 DB의
- Implement: 원본 CSV 데이터를 임시 테이블(Staging)에 로드한 후, 순수 SQL 쿼리 두세 번의 연속 실행만으로 중복 제거, 결측치 치환, 최종 집계 테이블(Target) 삽입을 연쇄 수행하는 미니 ELT 파이프라인.
Advanced
Core Topic 04: DAG 오케스트레이션과 멱등성 보장 (Orchestration & Idempotency)
- Why to Learn: 데이터가 A → B → C 순서로 가공되어야 할 때, B가 터지면 C는 대기하고 B부터 다시 시작하더라도 A의 데이터가 두 번 중복 적재되지 않는 '강철 같은 신뢰성'의 자동화 공장을 짓기 위해서입니다.
- What to Learn:
- Concepts: DAG(Directed Acyclic Graph), 오케스트레이션(Orchestration), 멱등성(Idempotency), 백필(Backfill).
- Skills: 센서(Sensor)를 이용한 선행 데이터 도착 대기 물리, Airflow
execution_date논리 타임스탬프와 템플릿 변수를 이용한 날짜 파티션 스와핑(Partition Swapping) 전략. - Tools: Apache Airflow, Prefect, Dagster.
- Trade-offs: 단순히
cron스케줄러로 매일 밤 12시에 쉘 스크립트를 돌리는 극도의 단순함 vs 노드 10개가 얽힌 DAG를 짜서 시각적 모니터링과 부분 재시도(Retry) 기능을 확보하는 프레임워크 학습/유지보수 비용.
- How to Learn:
- 1단계: 파이프라인 로직 내부에
DELETE FROM target WHERE date = '오늘'; INSERT INTO target ...(또는MERGE) 구문을 박아넣어, 이 코드가 100번 재실행되어도 타겟 테이블에는 '오늘' 데이터가 딱 한 벌만 존재하게 하는 멱등성 물리 법칙을 증명합니다. - 2단계: 과거 1년 치 로직이 통째로 바뀌었을 때,
Backfill명령 한 방에 365개의 일일 파이프라인 태스크가 과거 날짜 변수({{ ds }})를 머금고 동시에 병렬로 돌며 DW를 과거부터 다시 채워 올리는 장관을 시뮬레이션합니다.
- 1단계: 파이프라인 로직 내부에
- Implement: '추출' → '가공' → '적재' 3개의 함수를 선언하고, 이를 순차적으로 실행하되 중간에 일부러 예외(Exception)를 발생시킨 뒤, 재시작 시 이미 완료된 '추출'은 건너뛰고 '가공'부터 안전하게 실행되는 커스텀 워크플로우 제어 뼈대.
7. Terminology
8. References
Primary References
- [P4] DS-BoK - Data Ingestion & Transformation — Pipeline principles.
- [P2] SWEBOK - Software Construction — Data processing and transformation.
Secondary References
- [Fundamentals of Data Engineering] Joe Reis & Matt Housley — Modern holistic view.
- [Data Pipelines Pocket Reference] James Densmore — Practical design patterns.
Industry References
- [Confluent Documentation - CDC & Connector Guides] — Messaging ingestion standard.
- [dbt (data build tool) Best Practices] — Modern ELT transformation industry standard.
9. Final Checklist
Primary Checklist
- 데이터 파이프라인 설계 시 하류(Downstream) 시스템의 타입 제약 조건을 고려하여 소스 데이터를 물리적으로 캐스팅 가능한가? (P4)
- 대용량 배치 작업 중 특정 태스크가 실패했을 때, 전체를 다시 시작하지 않고 실패 지점부터 재개하도록 DAG를 설계할 수 있는가? (P4)
Secondary Checklist
- 데이터 소스가 텍스트 파일(CSV/JSON)인 경우와 DB 로그(CDC)인 경우의 네트워크 대역폭 부하를 정량적으로 추정할 수 있는가?
- 변환 로직 중 'Join' 연산이 메모리에서 일어날 때와 디스크 임시 공간을 사용하며 일어날 때의 속도 차이를 인지하는가?
Industry Checklist
- 실무 환경에서 Kafka와 같은 중간 버퍼링 시스템이 생산자(Producer)와 소비자(Consumer) 사이의 물리적 속도 차이를 어떻게 완충하는지 설명 가능한가? (SFIA)
- 규제 준수를 위해 데이터 유입 단계에서 개인정보 항목을 자동 식별하여 마스킹하는 전처리 로직을 파이프라인에 통합 가능한가?
태그
etl-pipelinesdata-ingestiondata-transformationbatch-streamingtransformation-engineeringdata-dbdatainformation-managementingestiontransformationengineeringdbdatabasesdata-engineeringstreaming