ETL 파이프라인: 웹·API 데이터를 정규화해 데이터 웨어하우스에 적재하는 방법

웹 스크래핑과 API 수집 데이터를 정규화하고 증분 적재하는 ETL 파이프라인의 설계, 운영, PostgreSQL 구현 예시

2026-08-14 · 최초 발행 2024-04-29

서로 다른 원천 데이터를 분석 가능한 형태로 바꾸는 흐름

ETL은 여러 원천에서 데이터를 꺼내고(Extract), 분석과 거버넌스 기준에 맞춰 바꾼 뒤(Transform), 데이터 웨어하우스(DWH)에 넣는(Load) 데이터 파이프라인이다.

여기서는 웹 스크래핑과 API를 통한 추출, 관계형 정규화와 표준화, 증분·트랜잭션 적재를 중심으로 다룬다. 스케줄링, 관측성, 데이터 품질 검증도 파이프라인의 일부로 본다.

웹 스크래핑은 HTML DOM을 파싱해 규칙으로 데이터를 뽑아낸다. 마크업은 자주 바뀔 수 있으므로 셀렉터를 추상화하고 회귀 테스트를 마련해야 한다. API 수집은 인증, 페이징, 스루풋과 쿼터 제약을 관리하는 일이 선행된다. 웹훅이나 증분 파라미터를 사용하면 네트워크 비용을 줄일 수 있다.

변환 단계에서는 스키마 검증, 형변환, 중복 제거, 키 표준화를 수행한다. 유일성과 참조 무결성 같은 도메인 규칙으로 품질을 검증하며, 1NF/2NF/3NF 정규화로 중복을 줄인 후 분석 목적에 맞춰 스타 스키마로 다시 구성한다.

적재는 스테이징 테이블을 거쳐 MERGE/UPSERT 방식으로 변경분을 반영한다. 트랜잭션, 락, 격리 수준을 다뤄 일관성을 보장하고, 배치·마이크로배치·스트리밍 가운데 서비스 SLO와 비용에 맞는 방식을 선택한다.

DAG 스케줄링에서는 의존성, 재시도, 백오프, 서킷브레이커를 관리한다. SLA 모니터링과 알림·치유 자동화, 데이터 카탈로그와 메타데이터, 데이터 라인리지와 품질 지표는 감사 가능성을 뒷받침한다.

추출부터 격리 처리까지의 데이터 경로

TransformExtractMERGE/UPSERTValidation ErrorBEGIN→COPY→MERGE→COMMITRollback on FailureWeb ScrapingStaging RawAPI PullSchema ValidationNormalization1NF/2NF/3NFDedup & Type CastingTransactional Load(Data Warehouse)DLQ/Quarantine

입력은 웹과 API 원천 데이터이며, 스키마 검증·정규화·증분 키 생성·트랜잭션 적재를 거쳐 DWH의 디멘전·팩트 테이블로 전달된다. 검증에 실패한 데이터는 격리 테이블(DLQ)로 보내 재처리 파이프라인에서 다룬다.

웹 스크래핑과 API 수집의 운영상 차이

항목 웹 스크래핑 API
성능 캐시·병렬 크롤링 시 중간, 구조 변경 시 재시도 비용 증가 서버 제한 내 고성능, 페이징·배치로 최적화 용이
확장성 셀렉터 관리 오버헤드 큼 수평 확장 용이, 토큰 분산·샤딩 가능
일관성 마크업 변동에 취약 스키마 안정적, 버전 관리 제공
안정성 차단/봇 디텍션 리스크 쿼터·429 처리 시 안정적
운영 편의 robots.txt·크롤링 에티켓 필요 인증·회전 토큰 관리 필요

데이터 원천이 섞인 파이프라인의 구성

리테일 가격 수집에서는 주요 이커머스 상품 페이지를 스크래핑하고, 변경을 감지하면 파서를 자동 재학습한다. 변환 단계에서는 통화와 부가세를 표준화하고 브랜드·카테고리를 매핑하며 중복 SKU를 제거한다. 적재 단계는 시간 단위 스냅샷 팩트를 넣고 가격 변동 알림을 트리거한다.

SaaS CRM API 적재는 OAuth2 인증과 updated_since 증분 호출을 사용하고, 429 응답에는 지수 백오프를 적용한다. 연락처·계정 디멘전을 정규화하고 이메일·도메인을 표준화한 뒤, 스테이징에서 MERGE로 변경분을 반영하며 SCD2로 이력을 관리한다.

뉴스·리뷰 스크래핑과 판매·재고 API를 결합하는 경우도 있다. 텍스트를 정규화·토큰화한 다음 메타데이터를 조인하고, 주제별 성능 대시보드를 구성한다.

웹 페이지에서 상품 데이터를 추출하는 코드

전제조건: Python 3.10+, requests 2.x, beautifulsoup4 4.x, PostgreSQL 14+, psycopg2 2.x

# pip install requests beautifulsoup4 tenacity psycopg2-binary
import time, os
import requests
from bs4 import BeautifulSoup
from tenacity import retry, wait_exponential, stop_after_attempt

HEADERS = {"User-Agent": "etl-bot/1.0 (+contact@example.com)"}

@retry(wait=wait_exponential(multiplier=1, min=1, max=30), stop=stop_after_attempt(5))
def fetch(url):
    r = requests.get(url, headers=HEADERS, timeout=15)
    r.raise_for_status()
    return r.text

