데이터는 어디서 오는가
| 소스 | 형태 | 수집 방법 |
|---|---|---|
| 파일 | CSV, JSON, Excel | 다운로드, SFTP, 공유 스토리지 |
| 데이터베이스 | 테이블 | SQL 쿼리, CDC(변경 데이터 캡처) |
| API | JSON over HTTP | requests, 페이지네이션, 인증 토큰 |
| 로그/이벤트 | 한 줄씩 계속 생성 | 에이전트 → 메시지 큐 |
| 센서/IoT | 스트림 | MQTT → Kafka |
API로 수집하기
import requests, pandas as pd
url = "https://api.example.com/items"
params = {"page": 1, "size": 100}
headers = {"Authorization": "Bearer <TOKEN>"}
rows = []
while True:
r = requests.get(url, params=params, headers=headers, timeout=10)
r.raise_for_status()
data = r.json()
rows.extend(data["items"])
if not data.get("next"):
break
params["page"] += 1
df = pd.DataFrame(rows)
df.to_parquet("items.parquet")
핵심 포인트: 페이지네이션, 재시도/타임아웃, rate limit, 원본 그대로 저장(raw) 후 가공.
DB에서 수집하기
import pandas as pd
from sqlalchemy import create_engine
engine = create_engine("postgresql://user:pw@host/db")
df = pd.read_sql("SELECT * FROM orders WHERE created_at >= '2026-01-01'", engine)
전체를 매번 읽지 않고 증분(incremental) 으로 — 마지막 수집 시각 이후만 가져오는 것이 배치 수집의 기본.
메시지 큐와 스트림
로그나 이벤트처럼 끝없이 생성되는 데이터는 파일로 모으기 어렵습니다. 그래서 중간에 버퍼(큐) 를 둡니다.
생산자(Producer) ──▶ [ Kafka Topic ] ──▶ 소비자(Consumer)
웹서버, 앱, 센서 파티션으로 분산 Spark, DB 적재기
- Topic: 이벤트의 카테고리 (예:
click-events) - Partition: 토픽을 나눠서 병렬 처리 — HDFS 블록과 같은 발상
- Offset: 소비자가 어디까지 읽었는지
- 생산자와 소비자가 분리(decoupling) 되어 서로 속도가 달라도 됨
연습 과제
- 공공데이터포털에서 API 키 발급 →
requests로 100건 이상 수집 → Parquet 저장 - 수집 코드를 "어제 이후 데이터만" 가져오도록 증분 수집으로 바꿔보기
- (선택) Docker로 Kafka 띄우고
kafka-python으로 메시지 1개 보내고 받기