ETL 단계와 실패 경계
120분 안팎
학습 목표
m04 정제 함수를 변환 단계로 재사용하고 로컬 SQLite 저장 전에 스키마 오류를 차단합니다.
개념
실패한 단계 다음으로 진행하지 않습니다
수집과 정제가 한 함수 안에 섞여 있으면 원본을 받다가 실패한 것인지 DB 저장 중 실패한 것인지 알기 어렵습니다. 이번 레슨에서는 앞 레슨의 원본 보존과 스키마 검사를 연결하고 검증을 통과한 행만 저장합니다. run은 collect, transform, load 순서로 호출합니다. 중간 함수에서 예외가 올라오면 다음 호출에 도달하지 않으므로 열 누락 파일이 기존 저장소를 지우는 일을 막습니다.
새 ETL은 m06의 지표나 보고서 로직을 다시 쓰지 않습니다. 미션 ZIP의 기존 파일과 테스트를 유지하면서 etl 패키지를 추가합니다. 입력은 이전 단계에서 만든 joined.csv이며 그림·지표·정제 결과는 근거 자료로 계속 남습니다. 여기서 새로 만드는 etl-out은 저장용 변환 결과입니다. 분석 작업을 단계로 나누는 일과 이미 검증된 해석을 바꾸는 일을 섞지 않습니다.
단계의 반환값을 다음 단계의 입력으로 연결합니다
collect_file은 joined.csv의 바이트를 내용 주소 폴더에 보존하고 payload.csv 경로를 반환합니다. transform은 그 경로에서 CSV를 읽어 열 목록을 확인하고 normalize로 타입을 변환합니다. errors로 숫자와 날짜 계약을 검사한 뒤 classify로 중복과 충돌을 구분합니다. 정상 행 목록을 반환하고 cleaning.run으로 정제 파일과 감사 자료를 etl-out에 기록합니다.
m04의 classify는 잘못된 행을 rejected 목록에 남기고 나머지를 계속 정제합니다. 이번 저장 흐름에서는 rejected가 하나라도 있으면 REJECTED 오류로 멈춥니다. 검증된 저장 스냅샷을 만드는 목적에 맞춰 정책을 강화한 것입니다. 완전히 같은 중복 행은 한 행으로 정리하지만 같은 지역·날짜의 서로 다른 값은 KEY_CONFLICT로 차단합니다. 원래 정제 함수의 구현을 복사해 수정하지 않고 호출자가 정책을 정합니다.
통행량이 없고 valid_observations가0인 경우는 정상 결측이며 정제 결과에 남습니다. 통행량0과 관측 수1은 실제 영점 관측입니다. rain_mm만 결측인 행도 이 저장 단계에서 보존합니다. 강수 그룹 분석에서 제외한 행을 저장에서도 모두 버리는 것으로 오해하지 않습니다. 저장 표는 후속 질문에 사용할 관측을 보존하며 분석별 포함 조건은 지표 계산이 담당합니다.
pipeline.json을 실행 계약으로 사용합니다
설정에는 source, snapshot, transformed, database 경로와 collect→transform→load 순서, version 및 입력 header를 기록합니다. run은 이 내용을 읽고 허용한 버전과 순서인지 확인합니다. 단순한 설명 문서로만 두면 코드 경로를 바꾼 뒤 설정이 낡아도 실행이 이어집니다. 설정 검사 테스트는 순서를 뒤집거나 입력과 출력 경로를 같게 했을 때 실행 전에 실패하는지 확인합니다.
경로는 프로젝트 루트 기준의 상대 경로를 씁니다. 절대 경로나 .. 경로를 거부하고 resolve한 결과가 루트 아래에 있는지도 확인합니다. 출력 폴더가 입력 파일의 부모이거나 다른 출력의 하위인 경우도 겹침으로 거부합니다. 경로가 모두 서로 다른 문자열이라고 안전한 것은 아닙니다. 같은 실제 위치를 가리키는 이름이나 입력을 포함하는 출력 폴더를 확인해야 원본 보존이 가능합니다.
설정의 schema.header는 m04의 아홉 열과 일치해야 합니다. 실제 CSV의 fieldnames도 같은 이름과 순서인지 검사합니다. 헤더가 맞아도 각 행의 열 수가 다를 수 있으므로 DictReader가 만든 None 키나 None 값을 확인합니다. HEADER 오류와 CSV 오류를 구분하면 제공자에게 열 이름이 달라진 것인지 한 행의 구분자가 깨진 것인지 구체적으로 문의할 수 있습니다.
저장 범위를 작고 명확하게 정합니다
load는 region, date, total_vehicles, rain_mm 네 열의 observations 테이블에 저장합니다. 지역·날짜가 기본키이며 None은 SQLite NULL로 전달합니다. 매 실행에서 검증된 전체 스냅샷으로 내용을 교체하는 정책입니다. 모든 과거 기간을 누적하는 수집기로 쓰면 안 됩니다. 이번 표본은 전체 교체이고 다음 모듈에서 증분 키, upsert와 재실행 정책을 설계합니다.
값은 SQL 문자열에 직접 붙이지 않고 자리 표시자와 executemany로 전달합니다. DB를 열기 전에 행 스키마를 다시 검사하고 연결의 트랜잭션 안에서 DELETE와 INSERT를 수행합니다. 중복 키 INSERT가 실패하면 기존 행 삭제까지 되돌아가는지 테스트합니다. with db는 트랜잭션의 완료·되돌림을 맡으며 연결 닫기는 closing으로 따로 보장합니다. 두 역할을 같은 것으로 이해하지 않습니다.
열 누락처럼 transform에서 발견하는 오류는 load 호출 전에 발생하므로 기존 DB 파일을 열지도 않습니다. 이 경우 DB 바이트를 전후 비교할 수 있습니다. 저장 중 제약 위반은 롤백 뒤 논리적 행이 유지되는지 확인합니다. 롤백된 SQLite 파일의 모든 바이트가 원래와 같다는 더 강한 주장을 하지 않습니다. 단계별로 무엇을 보존하는지에 맞게 검사 기준을 다르게 정합니다.
산출물마다 실패 시 상태를 설명합니다
수집이 성공한 뒤 변환이 실패하면 새 원본 스냅샷은 남고 DB는 기존 상태입니다. 이는 실패 원인을 재현할 증거입니다. 한 실행의 모든 파일이 동시에 교체되는 전체 파이프라인 트랜잭션을 구현한 것은 아닙니다. 변환 출력 기록 중 디스크 오류가 발생하면 일부 변환 파일이 남을 수 있으나 DB 저장은 호출되지 않습니다. 원본, 임시 결과, 소비자가 읽는 DB의 실패 경계를 구분해 인계합니다.
실행이 끝나면 loaded_rows와 실제 observations 행 수를 대조합니다. 기본 표본은 여섯 지역·날짜이며 통행량 결측 한 행과 강수 결측 두 행을 포함합니다. 같은 입력을 두 번 실행해 여섯 행이 유지되고 내용 주소 원본이 하나만 있는지도 확인합니다. 이 결과는 이번 전체 교체 정책의 반복 성질이며 모든 수집 방식의 중복 방지가 해결됐다는 뜻은 아닙니다.
미션 starter는 수집과 변환을 실행하지만 load를 호출하지 않아 저장 검사가 실패합니다. TODO를 채울 때 count만6으로 바꾸지 않습니다. 반환한 숫자와 실제 저장 행을 함께 검사하므로 표면적인 성공 메시지는 요구를 만족하지 못합니다. 기존 60개 회귀 검사와 새 실패 경계 검사까지 통과시키고 pipeline.json, 실행 명령, 정상·실패 상태를 설명하는 README를 제출합니다. 더 읽기는 모듈 간 의존성과 실행 진입점의 설계를 확장합니다.
따라하기
m04 정제 함수로 변환
미션 ZIP에서 제공하는 기존 normalize를 사용합니다. 다음 코드는 따라하기를 독립 실행하기 위해 함수 구현을 함께 싣습니다. 저장 단위에서는 cleaning.py에서 import합니다.
"""결합 표를 정제합니다. 입력은 보존하고 결측값은 None으로 다룹니다."""
import argparse
import csv
import json
import math
import re
from datetime import datetime
from pathlib import Path
HEADER = ['region','date','total_vehicles','valid_observations','weather_region','rain_mm','station_rows','valid_stations','join_status']
def integer(text, nullable=False):
text = text.strip()
if nullable and text == '':
return None
if not re.fullmatch(r'[0-9]+', text):
raise ValueError('INTEGER')
return int(text)
def day(text):
text = text.strip()
for fmt, pattern in [('%Y-%m-%d', r'[0-9]{4}-[0-9]{2}-[0-9]{2}'), ('%Y/%m/%d', r'[0-9]{4}/[0-9]{2}/[0-9]{2}')]:
if re.fullmatch(pattern, text):
try:
return datetime.strptime(text, fmt).date().isoformat()
except ValueError:
break
raise ValueError('DATE')
def normalize(row):
result = dict(row)
result['region'] = row['region'].strip()
if not result['region']:
raise ValueError('KEY')
result['date'] = day(row['date'])
for key in ['total_vehicles','valid_observations','station_rows','valid_stations']:
try:
result[key] = integer(row[key], key != 'valid_observations')
except ValueError:
raise ValueError('INTEGER:' + key) from None
rain = row['rain_mm'].strip()
try:
rain = None if rain == '' else float(rain)
except ValueError:
raise ValueError('RAIN') from None
if rain is not None and (not math.isfinite(rain) or rain < 0):
raise ValueError('RAIN')
result['rain_mm'] = rain
if (result['total_vehicles'] is None) != (result['valid_observations'] == 0):
raise ValueError('CONTRACT:traffic')
result['weather_region'] = row['weather_region'].strip()
result['join_status'] = row['join_status'].strip()
if result['join_status'] not in {'matched','mapping_missing','weather_missing','rain_missing'}:
raise ValueError('STATUS')
return result
raw = {'region':'A','date':'2026/09/01','total_vehicles':'0','valid_observations':'1','weather_region':'WA','rain_mm':'','station_rows':'1','valid_stations':'0','join_status':'rain_missing'}
row = normalize(raw)
print(row['date'],row['total_vehicles'],row['rain_mm'])실행 결과
2026-09-01 0 None
검증 실패 때 load 차단
예외가 발생한 뒤 저장 함수가 호출되지 않는지 확인합니다.
calls=[]
def collect():
calls.append('collect')
return 'snapshot.csv'
def transform(path):
calls.append('transform')
raise ValueError('HEADER: missing total_vehicles')
def load(rows):
calls.append('load')
try:
rows=transform(collect())
load(rows)
except ValueError as error:
print(error)
print('called', ','.join(calls))실행 결과
HEADER: missing total_vehicles called collect,transform
SQLite 전체 교체와 실패 되돌림
중복 키 실패 후 이전 논리적 행이 남습니다. 실패한 새 행을 저장 성공으로 보고하지 않습니다.
import sqlite3
from contextlib import closing
with closing(sqlite3.connect(':memory:')) as db:
db.execute('CREATE TABLE observations(region TEXT,date TEXT,total_vehicles INTEGER,PRIMARY KEY(region,date))')
with db:
db.execute('INSERT INTO observations VALUES(?,?,?)',('A','2026-09-01',100))
try:
with db:
db.execute('DELETE FROM observations')
db.executemany('INSERT INTO observations VALUES(?,?,?)',[('A','2026-09-02',0),('A','2026-09-02',80)])
except sqlite3.IntegrityError:
print('rollback')
print(db.execute('SELECT region,date,total_vehicles FROM observations ORDER BY region,date').fetchall())실행 결과
rollback
[('A', '2026-09-01', 100)]
ZIP 검사와 실패 위치 확인
starter.zip을 푼 루트에서 아래 명령을 실행합니다. 기존 60개와 새 20개, 총 80개 기본 테스트 통과입니다. HTTP 통합 9개는 기본 건너뜀입니다. starter는 실제 저장을 요구하는 검사가 실패합니다. 실패의 actual과 expected를 읽고 TODO를 채운 뒤 다시 검사합니다. HTTP 통합은 이 작성 샌드박스의 포트 권한 때문에 실행하지 못했으며 기본 검사와 구분합니다.
python3 -m unittest discover -s tests -v전체 표본 실행과 저장 행 수
solution 또는 TODO를 완성한 미션 ZIP 루트에서 python3 -m etl.run을 실행합니다. 기본 설정은 네트워크 없이 joined.csv를 읽으며 다음은 같은 ZIP을 임시 폴더에 복사해 실제 실행한 결과입니다. 수집 바이트 수와 저장 행 수를 확인하고 etl.sqlite의 observations를 조회해 결측과 영점을 대조합니다.
실행 결과
{"collected_bytes": 336, "loaded_rows": 6}
확인 문제
실습
m06 solution을 포함한 ZIP의 etl/run.py에서 load 호출을 연결합니다. pipeline.json에 지정한 경로를 사용합니다. 정상 6행, 영점과 NULL, 같은 입력 반복, 열 누락·날짜 오류·키 충돌 시 저장 차단, 기존 DB와 원본 보존 검사를 통과시킵니다. 기존 코드와 보고서는 유지합니다. 실제 HTTP 통합은 BOOTCAMP_HTTP=1로 별도 실행합니다.
실행 명령
python3 -m unittest discover -s tests -v
기대 결과
기존 60개와 새 20개, 총 80개 기본 테스트 통과입니다. HTTP 통합 9개는 기본 건너뜀입니다. starter는 실제 저장을 요구하는 검사가 실패합니다.
모범 답안
모범 답안 내려받기더 읽기
면접 질문
- 매일 같은 파일을 읽는 처리에서 중복 적재를 막는 방법을 설명해 주시면 됩니다.
- 보고서의 수치와 원본 데이터가 다를 때 확인할 순서를 설명해 주시면 됩니다.