← 목록으로
단계 04

데이터 수집

파일·DB·API·로그에서 데이터를 모으는 원리, 메시지 큐(Kafka)와 스트림 개념.

데이터는 어디서 오는가

소스형태수집 방법
파일CSV, JSON, Excel다운로드, SFTP, 공유 스토리지
데이터베이스테이블SQL 쿼리, CDC(변경 데이터 캡처)
APIJSON over HTTPrequests, 페이지네이션, 인증 토큰
로그/이벤트한 줄씩 계속 생성에이전트 → 메시지 큐
센서/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) 되어 서로 속도가 달라도 됨

연습 과제

  1. 공공데이터포털에서 API 키 발급 → requests로 100건 이상 수집 → Parquet 저장
  2. 수집 코드를 "어제 이후 데이터만" 가져오도록 증분 수집으로 바꿔보기
  3. (선택) Docker로 Kafka 띄우고 kafka-python으로 메시지 1개 보내고 받기