본문 바로가기
C.W.K.
Stream
Lesson 06 of 06 · published

Observability — 로그, metric, run log

~11 min · observability, logging, monitoring

Level 0구경꾼
0 XP0/47 lessons0/11 achievements
0/120 XP to next level120 XP to go0% complete

안 보이는 건 못 고쳐

"내 파이프라인 돌아가" 에서 "내 파이프라인은 들여다볼 수 있어" 로 가는 건, 스크립트에서 시스템으로 가는 것과 같은 걸음이야. Observability 는 세 겹이고, 겹마다 막아 주는 사고의 종류가 달라.

  • 로그 — 무슨 일이 있었는지, 순서대로, 사람이 읽는 형태로. 디버깅용.
  • Metric — stage 마다 흘려보내는 count, duration, rate. 추세와 알림용.
  • Run log — 시작 시간, 종료 시간, 성공/실패, stage 별 in/out row count, 에러를 run 단위로 남긴 기록. 사후 조사와 SLA 추적용.

구조화 로그엔 타협이 없어

일반 텍스트 로그 ('started extract') 는 터미널을 지켜보는 사람한테나 충분해. Production 급 로그는 구조화 (key=value 또는 JSON) 되어 있어서 log aggregator 가 인덱싱하고, 패턴에 알림을 걸고, run 을 가로질러 query 할 수 있어. Python logging 모듈로 가는 길은 짧아 — 설정 두 줄, 그리고 중요한 값은 꼭 실어 보내는 습관.

Code

규모가 커져도 버티는 구조화 로깅 설정·python
import json
import logging
from datetime import datetime, timezone

class JsonFormatter(logging.Formatter):
    def format(self, record: logging.LogRecord) -> str:
        # naive 말고 aware 로. utcnow() 는 Python 3.12 부터 deprecated 고,
        # timezone 을 아예 안 달고 오는 datetime 을 줘. aware 의 isoformat()
        # 은 +00:00 으로 끝나니까 그걸 로그 파이프라인이 원하는 Z 로 바꿔.
        # naive 문자열 뒤에 Z 만 붙여 놓고 UTC 라고 우기면 안 되고.
        ts = datetime.now(timezone.utc).isoformat().replace('+00:00', 'Z')
        payload = {
            'ts': ts,
            'level': record.levelname,
            'name': record.name,
            'msg': record.getMessage(),
        }
        # logger.info(..., extra={'kv': {...}}) 로 넘긴 추가물
        if hasattr(record, 'kv'):
            payload.update(record.kv)
        return json.dumps(payload, ensure_ascii=False)

handler = logging.StreamHandler()
handler.setFormatter(JsonFormatter())
logging.basicConfig(level=logging.INFO, handlers=[handler])
log = logging.getLogger('orders_pipeline')

log.info('extract complete', extra={'kv': {
    'stage': 'extract', 'rows': 12_345, 'ms': 1834,
}})
Run-log 레코드 — 대시보드가 query 하는 산출물·python
import time
from dataclasses import dataclass, field

@dataclass
class RunLog:
    pipeline: str
    window: str
    started_at: float = field(default_factory=time.time)
    stages: dict = field(default_factory=dict)
    success: bool = False
    error: str | None = None
    finished_at: float | None = None

    def stage(self, name: str, **kv) -> None:
        self.stages[name] = {'ms': int(kv.pop('ms', 0)), **kv}

    def to_row(self) -> dict:
        return {
            'pipeline':    self.pipeline,
            'window':      self.window,
            'started_at':  self.started_at,
            'finished_at': self.finished_at,
            'duration_s':  (self.finished_at or time.time()) - self.started_at,
            'success':     self.success,
            'error':       self.error,
            **{f'stage_{k}': v for k, v in self.stages.items()},
        }

External links

Exercise

본인 파이프라인 하나에 구조화 JSON 로깅과 run 단위 run-log 레코드를 붙여. 다음 run 이 끝나면 run-log row 를 Parquet 이나 DuckDB 테이블에 써. 그리고 duration 을 시간 위에 그려 봐. 그 그림이 바로 "이 파이프라인 느려지고 있어?" 의 답이야 — observability 없이는 답할 수 없는 질문이고.

Progress

Progress is local-only — sign in to sync across devices.
이 페이지에서 버그를 발견하셨거나 피드백이 있으세요?문제 신고

댓글 0

🔔 답글 알림 (로그인 필요)
로그인댓글을 남기려면 로그인해 주세요.

아직 댓글이 없어요. 첫 댓글을 남겨보세요.