def parse_products(html):
    soup = BeautifulSoup(html, "html.parser")
    items = []
    for card in soup.select(".product"):
        sku = card.get("data-sku")
        name = card.select_one(".name").get_text(strip=True)
        price = card.select_one(".price").get_text(strip=True).replace(",", "")
        items.append({"sku": sku, "name": name, "price": float(price)})
    return items

if __name__ == "__main__":
    html = fetch("https://example.com/products")
    rows = parse_products(html)
    print(f"extracted={len(rows)}")

robots.txt를 준수하고 요청 간 지연을 두며, 셀렉터 회귀 테스트를 구축한다.

페이징 API의 쿼터 응답 처리

# pip install requests
import requests, time

API = "https://api.example.com/v1/items"
TOKEN = "YOUR_TOKEN"

def pull_since(updated_since):
    url = API
    params = {"updated_since": updated_since, "limit": 200}
    headers = {"Authorization": f"Bearer {TOKEN}"}
    while url:
        r = requests.get(url, params=params, headers=headers, timeout=20)
        if r.status_code == 429:
            time.sleep(int(r.headers.get("Retry-After", "5")))
            continue
        r.raise_for_status()
        data = r.json()
        for item in data["results"]:
            yield item
        url = data.get("next")  # next page URL
        params = None  # next에 파라미터 포함되는 패턴 가정

for rec in pull_since("2025-01-01T00:00:00Z"):
    pass

PostgreSQL 스테이징에서 정규화와 UPSERT 수행

-- 스테이징 테이블
CREATE TABLE IF NOT EXISTS stg_orders(
  order_id TEXT, order_ts TIMESTAMPTZ, customer_email TEXT,
  product_sku TEXT, qty INT, price NUMERIC, _ingested_at TIMESTAMPTZ DEFAULT now()
);

-- 디멘전
CREATE TABLE IF NOT EXISTS dim_customer(
  customer_id BIGSERIAL PRIMARY KEY,
  email TEXT UNIQUE, domain TEXT, created_at TIMESTAMPTZ DEFAULT now()
);

CREATE TABLE IF NOT EXISTS dim_product(
  product_id BIGSERIAL PRIMARY KEY,
  sku TEXT UNIQUE, name TEXT, created_at TIMESTAMPTZ DEFAULT now()
);

-- 팩트
CREATE TABLE IF NOT EXISTS fct_order(
  order_id TEXT PRIMARY KEY, order_ts TIMESTAMPTZ,
  customer_id BIGINT REFERENCES dim_customer(customer_id),
  product_id BIGINT REFERENCES dim_product(product_id),
  qty INT, price NUMERIC
);

-- 정규화: 고객·제품 UPSERT
INSERT INTO dim_customer(email, domain)
SELECT DISTINCT customer_email, split_part(customer_email,'@',2) AS domain
FROM stg_orders s
ON CONFLICT (email) DO NOTHING;

INSERT INTO dim_product(sku, name)
SELECT DISTINCT product_sku, NULL
FROM stg_orders
ON CONFLICT (sku) DO NOTHING;

-- 팩트 MERGE(UPSERT)
INSERT INTO fct_order(order_id, order_ts, customer_id, product_id, qty, price)
SELECT s.order_id, s.order_ts, c.customer_id, p.product_id, s.qty, s.price
FROM stg_orders s
JOIN dim_customer c ON c.email = s.customer_email
JOIN dim_product p ON p.sku = s.product_sku
ON CONFLICT (order_id) DO UPDATE
SET order_ts = EXCLUDED.order_ts,
    customer_id = EXCLUDED.customer_id,
    product_id = EXCLUDED.product_id,
    qty = EXCLUDED.qty,
    price = EXCLUDED.price;

운영 절차는 BEGIN 트랜잭션으로 시작해 COPY로 스테이징에 적재하고, 검증 실패 시 ROLLBACK, 성공 시 COMMIT하는 방식이다.

품질과 복구 가능성을 함께 설계할 때

데이터 품질을 위해 스키마 레지스트리, 유닛·샘플 검증, DLQ 격리를 적용할 수 있다. 검증을 엄격하게 만들수록 지연은 증가하고, 완화할수록 오염 위험은 커진다.

증분 적재는 변경 시각(updated_at), 해시, CDC 로그를 기반으로 설계한다. 증분 키가 없으면 풀로드가 필요하며, 복구 용이성과 저장비용 사이의 선택도 남는다.

성능과 비용 측면에서는 컬럼형 DWH(Snowflake/BigQuery/Redshift)와 대량 적재(COPY/LOAD)를 조합할 수 있다. 대량 배치의 효율과 지연 요구를 함께 검토하고, 필요하면 마이크로배치와 스트리밍을 혼합한다.

비밀 관리(Vault/Secrets Manager), PII 마스킹·토큰화, 최소 권한은 보안과 컴플라이언스의 기본 요소다. 암호화와 권한 세분화는 운영 복잡성을 높일 수 있다.

재시도·백오프·아이들포턴시(idempotency) 키와 멱등 MERGE는 신뢰성을 위한 설계다. 멱등성 구현 비용과 장애 복구 속도 사이의 트레이드오프를 고려해야 한다.

데이터 품질 측면에서는 중복률 80% 이상 감소와 스키마 위반 90% 이상 조기 탐지를 기대할 수 있다. 증분 적재는 네트워크/쿼터 비용 4070% 절감, 배치 처리시간 3060% 단축으로 이어질 수 있다. 라인리지와 감사 체계를 갖추면 감사 대응 시간 50% 이상 단축과 서비스 가용성(SLA) 준수율 상승을 기대할 수 있다.

ETL데이터 파이프라인데이터 웨어하우스정규화API 연동