# DMF Crawler 데이터 모델과 변경 탐지 스펙 > **이 문서의 역할**: `docs/design/01-architecture.md` 가 확정한 모듈 경계 안에서, **무엇을 하나의 레코드로 볼 것인가 · 그것을 어떻게 저장할 것인가 · 어제와 오늘을 어떻게 비교할 것인가**를 실행 가능한 코드와 SQL 로 확정한다. `normalize.py` · `integrity.py` · `diff.py` · `storage/` 이하 전부의 구현 정본이며, 이 문서와 코드가 어긋나면 **코드가 틀린 것**이다. 데이터 소스의 필드 정의는 `docs/design/00-DATA-SOURCE-DECISION.md`, 도메인 근거는 `docs/research/01-dmf-domain-and-sources.md` §3·§10 을 따른다. --- ## 0. 한눈에 보기 이 문서가 확정하는 것. 1. **표준 엔티티는 `DmfRecord` 하나다.** 공식 Open API 가 주는 원천 7필드에 파생 21필드를 더해 28필드로 확정했다. 아키텍처 §3.5 의 `DmfRecord` 를 **파생 필드만 추가하는 방향으로 확장**했고, 추가분은 전부 `COMPARE_FIELDS` 밖이라 diff 판정에 영향을 주지 않는다. 2. **등록번호는 자연 키가 아니다.** 의미를 6토큰으로 완전히 해부했지만(`등록수리일자-별표1성분번호-시행일군-군내순번-동일성분순번(허여서순번)`), 신물질 포맷 공존·표기 흔들림·중복 가능성 때문에 **원문(`permit_no_raw`) / 정규화값(`permit_no`) / 대체키(`dmf_key`) 3층 구조**를 채택한다. 파서는 예외를 던지지 않고 `fmt="unknown"` 으로 안전 착지한다. 3. **스냅샷은 멤버십 테이블, 내용은 버전 테이블로 분리한다.** 매일 전량 1만 행을 통째로 적재하면 연 1.4GB 인데, `record_versions`(내용) + `snapshots`(run_id↔dmf_key↔content_hash 3열) 로 쪼개면 **연 150MB 이하**로 떨어지고 필드 이력 조회가 SQL 한 줄이 된다. 아키텍처 §3.8 의 "영구 보존" 을 실현 가능하게 만드는 유일한 형태다. 4. **해시는 두 개다.** `content_hash`(표시값 기준, 스냅샷 버전 식별)와 `identity_hash`(정규화 키 기준, 실질 변경 판정). 원본 오탈자 정정 같은 표기만의 변화는 `COSMETIC` 등급으로 기록만 하고 리포트 전면에 올리지 않는다 — 요구 R2.3 의 구현체다. 5. **이벤트 타입은 `NEW` / `CHANGED` / `WITHDRAWN` 3종으로 고정한다.** 리서치가 제안한 9종(`SITE_CHANGED`, `ANNUAL_REPORT` 등)은 **공식 API 가 그 필드를 주지 않으므로 채택하지 않는다.** 대신 `CHANGED` 이벤트의 **필드별 중요도 5등급**(`CRITICAL`/`HIGH`/`MEDIUM`/`INFO`/`COSMETIC`)으로 같은 정보를 표현한다. 재등장은 `NEW` + `is_reappearance=1` 로 표현한다. 6. **취하 오판을 막는 관문은 차단 게이트 5종 + 비차단 품질 지표 13종**이다. 게이트가 하나라도 실패하면 스냅샷 INSERT 도 하지 않는다(기준선 오염 방지). 급변 임계(`churn_ratio`)를 게이트 밖 6번째 안전장치로 신설한다. 7. **중복 등록번호는 집합 매칭으로 처리한다.** 같은 `permit_no` 가 둘 이상이면 `#2`·`#3` 접미사를 붙이되, diff 는 접미사를 짝짓지 않고 **그룹 대 그룹 탐욕 매칭**을 수행한다. 접미사 순서가 흔들려도 오탐이 나지 않는다. 8. **보존 정책은 자산별로 다르다.** DB 이벤트·버전은 영구, `data/raw/` 180일, `logs/` 90일, `reports/` 365일, `backup/` 30개. 삭제는 `finalize` 스테이지에서만, 트랜잭션 밖에서, 실패해도 실행 상태를 끌어내리지 않는다. --- ## 1. 도메인 엔티티 정의 ### 1.1 원천 필드 — 공식 Open API 가 실제로 주는 것 (7개) `https://apis.data.go.kr/1471000/MdcDmfInfoService01/getMdcDmfList01` 의 데이터 항목 전부다. 이 7개가 이 프로젝트의 **유일한 원천**이며, 나머지는 전부 여기서 계산된다. | # | 원본 컬럼 | 한글 항목명 | 명세 크기 | 샘플값 | 널 가능 | 우리 필드로의 매핑 | |---|---|---|---|---|---|---| | 1 | `DMF_PERMIT_NO` | 등록번호 | 200 | `20121228-168-I-169-04` | 이론상 가능 ⚠️ | `permit_no_raw` → `permit_no` → `dmf_key` | | 2 | `INGR_KOR_NAME` | 성분명 | 3000 | `포르모테롤푸마르산염수화물` | 이론상 가능 ⚠️ | `ingredient_name` (+`ingredient_key`, `ingredient_base`, `is_micronized`) | | 3 | `ENTP_NAME` | 업체명 | 200 | `(주)대웅제약` | 이론상 가능 ⚠️ | `applicant` (+`applicant_key`) | | 4 | `MNFCTR_NAME` | 제조소명 | 150 | `SICOR SOCIETA'ITALIANA CORTICOSTER OIDO S.R.L.` | 가능 | `manufacturer` (+`manufacturer_key`, `sites`) | | 5 | `MNFCTR_PLACE` | 제조소 소재지 | 2000 | `Rho(MI) - Via Terrazzano, 77, Italy` | 가능 | `manufacture_place` | | 6 | `MANUF_COUNTRY_CODE_NM` | 제조국가명 | 1000 | `이탈리아,스위스` | 가능 | `countries` (콤마 분해 후 정렬된 튜플) | | 7 | `DMF_PERMIT_DATE` | 발급일자 | 30 | `2015-02-26` | 가능 | `permit_date` (ISO 정규화) | > **명세 크기의 이중 용도**: 위 "명세 크기" 는 데이터 품질 검증(§7)에서 **스키마 드리프트 탐지 임계값**으로 쓴다. `MNFCTR_NAME` 이 150자를 넘는 값이 갑자기 다수 등장하면 그것은 값 이상이 아니라 **원본 스키마가 바뀐 신호**다. > ⚠️ **널 가능성 표기 근거**: 포털 명세는 응답 데이터 항목에 필수 여부를 표기하지 않았다. 따라서 코드는 **7필드 전부 널일 수 있다고 가정**하고 방어하되, `permit_no`·`ingredient_name`·`applicant` 세 개는 §7 게이트 4 에서 널 비율 1% 상한을 강제한다. ### 1.2 API 가 주지 않는 것 — 그리고 그 결과 `docs/research/01-dmf-domain-and-sources.md` §12.1 은 의약품안전나라 화면·엑셀 기준 14필드를 조사했다. 공식 API 는 그중 7개만 준다. **없는 7개가 데이터 모델에 미치는 영향을 명시적으로 기록한다.** | 화면 컬럼 | 화면 필드명 | API 제공 | 없어서 잃는 것 | 이 문서의 대응 | |---|---|---|---|---| | 대상의약품 | `noticeCode` | ✕ | `별표1` / `신물질` 구분 | 등록번호 `fmt` 로 대리 판정 (`new_substance` ⇒ 신물질) | | 최종변경일자 | `yrycReportNDate` | ✕ | 변경등록 발생 시점 | **스냅샷 diff 로 대체.** 오히려 변경 *내용*까지 잡는다 | | 최종연차보고년도 | `yrycReportYYear` | ✕ | 연차보고 이벤트 | **미지원.** 리서치의 `ANNUAL_REPORT` 이벤트 타입 폐기 | | 취소/취하구분 | `cancelCodeNm` | ✕ | 명시적 취하 상태 | **목록 이탈(`WITHDRAWN`)로만 판정.** 오판 위험이 높아 §7 게이트 전면 배치 | | 취소/취하일자 | `cancelDate` | ✕ | 취하 일자 | 이벤트 발생일(`event_date`)로 대리 | | 문서번호 | `dmfVersion` | ✕ | 재공고 버전 추적 | **미지원.** `DOC_VERSION_CHANGED` 폐기 | | 연계심사문서번호 | `cntcJdgmnNo` | ✕ | 원료-완제 연계심사 신호 | **미지원.** `LINKED_REVIEW_SET` 폐기 | **결론 3가지.** - 취하 판정이 **목록 이탈에만 의존**한다. 이것이 이 프로젝트의 단일 최대 위험이며 §5.4 안전장치 전체가 이 한 문장 때문에 존재한다. - 리서치의 이벤트 9종은 3종으로 줄어든다. 정보 손실은 `CHANGED` 의 **필드별 중요도 등급**(§5.2)이 상당 부분 흡수한다. - 화면 7필드는 **버리지 않고 스키마에 예약 컬럼으로 남긴다**(§3.3 `record_versions` 의 `reserved_*` 없음 — 대신 `raw_json` 이 원본 전체를 보존하므로 소스가 확장되면 컬럼 추가 마이그레이션만 하면 된다). ### 1.3 표준 엔티티 `DmfRecord` — 전체 필드 정의 | # | 필드명 (snake_case) | 한글 표시명 | 타입 | 널 | 예시값 | 원본 매핑 | 정규화 규칙 | 비교 대상 | |---|---|---|---|---|---|---|---|---| | 1 | `dmf_key` | 레코드 키 | `str` | ✕ | `20121228-168-I-169-04` | 파생 | §4.1 3층 키. 중복 시 `#2` 접미사 | **키** | | 2 | `permit_no` | 등록번호(정규화) | `str` | ✕ | `20121228-168-I-169-04` | `DMF_PERMIT_NO` | NFKC · 전각괄호/하이픈 통일 · 전 공백 제거 · 대문자화 | ✕ | | 3 | `permit_no_raw` | 등록번호(원문) | `str` | ✕ | `20121228-168-I-169-04 ` | `DMF_PERMIT_NO` | 없음(감사용 원문 보존) | ✕ | | 4 | `ingredient_name` | 성분명 | `str` | ✕(빈문자 허용) | `포르모테롤푸마르산염수화물` | `INGR_KOR_NAME` | NFKC · 연속공백 1개 · 앞뒤 공백 제거 · 전각→반각 | **○** | | 5 | `ingredient_key` | 성분 매칭키 | `str` | ✕ | `포르모테롤푸마르산염수화물` | 파생 | `ingredient_name` 에서 **모든 공백·구두점 제거 후 소문자화** | ✕(등급 판정용) | | 6 | `ingredient_base` | 성분 기본명 | `str` | ✕ | `포르모테롤` | 파생 | `ingredient_key` 에서 미분화 접두 · 염/수화물/양이온 접미 반복 제거 | ✕(**집계 전용**) | | 7 | `is_micronized` | 미분화 여부 | `bool` | ✕ | `False` | 파생 | 성분명 `미분화`/`초미분화` 접두 또는 제조소명 `[미분화공정 제조소]` | ✕ | | 8 | `applicant` | 신청인(업체명) | `str` | ✕(빈문자 허용) | `(주)대웅제약` | `ENTP_NAME` | NFKC · 공백 축약 · `㈜`→`(주)` 통일 | **○** | | 9 | `applicant_key` | 신청인 매칭키 | `str` | ✕ | `대웅제약` | 파생 | 법인격(`주식회사`/`(주)`/`유한회사`/`Co.,Ltd` 등) 제거 + 공백·구두점 제거 + 소문자화 | ✕(등급 판정용) | | 10 | `manufacturer` | 제조소명 | `str` | ✕(빈문자 허용) | `SICOR SOCIETA'ITALIANA CORTICOSTER OIDO S.R.L.` | `MNFCTR_NAME` | NFKC · 공백 축약 · 다중값은 ` , ` 로 재조립 | **○** | | 11 | `manufacturer_key` | 제조소 매칭키 | `str` | ✕ | `sicorsocietaitalianacorticosteroidosrl` | 파생 | 전 공백·구두점 제거 + 소문자화 (원본 오탈자 `CORTICOSTER OIDO` 흡수) | ✕(등급 판정용) | | 12 | `manufacture_place` | 제조소 소재지 | `str` | ✕(빈문자 허용) | `Rho(MI) - Via Terrazzano, 77, Italy` | `MNFCTR_PLACE` | NFKC · 공백 축약 | **○** | | 13 | `manufacture_place_key` | 소재지 매칭키 | `str` | ✕ | `rhomiviaterrazzano77italy` | 파생 | 전 공백·구두점 제거 + 소문자화 | ✕(등급 판정용) | | 14 | `countries` | 제조국가 | `tuple[str, ...]` | ✕(빈튜플 허용) | `("스위스", "이탈리아")` | `MANUF_COUNTRY_CODE_NM` | 콤마/슬래시 분해 · 별칭 정규화 · **정렬** | **○** | | 15 | `sites` | 제조소 분해 | `tuple[Site, ...]` | ✕(빈튜플 허용) | `(Site(name=..., role="미분화공정 제조소"),)` | 파생 | ` ,`(공백+콤마) 우선 분해 · 역할 접두어 추출 | ✕(**집계 전용**) | | 16 | `permit_date` | 발급일자 | `str` | ✕(빈문자 허용) | `2015-02-26` | `DMF_PERMIT_DATE` | 8자리·점·슬래시 표기 → `YYYY-MM-DD`. 파싱 실패 시 빈 문자열 | **○** | | 17 | `accepted_date` | 등록수리일자 | `str \| None` | ○ | `2012-12-28` | 등록번호 1토큰 | `permit.accept_date` 의 별칭(아키텍처 §3.5 명칭 보존) | ✕ | | 18 | `permit.fmt` | 등록번호 포맷 | `str` | ✕ | `standard` | 파생 | `standard`/`new_substance`/`unknown` | ✕ | | 19 | `permit.ingr_no` | 별표1 성분번호 | `int \| None` | ○ | `168` | 등록번호 2토큰 | — | ✕ | | 20 | `permit.group` | 시행일 알파벳군 | `str \| None` | ○ | `I` | 등록번호 3토큰 | 대문자 1자 `A`~`Z` | ✕ | | 21 | `permit.group_table` | 근거 별표 | `str \| None` | ○ | `별표1` | 파생 | `GROUP_TABLE` 조회 | ✕ | | 22 | `permit.group_effective_from` | 군 시행일 | `str \| None` | ○ | `2013-01-01` | 파생 | `GROUP_TABLE` 조회. **알파벳순 ≠ 시행일순이므로 정렬은 이 값으로** | ✕ | | 23 | `permit.serial` | 군내 접수순번 | `int \| None` | ○ | `169` | 등록번호 4토큰 | — | ✕ | | 24 | `permit.sub` | 동일성분 일련번호 | `int \| None` | ○ | `4` | 등록번호 5토큰 | J군은 항상 `None` | ✕ | | 25 | `permit.grant` | 허여서 순번 | `int \| None` | ○ | `None` | 등록번호 괄호 | 자료공유허여서 파생 등록 | ✕ | | 26 | `permit.base_permit_no` | 기준 등록번호 | `str \| None` | ○ | `20121228-168-I-169-04` | 파생 | 괄호 제거값. 허여서 파생을 원 등록과 묶는 축 | ✕ | | 27 | `permit.ingr_group_key` | 성분 코호트 키 | `str \| None` | ○ | `168-I` | 파생 | `{ingr_no}-{group}`. "이 성분에 제조원이 몇 개 붙었나" | ✕ | | 28 | `content_hash` | 내용 해시 | `str` | ✕ | `9f2c1ab4e7d05631` | 파생 | `COMPARE_FIELDS` **표시값** SHA-256 앞 16자 | ✕(판정 입력) | | 29 | `identity_hash` | 실질 해시 | `str` | ✕ | `31a0c8ee45b7f912` | 파생 | `COMPARE_KEY_FIELDS` **정규화 키** SHA-256 앞 16자 | ✕(등급 판정) | | 30 | `raw` | 원본 필드 | `dict[str, str]` | ✕ | `{"DMF_PERMIT_NO": "...", ...}` | 응답 그대로 | 없음 | ✕ | **비교 대상 6필드(`COMPARE_FIELDS`)를 이 6개로 고정한 이유**: 아키텍처 §3.5 가 확정한 목록 그대로다. `dmf_key`·`permit_no` 는 키이므로 비교 대상이 아니고, 나머지는 전부 이 6개에서 계산된 파생값이라 중복 비교가 된다. `raw` 는 원본 노이즈(응답 필드 순서·공백)를 그대로 담으므로 비교하면 매일 CHANGED 가 터진다. ### 1.4 아키텍처 §3.5 대비 변경점 이 문서는 아키텍처를 **확장**하되 **어떤 이름도 바꾸지 않는다.** 변경분을 전부 나열한다. | 항목 | 아키텍처 §3.5 | 이 문서 | 사유 | |---|---|---|---| | `DmfRecord.permit_no_raw` | 없음 | 추가 | 정규화 전 원문 보존. 감사·회귀 재현에 필요 | | `DmfRecord.ingredient_key` / `applicant_key` / `manufacturer_key` / `manufacture_place_key` | 없음 | 추가 | 요구 R2.3(표기 흔들림) 판정과 워치리스트 매칭의 축 | | `DmfRecord.ingredient_base` / `is_micronized` / `sites` | 없음 | 추가 | 리포트 `s03_ingredient` · `s04_company` 집계 전용. **diff 미참여** | | `DmfRecord.permit` (`PermitParts`) | `accepted_date` 만 | 구조체로 확장, `accepted_date` 는 별칭으로 유지 | 등록번호 6토큰이 리포트 파생 지표 5종의 원천(리서치 §3.6) | | `DmfRecord.identity_hash` | 없음 | 추가 | 표기만의 변화(`COSMETIC`)를 실질 변경과 분리 | | `make_dmf_key(permit_no, dup_index)` | 2인자 | `fallback_seed` 3번째 인자 추가(기본 `None`) | 등록번호가 빈 값인 레코드의 합성키 경로 | | `normalize.py` 의 책임 | 표준화 + 키·해시 | + 등록번호 파싱 | 파서를 별도 모듈로 만들면 아키텍처의 디렉터리 트리를 깬다. `normalize.py` 안에 둔다 | | 마이그레이션 `0001_init.sql` 테이블 | runs/stage_status/fetch_stats/snapshots/records/events | + `schema_version`, `sources`, `integrity_gates`, `record_versions` | §3.2 의 저장량 산정과 §7 결과 조회 필요. 전부 코어 테이블이라 0001 에 둔다 | | 보존 정리 함수 위치 | 명시 없음 | `backup.py` 의 `prune_*` 함수군 + `repo.prune_snapshots` | `backup.py` 가 이미 "보존 개수 정리" 를 맡는다(§8.2) | ### 1.5 `DmfRecord` 정의 코드 ```python # src/dmf_crawler/normalize.py (1/3 — 엔티티 정의) """원본 7필드를 표준 DmfRecord 로 변환한다. 이 모듈은 순수 함수만 담는다. DB·네트워크·파일에 접근하지 않는다. 표준 라이브러리 외 의존성이 없다(unicodedata, hashlib, re, dataclasses). """ from __future__ import annotations import hashlib import re import unicodedata from dataclasses import dataclass, field from typing import Iterable, Mapping, Optional # -------------------------------------------------------------------------- # 원본 컬럼명 — API 응답 키. 이 상수 밖에서 문자열 리터럴로 쓰지 않는다. # -------------------------------------------------------------------------- COL_PERMIT_NO = "DMF_PERMIT_NO" COL_INGREDIENT = "INGR_KOR_NAME" COL_APPLICANT = "ENTP_NAME" COL_MANUFACTURER = "MNFCTR_NAME" COL_PLACE = "MNFCTR_PLACE" COL_COUNTRY = "MANUF_COUNTRY_CODE_NM" COL_PERMIT_DATE = "DMF_PERMIT_DATE" SOURCE_COLUMNS: tuple[str, ...] = ( COL_PERMIT_NO, COL_INGREDIENT, COL_APPLICANT, COL_MANUFACTURER, COL_PLACE, COL_COUNTRY, COL_PERMIT_DATE, ) # 포털 명세의 항목 크기. §7 스키마 드리프트 탐지 임계값으로 쓴다. SOURCE_MAX_LEN: Mapping[str, int] = { COL_PERMIT_NO: 200, COL_INGREDIENT: 3000, COL_APPLICANT: 200, COL_MANUFACTURER: 150, COL_PLACE: 2000, COL_COUNTRY: 1000, COL_PERMIT_DATE: 30, } # 널 비율을 강제하는 필수 필드 (게이트 4) REQUIRED_FIELDS: tuple[str, ...] = ("permit_no", "ingredient_name", "applicant") # 아키텍처 §3.5 확정 — 표시값 기준 비교 대상. 순서 고정(해시 입력 순서다). COMPARE_FIELDS: tuple[str, ...] = ( "ingredient_name", "applicant", "manufacturer", "manufacture_place", "countries", "permit_date", ) # 실질 변경 판정용 — 정규화 키 기준. COMPARE_FIELDS 와 1:1 대응. COMPARE_KEY_FIELDS: tuple[str, ...] = ( "ingredient_key", "applicant_key", "manufacturer_key", "manufacture_place_key", "countries", "permit_date", ) # 표시 필드 -> 대응하는 키 필드. §5.2 등급 판정에서 쓴다. FIELD_KEY_MAP: Mapping[str, str] = { "ingredient_name": "ingredient_key", "applicant": "applicant_key", "manufacturer": "manufacturer_key", "manufacture_place": "manufacture_place_key", "countries": "countries", "permit_date": "permit_date", } FIELD_LABELS_KO: Mapping[str, str] = { "ingredient_name": "성분명", "applicant": "신청인", "manufacturer": "제조소명", "manufacture_place": "제조소 소재지", "countries": "제조국가", "permit_date": "발급일자", } @dataclass(frozen=True, slots=True) class Site: """제조소 1건. manufacturer/manufacture_place/countries 를 축별로 분해한 결과.""" seq: int name: str address: str country: str role: str = "" # '미분화공정 제조소' 등 대괄호 역할 표기 @dataclass(frozen=True, slots=True) class PermitParts: """등록번호 파싱 결과. 어떤 입력이든 예외 없이 생성된다.""" raw: str normalized: str fmt: str # standard | new_substance | unknown accept_date: Optional[str] = None # ISO YYYY-MM-DD ingr_no: Optional[int] = None group: Optional[str] = None group_table: Optional[str] = None group_range: Optional[str] = None group_effective_from: Optional[str] = None serial: Optional[int] = None sub: Optional[int] = None grant: Optional[int] = None is_grant_derived: bool = False base_permit_no: Optional[str] = None ingr_group_key: Optional[str] = None @dataclass(frozen=True, slots=True) class DmfRecord: """DMF 등록 레코드 표준형. 아키텍처 §3.5 확장(파생 필드만 추가).""" dmf_key: str permit_no: str permit_no_raw: str ingredient_name: str ingredient_key: str ingredient_base: str is_micronized: bool applicant: str applicant_key: str manufacturer: str manufacturer_key: str manufacture_place: str manufacture_place_key: str countries: tuple[str, ...] sites: tuple[Site, ...] permit_date: str permit: PermitParts raw: dict[str, str] = field(default_factory=dict) dup_index: int = 0 content_hash: str = "" identity_hash: str = "" @property def accepted_date(self) -> Optional[str]: """아키텍처 §3.5 의 필드명을 보존하는 별칭.""" return self.permit.accept_date def compare_view(self) -> dict[str, str]: """COMPARE_FIELDS 를 문자열로 평탄화한다. 이벤트 before/after 의 정본 형태.""" return { "ingredient_name": self.ingredient_name, "applicant": self.applicant, "manufacturer": self.manufacturer, "manufacture_place": self.manufacture_place, "countries": ", ".join(self.countries), "permit_date": self.permit_date, } def key_view(self) -> dict[str, str]: """COMPARE_KEY_FIELDS 를 문자열로 평탄화한다. COSMETIC 판정 입력.""" return { "ingredient_key": self.ingredient_key, "applicant_key": self.applicant_key, "manufacturer_key": self.manufacturer_key, "manufacture_place_key": self.manufacture_place_key, "countries": "|".join(self.countries), "permit_date": self.permit_date, } @dataclass(frozen=True, slots=True) class NormalizeStats: """정규화 통계. integrity 가 게이트 4·5 판정에 그대로 쓴다.""" total_in: int total_out: int rejected: tuple[dict[str, str], ...] # 파싱 자체가 불가능했던 원본 null_counts: Mapping[str, int] # 필수 필드별 널 건수 duplicate_permit_no: int # 중복 등록번호로 접미사가 붙은 건수 duplicate_groups: int # 중복이 발생한 등록번호 종류 수 unparsed_permit_no: int # fmt == 'unknown' new_substance_count: int # fmt == 'new_substance' synthetic_key_count: int # 등록번호가 없어 합성키를 쓴 건수 invalid_permit_date: int # 발급일자 파싱 실패 oversize_fields: Mapping[str, int] # 명세 크기 초과 건수(스키마 드리프트) replacement_char_count: int # U+FFFD 포함 건수(인코딩 사고) ``` --- ## 2. 등록번호 파싱과 키 설계 ### 2.1 등록번호 체계 해부 근거는 식약처 FAQ(해설서 제5개정판 Ⅳ장 Q51) 원문이다. `20110531-71-B-317-05(1)` 을 6토큰으로 해부한다. | 위치 | 토큰 | 예시 | 의미 | 타입 | 규칙 | 파생 활용 | |---|---|---|---|---|---|---| | 1 | `accept_date` | `20110531` | **등록수리일자** | date | `(19\|20)\d{2}` + 월 + 일 | 접수 시점 시계열. `permit_date` 와 별개 필드로 저장 | | 2 | `ingr_no` | `71` | 「원료의약품 등록에 관한 규정」 **[별표1] 성분 일련번호** | int | 1~211+ | 성분 코호트 조인 키 | | 3 | `group` | `B` | **부칙 시행일 군** | char | `A`~`K` 관측 | 등록 의무 코호트 | | 4 | `serial` | `317` | 동일 군 내 **접수 순번** | int | 1~6자리 | 군 내 누적 순서 | | 5 | `sub` | `05` | **동일 성분 내 일련번호** | int? | **J군은 미부여** | 성분별 등록 경쟁 강도 | | 6 | `grant` | `(1)` | **허여서(자료공유허여서) 등록 순번** | int? | `(1)`~`(9)` | 원 등록번호와 반드시 묶어 표시 | **군 → 별표 → 시행일 매핑** (파서 내장 상수). 알파벳 순서와 시행일 순서가 일치하지 않으므로 **정렬 키로 알파벳을 쓰지 말고 시행일을 쓴다.** | 군 | 별표 | 호수 범위 | 시행일 | |---|---|---|---| | A~D | 별표1 | 제1호~제99호 | 2003-01-01 | | E | 별표1 | 제1호~제99호 | 2008-01-01 | | F | 별표1 | 제100호~제113호 | 2009-01-01 | | G | 별표1 | 제114호~제123호 | 2010-01-01 | | H | 별표1 | 제124호~제141호 | 2011-01-01 | | I | 별표1 | 제142호~제208호 | 2013-01-01 | | J | 별표1 | 제209호·제210호·제211호 | 2017-12-25 (211호는 2018-01-01) | | K | 별표1의2 | 제1호~제19호 | 2018-01-01 | **두 번째 포맷 — 신물질 계열.** 실측 `수6580-16-ND(20)` (대상의약품 = 신물질, 성분 = 프레가발린). FAQ Q51 은 별표1 포맷만 설명하므로 구조 해석은 ⚠️ **미검증**이다. 파서는 이 포맷을 표준 포맷으로 강제 파싱하지 않고 `fmt="new_substance"` 로 분기해 **토큰 의미를 부여하지 않는다.** 잘못된 의미 부여보다 값 없음이 낫다. ### 2.2 등록번호는 자연 키인가 — 판정 자연 키의 조건은 **유일성 · 불변성 · 존재성** 세 가지다. 각각을 검증한다. | 조건 | 검증 | 판정 | |---|---|---| | **유일성** | 등록번호는 접수 순번을 포함하므로 설계상 유일하다. 다만 API 가 실제로 중복을 반환하지 않는지는 ⚠️ **미검증**(serviceKey 발급 후 실측). 허여서 파생 `(1)` 은 괄호가 붙어 원 번호와 구별되므로 충돌하지 않는다 | **조건부 충족** | | **불변성** | 앞 8자리가 등록수리일자로 고정돼 있고, 재공고 시에도 등록번호는 유지된다(리서치 §2.9 의 "재공고" 건이 같은 번호로 재등장). 발급일자(`permit_date`)와 수리일자가 크게 벌어지는 사례(`20121228` vs `2015-02-26`)가 있으나 **번호 자체는 바뀌지 않는다** | **충족** | | **존재성** | 명세상 널 여부 표기가 없다. 빈 값이 올 가능성을 배제할 근거가 없다 | **미보장** | **결론: 등록번호는 "거의 자연 키"이지 자연 키가 아니다.** 세 가지가 걸린다. 1. **표기 흔들림** — 전각 하이픈(`-`), 전각 괄호, 비단절 공백(U+00A0), 앞뒤 공백이 섞이면 같은 등록이 다른 키가 된다. 그 즉시 `WITHDRAWN` + `NEW` 쌍으로 오탐한다. 2. **중복 가능성** — 미검증. 중복이 실재하면 PK 제약이 INSERT 를 깨뜨리고 그날 실행 전체가 실패한다. 3. **부재 가능성** — 빈 등록번호 레코드가 하나라도 오면 PK 가 성립하지 않는다. 따라서 **정규화 계층과 대체 키를 둔다.** 이것이 §2.3 의 3층 구조다. ### 2.3 3층 키 구조 ``` ┌ permit_no_raw ─ API 가 준 문자열 그대로. 절대 가공하지 않는다. 감사·회귀 재현용. │ │ NFKC · 전각→반각 · 괄호/하이픈 통일 · 전 공백 제거 · 대문자화 │ ▼ ├ permit_no ───── 정규화된 등록번호. 표기 흔들림이 제거된 "의미상의 등록번호". │ │ 중복 그룹 내 결정론적 접미사(#2, #3) / 부재 시 합성키 │ ▼ └ dmf_key ─────── 스냅샷 안에서 유일함이 보장되는 레코드 키. diff · PK · 리포트 링크의 축. ``` | 층 | 유일성 보장 | 사람이 읽는가 | 저장 위치 | 용도 | |---|---|---|---|---| | `permit_no_raw` | ✕ | ○ | `record_versions.permit_no_raw`, `raw_json` | 원본 대조, 재파싱 | | `permit_no` | 거의 (미검증) | ○ | 인덱스 있음 | 사람의 검색, 허여서 묶기(`base_permit_no`) | | `dmf_key` | **○ (강제)** | 대체로 | 모든 테이블의 조인 키 | **diff 키 · PK · 이벤트 키** | **대체 키 규칙 3가지.** - **정상**: `dmf_key = permit_no` - **중복**: 같은 `permit_no` 가 스냅샷 안에 n개면 결정론적 정렬 후 `permit_no`, `permit_no#2`, … `permit_no#n`. 정렬 키는 `(manufacturer_key, manufacture_place_key, countries 결합문자열, ingredient_key, permit_date)` 튜플이다. 아키텍처는 "제조소명 사전순" 이라고만 했는데, 제조소명 하나로는 동률이 나면 순서가 흔들리므로 **5요소 전체 튜플로 확장**한다. 그래도 흔들릴 수 있으므로 diff 에서 §4.5 그룹 매칭으로 이중 방어한다. - **부재**: `permit_no` 가 빈 문자열이면 `dmf_key = "SYN-" + sha1(ingredient_key|applicant_key|manufacturer_key|permit_date)` 앞 12자. 접두어 `SYN-` 로 합성키임을 리포트에서 즉시 식별할 수 있게 한다. 합성키 레코드는 **`applicant` 표기만 바뀌어도 신규+취하로 오탐**하므로 §7 품질 지표에서 건수를 별도 추적하고 리포트 메타 시트에 표시한다. ### 2.4 파서 완결 코드 ```python # src/dmf_crawler/normalize.py (2/3 — 등록번호 파서 · 정규화 원자 함수) # -------------------------------------------------------------------------- # 포맷 A: 일반(별표1 / 별표1의2) # 20110531-71-B-317-05 기본 # 20110531-71-B-317-05(1) 허여서 파생 # 20260901-209-J-2270 J군: sub 없음 # -------------------------------------------------------------------------- RE_PERMIT_STANDARD = re.compile( r"^(?P(?:19|20)\d{2}(?:0[1-9]|1[0-2])(?:0[1-9]|[12]\d|3[01]))" r"-(?P\d{1,4})" r"-(?P[A-Z])" r"-(?P\d{1,6})" r"(?:-(?P\d{1,3}))?" r"(?:\((?P\d{1,3})\))?$" ) # -------------------------------------------------------------------------- # 포맷 B: 신물질(신약 원료) 계열 — 예: 수6580-16-ND(20) # 토큰 의미는 미검증이므로 구조 인식만 하고 값 해석은 하지 않는다. # -------------------------------------------------------------------------- RE_PERMIT_NEW_SUBSTANCE = re.compile( r"^(?P[가-힣]{1,2})(?P\d{3,7})" r"-(?P\d{1,3})" r"-(?PND)" r"(?:\((?P\d{1,3})\))?$" ) # 군 -> (별표, 호수범위, 시행일) GROUP_TABLE: Mapping[str, tuple[str, str, str]] = { "A": ("별표1", "제1호~제99호", "2003-01-01"), "B": ("별표1", "제1호~제99호", "2003-01-01"), "C": ("별표1", "제1호~제99호", "2003-01-01"), "D": ("별표1", "제1호~제99호", "2003-01-01"), "E": ("별표1", "제1호~제99호", "2008-01-01"), "F": ("별표1", "제100호~제113호", "2009-01-01"), "G": ("별표1", "제114호~제123호", "2010-01-01"), "H": ("별표1", "제124호~제141호", "2011-01-01"), "I": ("별표1", "제142호~제208호", "2013-01-01"), "J": ("별표1", "제209호·제210호·제211호", "2017-12-25"), "K": ("별표1의2", "제1호~제19호", "2018-01-01"), } # 전각·유사 문자 → 반각 표준 문자 _PUNCT_FOLD = { "(": "(", ")": ")", "[": "[", "]": "]", "-": "-", "–": "-", "—": "-", "―": "-", "‐": "-", "‑": "-", "ー": "-", ",": ",", ".": ".", "/": "/", ":": ":", ";": ";", " ": " ", " ": " ", " ": " ", " ": " ", "": "", "​": "", "‌": "", "‍": "", } _WS_RE = re.compile(r"\s+") _NON_ALNUM_KO_RE = re.compile(r"[^0-9A-Za-z가-힣]+") def fold_punct(s: str) -> str: """전각·유사 구두점을 반각 표준으로 접는다.""" if not s: return "" return "".join(_PUNCT_FOLD.get(ch, ch) for ch in s) def clean_text(s: object) -> str: """모든 문자열 필드의 1차 정규화. NFKC -> 구두점 접기 -> 제어문자 제거 -> 연속 공백 1개 -> 앞뒤 공백 제거. None 이나 비문자열은 빈 문자열이 된다(널 방어). """ if s is None: return "" text = s if isinstance(s, str) else str(s) text = unicodedata.normalize("NFKC", text) text = fold_punct(text) text = "".join(ch for ch in text if unicodedata.category(ch) != "Cc") text = _WS_RE.sub(" ", text) return text.strip() def matching_key(s: str) -> str: """매칭 전용 키. 공백·구두점을 전부 지우고 소문자화한다. 원본 오탈자 CORTICOSTER OIDO 와 CORTICOSTEROIDO 를 같은 키로 만든다. """ if not s: return "" text = unicodedata.normalize("NFKC", s) text = _NON_ALNUM_KO_RE.sub("", text) return text.lower() def normalize_permit_no(raw: str) -> str: """등록번호 정규화 — 표기 흔들림 제거. 공백은 전부 삭제한다.""" text = clean_text(raw).replace(" ", "") return text.upper() def _is_valid_iso_date(iso: str) -> bool: """YYYY-MM-DD 문자열이 실재하는 날짜인지 확인한다.""" try: y, mo, d = (int(x) for x in iso.split("-")) except (ValueError, AttributeError): return False if not (1900 <= y <= 2199 and 1 <= mo <= 12): return False if mo in (1, 3, 5, 7, 8, 10, 12): last = 31 elif mo in (4, 6, 9, 11): last = 30 else: leap = (y % 4 == 0 and y % 100 != 0) or (y % 400 == 0) last = 29 if leap else 28 return 1 <= d <= last def parse_permit_no(raw: str) -> PermitParts: """등록번호를 구조화한다. 어떤 입력에서도 예외를 던지지 않는다.""" normalized = normalize_permit_no(raw) if not normalized: return PermitParts(raw=raw or "", normalized="", fmt="unknown") m = RE_PERMIT_STANDARD.match(normalized) if m: g = m.groupdict() d = g["accept_date"] accept_iso = f"{d[0:4]}-{d[4:6]}-{d[6:8]}" if not _is_valid_iso_date(accept_iso): # 20260231 처럼 형식만 맞는 날짜. 포맷은 standard 로 두되 날짜는 버린다. accept_iso = None table, rng, eff = GROUP_TABLE.get(g["group"], (None, None, None)) grant = int(g["grant"]) if g["grant"] else None ingr_no = int(g["ingr_no"]) return PermitParts( raw=raw or "", normalized=normalized, fmt="standard", accept_date=accept_iso, ingr_no=ingr_no, group=g["group"], group_table=table, group_range=rng, group_effective_from=eff, serial=int(g["serial"]), sub=int(g["sub"]) if g["sub"] else None, grant=grant, is_grant_derived=grant is not None, base_permit_no=normalized.split("(")[0], ingr_group_key=f"{ingr_no}-{g['group']}", ) m = RE_PERMIT_NEW_SUBSTANCE.match(normalized) if m: seq = m.group("seq") return PermitParts( raw=raw or "", normalized=normalized, fmt="new_substance", grant=int(seq) if seq else None, is_grant_derived=seq is not None, base_permit_no=normalized.split("(")[0], ) return PermitParts( raw=raw or "", normalized=normalized, fmt="unknown", base_permit_no=normalized.split("(")[0] or None, ) ``` ### 2.5 파서 회귀 케이스 (`tests/test_normalize.py` 의 정본 표) | 입력 | `fmt` | `accept_date` | `ingr_no` | `group` | `serial` | `sub` | `grant` | `ingr_group_key` | |---|---|---|---|---|---|---|---|---| | `20110531-71-B-317-05` | standard | 2011-05-31 | 71 | B | 317 | 5 | None | `71-B` | | `20110531-71-B-317-05(1)` | standard | 2011-05-31 | 71 | B | 317 | 5 | 1 | `71-B` | | `20260901-209-J-2270` | standard | 2026-09-01 | 209 | J | 2270 | **None** | None | `209-J` | | `20121228-168-I-169-04` | standard | 2012-12-28 | 168 | I | 169 | 4 | None | `168-I` | | 전각 하이픈·전각 공백이 섞인 위 값 | standard | 2012-12-28 | 168 | I | 169 | 4 | None | `168-I` | | `수6580-16-ND(20)` | new_substance | None | None | None | None | None | 20 | None | | `20260231-1-A-1-1` (2월 31일) | standard | **None** | 1 | A | 1 | 1 | None | `1-A` | | `20260901-209-Z-1` (미지 군) | standard | 2026-09-01 | 209 | Z | 1 | None | None | `209-Z` | | `이상한값` | unknown | None | None | None | None | None | None | None | | 빈 문자열 | unknown | None | None | None | None | None | None | None | > `Z` 군처럼 `GROUP_TABLE` 에 없는 알파벳은 **파싱은 성공시키되 `group_table`·`group_effective_from` 을 `None`** 으로 둔다. 제도 확대로 새 군이 생겨도 파이프라인이 멈추지 않는다. 대신 §7 품질 지표 `unknown_group_count` 가 잡아 리포트 메타 시트에 표시한다. --- ## 3. SQLite 스키마 전문 ### 3.1 설계 원칙 7가지 1. **append-only 를 기본으로 한다.** `snapshots`·`events`·`event_changes`·`agy_calls`·`alerts`·`quality_checks` 는 UPDATE·DELETE 하지 않는다. 유일한 예외는 §8 보존 정책에 따른 오래된 파티션 삭제다. 2. **현재 상태는 `records` 하나뿐이다.** UPDATE 가 허용되는 유일한 테이블이며, 그마저도 라이프사이클 컬럼(`last_seen_*`, `status`, `current_version_id`, 카운터)만 바뀐다. 3. **내용과 멤버십을 분리한다.** 같은 내용을 매일 다시 쓰지 않는다(§3.2 산정). 4. **자연 키를 조인 축으로 쓰지 않는다.** `dmf_key` 는 사람과 리포트가 쓰고, 테이블 간 조인은 정수 대리키(`record_id`·`version_id`·`run_seq`)로 한다. 등록번호가 21바이트 TEXT 라 1만 행 × 365일 조인 축으로 쓰면 그 자체가 수백 MB다. 5. **모든 시각은 ISO 8601 문자열, 모든 날짜는 `YYYY-MM-DD`.** SQLite 에 네이티브 날짜 타입이 없으므로 `CHECK ... GLOB` 로 형식을 강제한다. 6. **외래키를 실제로 켠다.** `PRAGMA foreign_keys=ON`(아키텍처 §3.8). 참조 무결성을 앱 코드가 아니라 DB 가 지킨다. 7. **비즈니스 규칙을 DB 제약으로 내린다.** idempotency 가드(하루 1건 SUCCESS)를 앱 코드의 `if` 가 아니라 **부분 유니크 인덱스**로 강제한다. 코드가 버그를 내도 데이터는 깨지지 않는다. ### 3.2 저장량 산정 — 왜 `record_versions` 를 분리하는가 관측 기준 전체 등록 건수 ≈ **9,840행**(리서치 §12.4). 안전하게 10,000행으로 잡는다. | 설계 | 하루 증가 | 1년 누적 | 5년 누적 | |---|---|---|---| | **A. 매일 전량 통 적재** (모든 필드 × 10,000행 × ~850B) | 8.5 MB | **3.1 GB** | 15.5 GB | | **B. 버전 분리 + TEXT 멤버십** (내용 델타 + `(run_id TEXT, dmf_key TEXT, hash TEXT)` ~90B) | 0.9 MB | 329 MB | 1.6 GB | | **C. 버전 분리 + 정수 멤버십** ← **채택** (`(run_seq, record_id, version_id)` 정수 3개 ~35B) | 0.35 MB | **128 MB** | 640 MB | - **내용 델타**: 기준선 수립일에 `record_versions` 10,000행(≈8.5MB)이 한 번 들어가고, 이후에는 **실제로 바뀐 레코드만** 새 버전을 만든다. 하루 신규 20 + 변경 30 = 50행 가정 시 42KB/일 → **연 15MB**. - 설계 C 는 A 대비 **24배** 작다. 아키텍처가 `storage.snapshot_retain_days = 0`(영구 보존)을 기본값으로 정한 것은 C 를 전제해야만 성립한다. - 그래도 5년에 640MB 다. `doctor` 진단에 **DB 크기 임계(기본 2GB)** 체크를 넣고, 초과하면 `snapshot_retain_days` 를 400(전년 동기 비교 가능 최소값)으로 낮추라고 안내한다. `records` + `events` 만으로 전체 이력 재구성이 가능하므로 멤버십 삭제는 **정보 손실이 아니라 인덱스 손실**이다. ### 3.3 `0001_init.sql` — 코어 스키마 전문 ```sql -- storage/migrations/0001_init.sql -- DMF Crawler 코어 스키마. 실행 · 수집 · 스냅샷 · 이벤트. -- 규칙: 이 파일은 한 번 배포되면 절대 수정하지 않는다. 변경은 새 번호 파일로만. -- --------------------------------------------------------------------------- -- 0. 마이그레이션 원장 -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS schema_version ( version INTEGER NOT NULL PRIMARY KEY, name TEXT NOT NULL, applied_at TEXT NOT NULL, checksum TEXT NOT NULL, -- 적용된 SQL 파일의 sha256. 사후 변조 탐지 app_version TEXT NOT NULL ); -- --------------------------------------------------------------------------- -- 1. 소스 메타 — 출처 표시(요구 N1)와 리포트 메타 시트의 원천 -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS sources ( source_id TEXT NOT NULL PRIMARY KEY, -- 'mfds_open_api' display_name TEXT NOT NULL, -- '식품의약품안전처_원료의약품등록(DMF)현황' provider TEXT NOT NULL, -- '식품의약품안전처' portal_url TEXT NOT NULL, -- data.go.kr 상세 페이지 endpoint TEXT NOT NULL, -- apis.data.go.kr/.../getMdcDmfList01 license TEXT NOT NULL, -- '이용허락범위 제한 없음' attribution TEXT NOT NULL, -- 리포트에 그대로 인쇄할 출처 문구 daily_quota INTEGER, -- 개발계정 10000 first_used_date TEXT, last_ok_date TEXT, last_ok_run_seq INTEGER, notes TEXT NOT NULL DEFAULT '', CHECK (first_used_date IS NULL OR first_used_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]'), CHECK (last_ok_date IS NULL OR last_ok_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]') ); INSERT OR IGNORE INTO sources (source_id, display_name, provider, portal_url, endpoint, license, attribution, daily_quota) VALUES ( 'mfds_open_api', '식품의약품안전처_원료의약품등록(DMF)현황', '식품의약품안전처', 'https://www.data.go.kr/data/15057075/openapi.do', 'https://apis.data.go.kr/1471000/MdcDmfInfoService01/getMdcDmfList01', '이용허락범위 제한 없음', '출처: 식품의약품안전처 「원료의약품등록(DMF)현황」 공공데이터포털 오픈API', 10000 ); -- --------------------------------------------------------------------------- -- 2. 실행 — 파이프라인 1회 = 1행 -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS runs ( run_seq INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, -- 조인 축 run_id TEXT NOT NULL UNIQUE, -- 'run_20260902_060013' run_date TEXT NOT NULL, -- KST 기준 실행일자 trigger TEXT NOT NULL, started_at TEXT NOT NULL, -- ISO8601 +09:00 finished_at TEXT, duration_ms INTEGER, status TEXT NOT NULL, exit_code INTEGER, is_baseline INTEGER NOT NULL DEFAULT 0, -- 기준선 수립 실행인가 integrity_ok INTEGER, -- NULL = 아직 판정 전 diff_performed INTEGER NOT NULL DEFAULT 0, report_path TEXT, app_version TEXT NOT NULL, notes TEXT NOT NULL DEFAULT '', CHECK (status IN ('RUNNING','SUCCESS','PARTIAL','FAILED','SKIPPED','BLOCKED')), CHECK (trigger IN ('scheduled','startup','manual','retry','backfill')), CHECK (run_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]'), CHECK (is_baseline IN (0,1)), CHECK (diff_performed IN (0,1)), CHECK (integrity_ok IS NULL OR integrity_ok IN (0,1)) ); -- ★ idempotency 가드를 DB 제약으로 내린다 (아키텍처 §4.1). -- 하루에 SUCCESS 실행은 최대 1건. 앱 코드가 last_success_run_on() 검사를 빼먹어도 막힌다. CREATE UNIQUE INDEX IF NOT EXISTS ux_runs_success_per_day ON runs(run_date) WHERE status = 'SUCCESS'; CREATE INDEX IF NOT EXISTS ix_runs_date ON runs(run_date DESC); CREATE INDEX IF NOT EXISTS ix_runs_status ON runs(status, run_date DESC); -- --------------------------------------------------------------------------- -- 3. 스테이지 체크포인트 — 재시도 시 완료 스테이지 건너뛰기의 근거 -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS stage_status ( run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, stage TEXT NOT NULL, attempt INTEGER NOT NULL DEFAULT 1, status TEXT NOT NULL, started_at TEXT NOT NULL, finished_at TEXT, duration_ms INTEGER, artifact_path TEXT, error TEXT, PRIMARY KEY (run_seq, stage), CHECK (status IN ('RUNNING','SUCCESS','FAILED','SKIPPED')), CHECK (stage IN ('preflight','fetch','normalize','integrity','diff', 'persist','enrich','report','backup','finalize')) ); -- --------------------------------------------------------------------------- -- 4. 수집 통계 — 실행 1회당 1행. 급감 판정(게이트 3)의 전일 기준값 공급원 -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS fetch_stats ( run_seq INTEGER NOT NULL PRIMARY KEY REFERENCES runs(run_seq) ON DELETE CASCADE, source_id TEXT NOT NULL REFERENCES sources(source_id), endpoint TEXT NOT NULL, page_size INTEGER NOT NULL, total_count_reported INTEGER, -- 응답의 totalCount records_received INTEGER NOT NULL, -- 실제로 받은 item 수 records_normalized INTEGER NOT NULL, -- DmfRecord 로 살아남은 수 pages_expected INTEGER, pages_fetched INTEGER NOT NULL, http_calls INTEGER NOT NULL, retry_count INTEGER NOT NULL DEFAULT 0, http_429_count INTEGER NOT NULL DEFAULT 0, bytes_received INTEGER NOT NULL DEFAULT 0, elapsed_seconds REAL NOT NULL DEFAULT 0, result_code TEXT, result_msg TEXT, body_signature_ok INTEGER NOT NULL DEFAULT 1, archive_dir TEXT, payload_sha256 TEXT NOT NULL, -- 정규화 레코드 집합 전체 해시 CHECK (body_signature_ok IN (0,1)) ); -- payload_sha256 이 전일과 같으면 '소스 미갱신'. 요구 SSOT 7.4(갱신 주기 실측)의 계측점. CREATE INDEX IF NOT EXISTS ix_fetch_payload ON fetch_stats(payload_sha256); -- --------------------------------------------------------------------------- -- 5. 품질 검증 결과 — 차단 게이트 5종과 비차단 지표 13종을 한 테이블에 -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS quality_checks ( run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, check_name TEXT NOT NULL, blocking INTEGER NOT NULL, -- 1 = 실패 시 diff 차단 passed INTEGER NOT NULL, severity TEXT NOT NULL, -- OK | INFO | WARN | CRITICAL observed REAL, -- 관측값(비율·건수) threshold REAL, -- 임계값 detail TEXT NOT NULL DEFAULT '', PRIMARY KEY (run_seq, check_name), CHECK (blocking IN (0,1)), CHECK (passed IN (0,1)), CHECK (severity IN ('OK','INFO','WARN','CRITICAL')) ); CREATE INDEX IF NOT EXISTS ix_quality_failed ON quality_checks(check_name, passed); -- --------------------------------------------------------------------------- -- 6. 레코드 버전 — 내용의 정본. 같은 내용은 두 번 저장하지 않는다. -- (dmf_key, content_hash) 가 논리 키. version_id 는 조인용 대리키. -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS record_versions ( version_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, dmf_key TEXT NOT NULL, content_hash TEXT NOT NULL, -- COMPARE_FIELDS 표시값 sha256[:16] identity_hash TEXT NOT NULL, -- COMPARE_KEY_FIELDS 정규화키 sha256[:16] -- 원천 7필드의 표준화 결과 permit_no TEXT NOT NULL, permit_no_raw TEXT NOT NULL, ingredient_name TEXT NOT NULL DEFAULT '', applicant TEXT NOT NULL DEFAULT '', manufacturer TEXT NOT NULL DEFAULT '', manufacture_place TEXT NOT NULL DEFAULT '', countries_json TEXT NOT NULL DEFAULT '[]', -- 정렬된 JSON 배열 permit_date TEXT NOT NULL DEFAULT '', -- 매칭 키 ingredient_key TEXT NOT NULL DEFAULT '', ingredient_base TEXT NOT NULL DEFAULT '', applicant_key TEXT NOT NULL DEFAULT '', manufacturer_key TEXT NOT NULL DEFAULT '', manufacture_place_key TEXT NOT NULL DEFAULT '', is_micronized INTEGER NOT NULL DEFAULT 0, sites_json TEXT NOT NULL DEFAULT '[]', country_count INTEGER NOT NULL DEFAULT 0, site_count INTEGER NOT NULL DEFAULT 0, -- 등록번호 파싱 결과 permit_fmt TEXT NOT NULL DEFAULT 'unknown', accept_date TEXT, ingr_no INTEGER, permit_group TEXT, group_table TEXT, group_effective_from TEXT, permit_serial INTEGER, permit_sub INTEGER, grant_seq INTEGER, is_grant_derived INTEGER NOT NULL DEFAULT 0, base_permit_no TEXT, ingr_group_key TEXT, -- 원본 보존과 계보 raw_json TEXT NOT NULL, first_run_seq INTEGER NOT NULL REFERENCES runs(run_seq), first_seen_date TEXT NOT NULL, UNIQUE (dmf_key, content_hash), CHECK (permit_fmt IN ('standard','new_substance','unknown')), CHECK (is_micronized IN (0,1)), CHECK (is_grant_derived IN (0,1)), CHECK (permit_date = '' OR permit_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]'), CHECK (accept_date IS NULL OR accept_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]') ); CREATE INDEX IF NOT EXISTS ix_ver_key ON record_versions(dmf_key, version_id DESC); CREATE INDEX IF NOT EXISTS ix_ver_ingr ON record_versions(ingredient_key); CREATE INDEX IF NOT EXISTS ix_ver_ingr_base ON record_versions(ingredient_base); CREATE INDEX IF NOT EXISTS ix_ver_applicant ON record_versions(applicant_key); CREATE INDEX IF NOT EXISTS ix_ver_mnf ON record_versions(manufacturer_key); CREATE INDEX IF NOT EXISTS ix_ver_cohort ON record_versions(ingr_group_key); -- --------------------------------------------------------------------------- -- 7. 레코드 현재 상태 — UPDATE 가 허용되는 유일한 테이블. 원장 시트의 원천. -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS records ( record_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, dmf_key TEXT NOT NULL UNIQUE, permit_no TEXT NOT NULL, base_permit_no TEXT, is_synthetic_key INTEGER NOT NULL DEFAULT 0, -- 'SYN-' 합성키 여부 dup_index INTEGER NOT NULL DEFAULT 0, current_version_id INTEGER REFERENCES record_versions(version_id), status TEXT NOT NULL DEFAULT 'ACTIVE', first_seen_date TEXT NOT NULL, first_seen_run_seq INTEGER NOT NULL REFERENCES runs(run_seq), last_seen_date TEXT NOT NULL, last_seen_run_seq INTEGER NOT NULL REFERENCES runs(run_seq), withdrawn_date TEXT, version_count INTEGER NOT NULL DEFAULT 1, change_count INTEGER NOT NULL DEFAULT 0, reappear_count INTEGER NOT NULL DEFAULT 0, CHECK (status IN ('ACTIVE','WITHDRAWN')), CHECK (is_synthetic_key IN (0,1)), CHECK (first_seen_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]'), CHECK (last_seen_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]'), CHECK (withdrawn_date IS NULL OR withdrawn_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]') ); CREATE INDEX IF NOT EXISTS ix_rec_status ON records(status, last_seen_date DESC); CREATE INDEX IF NOT EXISTS ix_rec_permit ON records(permit_no); CREATE INDEX IF NOT EXISTS ix_rec_base ON records(base_permit_no); CREATE INDEX IF NOT EXISTS ix_rec_first ON records(first_seen_date DESC); -- --------------------------------------------------------------------------- -- 8. 스냅샷 멤버십 — "이 실행에서 이 레코드가 이 버전으로 존재했다" 3열. -- 전량 재적재가 아니라 포인터만 쌓는다(§3.2 설계 C). -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS snapshots ( run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, record_id INTEGER NOT NULL REFERENCES records(record_id), version_id INTEGER NOT NULL REFERENCES record_versions(version_id), PRIMARY KEY (run_seq, record_id) ) WITHOUT ROWID; CREATE INDEX IF NOT EXISTS ix_snap_record ON snapshots(record_id, run_seq DESC); -- --------------------------------------------------------------------------- -- 9. 도메인 이벤트 — append only. 절대 UPDATE/DELETE 금지. -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS events ( event_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, event_date TEXT NOT NULL, occurred_at TEXT NOT NULL, record_id INTEGER NOT NULL REFERENCES records(record_id), dmf_key TEXT NOT NULL, -- 비정규화(리포트 조인 절약) permit_no TEXT NOT NULL, event_type TEXT NOT NULL, severity TEXT NOT NULL, is_reappearance INTEGER NOT NULL DEFAULT 0, change_count INTEGER NOT NULL DEFAULT 0, changed_fields TEXT NOT NULL DEFAULT '', -- 정렬된 CSV. 'countries,manufacturer' before_version_id INTEGER REFERENCES record_versions(version_id), after_version_id INTEGER REFERENCES record_versions(version_id), changes_json TEXT NOT NULL DEFAULT '[]', before_json TEXT, after_json TEXT, UNIQUE (run_seq, record_id, event_type), CHECK (event_type IN ('NEW','CHANGED','WITHDRAWN')), CHECK (severity IN ('CRITICAL','HIGH','MEDIUM','INFO','COSMETIC')), CHECK (is_reappearance IN (0,1)), CHECK (event_date GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9]'), -- NEW 는 before 가 없고, WITHDRAWN 은 after 가 없다. CHECK (event_type <> 'NEW' OR before_version_id IS NULL), CHECK (event_type <> 'WITHDRAWN' OR after_version_id IS NULL), CHECK (event_type <> 'CHANGED' OR (before_version_id IS NOT NULL AND after_version_id IS NOT NULL AND change_count > 0)) ); CREATE INDEX IF NOT EXISTS ix_ev_run ON events(run_seq, event_type); CREATE INDEX IF NOT EXISTS ix_ev_date ON events(event_date DESC, event_type); CREATE INDEX IF NOT EXISTS ix_ev_record ON events(record_id, event_date DESC); CREATE INDEX IF NOT EXISTS ix_ev_severity ON events(severity, event_date DESC); -- --------------------------------------------------------------------------- -- 10. 필드 단위 변경 — '오늘 변경분' 시트가 조인 없이 그대로 읽는다. -- --------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS event_changes ( change_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, event_id INTEGER NOT NULL REFERENCES events(event_id) ON DELETE CASCADE, run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, event_date TEXT NOT NULL, dmf_key TEXT NOT NULL, field TEXT NOT NULL, field_label TEXT NOT NULL, -- '제조소명' 등 한글 표시명 change_class TEXT NOT NULL, before_value TEXT NOT NULL DEFAULT '', after_value TEXT NOT NULL DEFAULT '', UNIQUE (event_id, field), CHECK (field IN ('ingredient_name','applicant','manufacturer', 'manufacture_place','countries','permit_date')), CHECK (change_class IN ('CRITICAL','HIGH','MEDIUM','INFO','COSMETIC')) ); CREATE INDEX IF NOT EXISTS ix_chg_field ON event_changes(field, event_date DESC); CREATE INDEX IF NOT EXISTS ix_chg_run ON event_changes(run_seq, change_class); -- --------------------------------------------------------------------------- -- 11. 읽기 편의 뷰 — 리포트 SQL 이 조인을 반복하지 않게 한다. -- --------------------------------------------------------------------------- CREATE VIEW IF NOT EXISTS v_current_records AS SELECT r.record_id, r.dmf_key, r.permit_no, r.base_permit_no, r.status, r.first_seen_date, r.last_seen_date, r.withdrawn_date, r.version_count, r.change_count, r.reappear_count, r.is_synthetic_key, v.ingredient_name, v.ingredient_base, v.is_micronized, v.applicant, v.manufacturer, v.manufacture_place, v.countries_json, v.country_count, v.site_count, v.sites_json, v.permit_date, v.accept_date, v.permit_fmt, v.ingr_no, v.permit_group, v.group_table, v.group_effective_from, v.permit_serial, v.permit_sub, v.grant_seq, v.is_grant_derived, v.ingr_group_key, v.content_hash, v.identity_hash FROM records r JOIN record_versions v ON v.version_id = r.current_version_id; CREATE VIEW IF NOT EXISTS v_run_summary AS SELECT ru.run_seq, ru.run_id, ru.run_date, ru.status, ru.duration_ms, ru.integrity_ok, ru.diff_performed, ru.is_baseline, ru.report_path, fs.total_count_reported, fs.records_normalized, fs.pages_fetched, fs.http_calls, fs.elapsed_seconds, fs.payload_sha256, (SELECT COUNT(*) FROM events e WHERE e.run_seq = ru.run_seq AND e.event_type = 'NEW') AS n_new, (SELECT COUNT(*) FROM events e WHERE e.run_seq = ru.run_seq AND e.event_type = 'CHANGED') AS n_changed, (SELECT COUNT(*) FROM events e WHERE e.run_seq = ru.run_seq AND e.event_type = 'WITHDRAWN') AS n_withdrawn, (SELECT COUNT(*) FROM events e WHERE e.run_seq = ru.run_seq AND e.severity = 'CRITICAL') AS n_critical FROM runs ru LEFT JOIN fetch_stats fs ON fs.run_seq = ru.run_seq; ``` ### 3.4 `0002_enrichment.sql` — AI 계층 (스키마 레벨 격리) `agy` 산출물은 **정본 데이터와 물리적으로 다른 테이블**에 둔다. AI 가 죽어도, 쿼터가 끊겨도, 잘못된 값을 내놔도 §3.3 의 어떤 테이블도 오염되지 않는다(아키텍처 ADR). ```sql -- storage/migrations/0002_enrichment.sql CREATE TABLE IF NOT EXISTS agy_calls ( call_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, purpose TEXT NOT NULL, -- 'daily_briefing' 등 started_at TEXT NOT NULL, finished_at TEXT, duration_ms INTEGER, model TEXT NOT NULL, effort TEXT NOT NULL, exit_code INTEGER, envelope_status TEXT, -- agy JSON 봉투의 status input_tokens INTEGER NOT NULL DEFAULT 0, output_tokens INTEGER NOT NULL DEFAULT 0, total_tokens INTEGER NOT NULL DEFAULT 0, json_extract_ok INTEGER NOT NULL DEFAULT 0, schema_valid INTEGER NOT NULL DEFAULT 0, retry_index INTEGER NOT NULL DEFAULT 0, error TEXT, stdout_path TEXT, -- logs/run_*/agy.stdout.json CHECK (json_extract_ok IN (0,1)), CHECK (schema_valid IN (0,1)) ); CREATE INDEX IF NOT EXISTS ix_agy_day ON agy_calls(started_at); CREATE TABLE IF NOT EXISTS enrichment_run ( run_seq INTEGER NOT NULL PRIMARY KEY REFERENCES runs(run_seq) ON DELETE CASCADE, call_id INTEGER REFERENCES agy_calls(call_id), status TEXT NOT NULL, -- OK | SKIPPED | FAILED skip_reason TEXT, -- 'diff_empty' | 'circuit_open' | 'token_cap' | 'disabled' headline TEXT, -- 대시보드 한 줄 요약 summary_md TEXT, -- 브리핑 본문(마크다운) risk_note TEXT, generated_at TEXT, CHECK (status IN ('OK','SKIPPED','FAILED')) ); CREATE TABLE IF NOT EXISTS enrichment ( enrichment_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, event_id INTEGER REFERENCES events(event_id) ON DELETE CASCADE, dmf_key TEXT, kind TEXT NOT NULL, -- 'event_comment' | 'ingredient_alias' | 'anomaly_hypothesis' payload_json TEXT NOT NULL, confidence REAL, applied INTEGER NOT NULL DEFAULT 0, -- 사람 승인 전에는 항상 0. 자동 적용 금지(ADR) CHECK (applied IN (0,1)), CHECK (confidence IS NULL OR (confidence >= 0.0 AND confidence <= 1.0)) ); CREATE INDEX IF NOT EXISTS ix_enrich_event ON enrichment(event_id); CREATE INDEX IF NOT EXISTS ix_enrich_run ON enrichment(run_seq, kind); ``` ### 3.5 `0003_ops.sql` — 운영 계층 ```sql -- storage/migrations/0003_ops.sql CREATE TABLE IF NOT EXISTS component_health ( component TEXT NOT NULL PRIMARY KEY, -- 'source_mfds' | 'agy' state TEXT NOT NULL DEFAULT 'CLOSED', consecutive_failures INTEGER NOT NULL DEFAULT 0, last_success_at TEXT, last_success_run TEXT, last_failure_at TEXT, last_reason TEXT, cooldown_until TEXT, transitions INTEGER NOT NULL DEFAULT 0, CHECK (state IN ('CLOSED','OPEN','HALF_OPEN')) ); INSERT OR IGNORE INTO component_health (component) VALUES ('source_mfds'), ('agy'); CREATE TABLE IF NOT EXISTS alerts ( alert_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, created_at TEXT NOT NULL, run_seq INTEGER REFERENCES runs(run_seq) ON DELETE SET NULL, level TEXT NOT NULL, code TEXT NOT NULL, -- 'INTEGRITY_BLOCKED' 등 dedup_key TEXT NOT NULL, -- 쿨다운 억제 축(요구 R7.8) title TEXT NOT NULL, what TEXT NOT NULL, -- 알림 문구 4요소 why TEXT NOT NULL, how TEXT NOT NULL, next_action TEXT NOT NULL, shown_at TEXT, -- 대화형 에이전트가 실제로 띄운 시각 resolved_at TEXT, CHECK (level IN ('INFO','WARN','CRITICAL')) ); CREATE INDEX IF NOT EXISTS ix_alert_open ON alerts(resolved_at, created_at DESC); CREATE INDEX IF NOT EXISTS ix_alert_dedup ON alerts(dedup_key, created_at DESC); CREATE TABLE IF NOT EXISTS watchlist ( watch_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, axis TEXT NOT NULL, -- ingredient | applicant | manufacturer | country pattern TEXT NOT NULL, -- 사람이 입력한 원문 pattern_key TEXT NOT NULL, -- matching_key() 적용 결과 match_mode TEXT NOT NULL DEFAULT 'contains', label TEXT NOT NULL DEFAULT '', enabled INTEGER NOT NULL DEFAULT 1, created_at TEXT NOT NULL, UNIQUE (axis, pattern_key, match_mode), CHECK (axis IN ('ingredient','applicant','manufacturer','country')), CHECK (match_mode IN ('contains','exact')), CHECK (enabled IN (0,1)) ); CREATE TABLE IF NOT EXISTS watchlist_hits ( hit_id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, run_seq INTEGER NOT NULL REFERENCES runs(run_seq) ON DELETE CASCADE, watch_id INTEGER NOT NULL REFERENCES watchlist(watch_id) ON DELETE CASCADE, event_id INTEGER REFERENCES events(event_id) ON DELETE CASCADE, record_id INTEGER NOT NULL REFERENCES records(record_id), matched_on TEXT NOT NULL, -- 실제로 매칭된 값 UNIQUE (run_seq, watch_id, record_id) ); CREATE INDEX IF NOT EXISTS ix_hit_run ON watchlist_hits(run_seq); -- 아키텍처 §2 가 지정한 'schema_version 인덱스' CREATE INDEX IF NOT EXISTS ix_schema_applied ON schema_version(applied_at DESC); ``` ### 3.6 마이그레이션 전략 **규칙 6가지.** 1. 파일명은 `NNNN_snake_name.sql`. 번호는 4자리 0-패딩, 결번 없이 1씩 증가. 2. **배포된 파일은 절대 수정하지 않는다.** `schema_version.checksum` 이 sha256 을 보관하므로 수정하면 다음 실행이 `StorageError` 로 멈춘다. 3. 각 파일은 **단일 트랜잭션 안에서** 실행된다. 파일 안에 `BEGIN`/`COMMIT` 을 쓰지 않는다(러너가 감싼다). 4. `ALTER TABLE ... ADD COLUMN` 은 허용. 컬럼 삭제·타입 변경이 필요하면 **새 테이블 생성 → INSERT SELECT → DROP → RENAME** 4단계를 한 파일에 쓴다. 5. **적용 전에 항상 `VACUUM INTO` 백업을 만든다.** 이것이 유일한 롤백 경로다(아키텍처 ADR-04). 실패 안내 문구에 백업본의 **절대경로**를 반드시 포함한다. 6. 데이터 이관이 필요한 마이그레이션은 SQL 만으로 끝낸다. 파이썬 후처리를 요구하는 마이그레이션은 만들지 않는다(부분 적용 상태가 생긴다). ```python # src/dmf_crawler/storage/db.py """연결 생성과 마이그레이션. 이 모듈 밖에서 sqlite3.connect 를 직접 부르지 않는다.""" from __future__ import annotations import hashlib import re import sqlite3 from datetime import datetime from pathlib import Path from dmf_crawler import __version__ from dmf_crawler.errors import StorageError MIGRATION_RE = re.compile(r"^(?P\d{4})_(?P[a-z0-9_]+)\.sql$") def connect(db_path: Path, *, read_only: bool = False, busy_timeout_ms: int = 15000) -> sqlite3.Connection: db_path.parent.mkdir(parents=True, exist_ok=True) if read_only: uri = f"file:{db_path.as_posix()}?mode=ro" conn = sqlite3.connect(uri, uri=True, timeout=busy_timeout_ms / 1000) else: conn = sqlite3.connect(db_path, timeout=busy_timeout_ms / 1000) conn.row_factory = sqlite3.Row conn.execute("PRAGMA journal_mode=WAL") conn.execute("PRAGMA foreign_keys=ON") conn.execute("PRAGMA synchronous=NORMAL") conn.execute(f"PRAGMA busy_timeout={int(busy_timeout_ms)}") conn.execute("PRAGMA temp_store=MEMORY") return conn def current_version(conn: sqlite3.Connection) -> int: row = conn.execute( "SELECT name FROM sqlite_master WHERE type='table' AND name='schema_version'" ).fetchone() if row is None: return 0 got = conn.execute("SELECT COALESCE(MAX(version), 0) AS v FROM schema_version").fetchone() return int(got["v"]) def _discover(migrations_dir: Path) -> list[tuple[int, str, Path]]: found: list[tuple[int, str, Path]] = [] for path in sorted(migrations_dir.glob("*.sql")): m = MIGRATION_RE.match(path.name) if not m: raise StorageError(f"마이그레이션 파일명 규칙 위반: {path.name}") found.append((int(m.group("num")), m.group("name"), path)) numbers = [n for n, _, _ in found] if numbers != list(range(1, len(numbers) + 1)): raise StorageError(f"마이그레이션 번호에 결번 또는 중복이 있다: {numbers}") return found def _verify_checksums(conn: sqlite3.Connection, found: list[tuple[int, str, Path]]) -> None: """이미 적용된 마이그레이션 파일이 사후 수정되지 않았는지 확인한다.""" if current_version(conn) == 0: return applied = {int(r["version"]): r["checksum"] for r in conn.execute("SELECT version, checksum FROM schema_version")} for num, _, path in found: if num not in applied: continue digest = hashlib.sha256(path.read_bytes()).hexdigest() if digest != applied[num]: raise StorageError( f"이미 적용된 마이그레이션 {path.name} 이 수정됐다. " f"기록된 체크섬={applied[num][:12]}… 현재={digest[:12]}… " f"배포된 마이그레이션 파일은 수정하지 말고 새 번호 파일을 추가하라." ) def apply_migrations(conn: sqlite3.Connection, migrations_dir: Path, backup_dir: Path) -> list[int]: """적용 전 VACUUM INTO 백업 -> 번호순 적용 -> schema_version 기록.""" found = _discover(migrations_dir) _verify_checksums(conn, found) have = current_version(conn) todo = [item for item in found if item[0] > have] if not todo: return [] backup_path: Path | None = None if have > 0: # 최초 생성이 아니면 반드시 백업부터 backup_dir.mkdir(parents=True, exist_ok=True) stamp = datetime.now().strftime("%Y%m%d_%H%M%S") backup_path = backup_dir / f"premigrate_v{have}_{stamp}.sqlite3" conn.execute("VACUUM INTO ?", (str(backup_path),)) applied: list[int] = [] for num, name, path in todo: sql = path.read_text(encoding="utf-8") digest = hashlib.sha256(path.read_bytes()).hexdigest() try: conn.execute("BEGIN") conn.executescript(sql) conn.execute( "INSERT INTO schema_version (version, name, applied_at, checksum, app_version)" " VALUES (?, ?, ?, ?, ?)", (num, name, datetime.now().astimezone().isoformat(timespec="seconds"), digest, __version__), ) conn.execute("COMMIT") except Exception as exc: # noqa: BLE001 — 어떤 예외든 롤백 후 승격한다 conn.execute("ROLLBACK") hint = (f" 복원본: {backup_path}" if backup_path else " (최초 생성이라 백업본이 없다. data/dmf.sqlite3 를 지우고 다시 실행하라)") raise StorageError(f"마이그레이션 {path.name} 적용 실패: {exc}.{hint}") from exc applied.append(num) return applied def integrity_check(conn: sqlite3.Connection) -> tuple[bool, str]: rows = conn.execute("PRAGMA integrity_check").fetchall() messages = [r[0] for r in rows] ok = messages == ["ok"] fk = conn.execute("PRAGMA foreign_key_check").fetchall() if fk: ok = False messages.append(f"foreign_key_check 위반 {len(fk)}건") return ok, "; ".join(messages) ``` ### 3.7 자주 쓰는 조회 SQL ```sql -- (1) 오늘 변경분 — s01_changes 시트가 그대로 읽는다. SELECT e.event_type, e.severity, e.is_reappearance, e.dmf_key, e.permit_no, v.ingredient_name, v.applicant, v.manufacturer, v.countries_json, e.changed_fields, e.changes_json FROM events e JOIN runs r ON r.run_seq = e.run_seq LEFT JOIN record_versions v ON v.version_id = COALESCE(e.after_version_id, e.before_version_id) WHERE r.run_id = :run_id ORDER BY CASE e.severity WHEN 'CRITICAL' THEN 0 WHEN 'HIGH' THEN 1 WHEN 'MEDIUM' THEN 2 WHEN 'INFO' THEN 3 ELSE 4 END, e.event_type, e.dmf_key; -- (2) 특정 시점의 전체 상태 (temporal query) SELECT v.* FROM snapshots s JOIN runs r ON r.run_seq = s.run_seq JOIN record_versions v ON v.version_id = s.version_id WHERE r.run_seq = (SELECT MAX(run_seq) FROM runs WHERE status = 'SUCCESS' AND run_date <= :as_of_date); -- (3) 한 등록건의 전체 변경 이력 SELECT e.event_date, e.event_type, e.severity, e.changed_fields, e.changes_json FROM events e JOIN records r ON r.record_id = e.record_id WHERE r.dmf_key = :dmf_key ORDER BY e.event_date, e.event_id; -- (4) 성분별 소싱 대안 폭 — s03_ingredient 시트 SELECT ingredient_base, COUNT(*) AS n_records, COUNT(DISTINCT manufacturer_key) AS n_manufacturers, COUNT(DISTINCT applicant_key) AS n_applicants, SUM(is_micronized) AS n_micronized FROM v_current_records WHERE status = 'ACTIVE' GROUP BY ingredient_base ORDER BY n_records DESC LIMIT :top_n; -- (5) 제조국가 점유 — 다중 국가를 JSON1 로 분해해 집계한다. SELECT j.value AS country, COUNT(*) AS n_records FROM v_current_records c, json_each(c.countries_json) j WHERE c.status = 'ACTIVE' GROUP BY j.value ORDER BY n_records DESC; -- (6) 일자별 건수 추이 — s06_trend 시트 · 스파크라인 원본 SELECT r.run_date, f.total_count_reported, f.records_normalized, (SELECT COUNT(*) FROM events e WHERE e.run_seq = r.run_seq AND e.event_type='NEW') AS n_new, (SELECT COUNT(*) FROM events e WHERE e.run_seq = r.run_seq AND e.event_type='CHANGED') AS n_changed, (SELECT COUNT(*) FROM events e WHERE e.run_seq = r.run_seq AND e.event_type='WITHDRAWN') AS n_withdrawn FROM runs r JOIN fetch_stats f ON f.run_seq = r.run_seq WHERE r.status = 'SUCCESS' AND r.run_date >= date('now', '-' || :trend_days || ' days') ORDER BY r.run_date; -- (7) 소스가 실제로 며칠마다 갱신되는가 (요구 SSOT 7.4 실측) SELECT payload_sha256, COUNT(*) AS days_identical, MIN(r.run_date) AS first_date, MAX(r.run_date) AS last_date FROM fetch_stats f JOIN runs r ON r.run_seq = f.run_seq WHERE r.status = 'SUCCESS' GROUP BY payload_sha256 ORDER BY first_date; ``` --- ## 4. 레코드 동일성 판정 ### 4.1 무엇을 같은 레코드로 볼 것인가 — 결정 **같은 레코드 = 같은 `dmf_key`.** 그 이상도 이하도 아니다. 성분·업체가 같아도 `dmf_key` 가 다르면 다른 등록이고, 성분·업체·제조소가 전부 바뀌어도 `dmf_key` 가 같으면 같은 등록의 변경이다. 이 단순한 규칙이 성립하려면 `dmf_key` 생성이 **완전히 결정론적**이어야 한다. 같은 입력에서 언제나 같은 키가 나와야 하고, 원본의 표기가 흔들려도 키는 흔들리지 않아야 한다. §2.3 의 3층 구조와 아래 정규화 알고리즘이 그 보장이다. **동일성 판정에 참여하지 않는 것을 명시한다** — 실수를 막기 위해서다. | 참여하지 않는 것 | 이유 | |---|---| | `ingredient_base` (염·수화물 제거형) | `탐스로신염산염` 과 `탐스로신메실산염`은 **다른 등록**이다. 규정 제2조 2호가 염류·수화물을 각각 등록 대상으로 규정한다. 집계 축으로만 쓴다 | | `is_micronized` 를 뺀 성분명 | `미분화부데소니드` 는 `부데소니드` 와 **별도 등록**이다(리서치 §10.6-1) | | `base_permit_no` (허여서 괄호 제거형) | `…-05` 와 `…-05(1)` 은 **다른 등록건**이다. 묶어서 보여주기만 한다 | | `sites` 개별 항목 | 제조소는 1:N 이지만 등록건의 하위 속성이다. 제조소 단위로 레코드를 쪼개면 취하 판정이 붕괴한다 | ### 4.2 표기 흔들림 정규화 알고리즘 정규화는 **2단 파이프라인**이다. ``` 원본 문자열 │ ├─[1단] clean_text() → 표시값(display). 사람이 읽고 리포트에 인쇄된다. │ NFKC · 전각→반각 · 구두점 접기 · 제어문자 제거 · 연속공백 1개 · trim │ + 필드별 추가 규칙(법인격 표기 통일, 국가 별칭, 날짜 ISO화) │ └─[2단] matching_key() → 매칭키(key). 기계가 비교하고 워치리스트가 찾는다. 1단 결과에서 공백·구두점 전부 삭제 + 소문자화 + 필드별 추가 규칙(법인격 어휘 제거) ``` **두 단을 나누는 이유**: 1단만 있으면 `CORTICOSTER OIDO`(원본 오탈자) → `CORTICOSTEROIDO`(정정) 이 변경으로 잡힌다. 2단만 있으면 리포트에 `sicorsocietaitaliana…` 같은 읽을 수 없는 문자열이 인쇄된다. 둘 다 필요하다. `content_hash` 는 1단 결과로, `identity_hash` 는 2단 결과로 만든다(§1.3). **흔들림 유형별 처리표.** | 흔들림 유형 | 실례 | 처리 단계 | 처리 방법 | |---|---|---|---| | 전각/반각 | `SICOR` vs `SICOR` | 1단 | `unicodedata.normalize("NFKC")` | | 전각 하이픈·괄호 | `20121228-168` / `(1)` | 1단 | `_PUNCT_FOLD` 치환표 | | 비단절 공백·제로폭 | U+00A0, U+200B, U+FEFF | 1단 | `_PUNCT_FOLD` + `Cc` 카테고리 제거 | | 연속 공백·앞뒤 공백 | `(주) 대웅제약 ` | 1단 | `\s+` → 한 칸, `strip()` | | 법인격 표기 | `㈜하이플` / `(주)하이플` / `주식회사 하이플` | 1단(표시 통일) + 2단(어휘 제거) | `㈜`→`(주)` 통일 후, 키에서는 법인격 어휘 삭제 | | 원본 오탈자 공백 | `CORTICOSTER OIDO` | 2단 | 공백 전부 삭제 → 정정본과 동일 키 | | 대소문자 | `S.R.L.` vs `s.r.l.` | 2단 | `lower()` | | 구두점 | `Co., Ltd.` vs `Co.,Ltd` | 2단 | 영숫자·한글 외 전부 삭제 | | 국가 다중값 순서 | `이탈리아,스위스` vs `스위스,이탈리아` | 1단 | 분해 후 **정렬** | | 국가 별칭 | `미국` / `미합중국` / `USA` | 1단 | `COUNTRY_ALIASES` 치환표 | | 날짜 표기 | `2015-02-26` / `20150226` / `2015.02.26` | 1단 | `normalize_date()` → ISO | | 다중 제조소 구분자 | `A , B` vs `A, B` | 1단 | ` ,` 우선 분해 후 ` , ` 로 재조립 | ### 4.3 정규화 코드 전문 ```python # src/dmf_crawler/normalize.py (3/3 — 필드 정규화 · 레코드 조립) # -------------------------------------------------------------------------- # 법인격 어휘 — 매칭키에서 제거한다. 긴 것부터 지워야 부분 삭제가 안 생긴다. # -------------------------------------------------------------------------- LEGAL_FORMS_KO: tuple[str, ...] = ( "주식회사", "유한책임회사", "유한회사", "합자회사", "합명회사", "재단법인", "사단법인", "의료법인", "학교법인", "(주)", "(유)", "(재)", "(사)", ) LEGAL_FORMS_EN: tuple[str, ...] = ( "coltd", "companylimited", "limited", "ltd", "llc", "llp", "inc", "incorporated", "corporation", "corp", "gmbh", "ag", "sa", "srl", "spa", "bv", "nv", "as", "pvtltd", "pvt", "plc", "kg", "oy", "ab", "sas", "sl", "pte", "sdnbhd", ) # 표시값 단계의 법인격 표기 통일 _LEGAL_DISPLAY_FOLD = { "㈜": "(주)", "㈔": "(사)", "㈖": "(재)", "㈕": "(유)", } # -------------------------------------------------------------------------- # 성분명 접두/접미 — ingredient_base(집계 전용) 산출에만 쓴다. # 동일성 판정에는 절대 쓰지 않는다(염 형태가 다르면 다른 등록이다). # -------------------------------------------------------------------------- MICRONIZE_PREFIXES: tuple[str, ...] = ("초미분화", "미분화") SALT_SUFFIXES: tuple[str, ...] = ( "브롬화수소산염", "메탄술폰산염", "메탄설폰산염", "메실산염", "베실산염", "토실산염", "에실산염", "이세티온산염", "파모산염", "팜산염", "푸마르산염", "말레산염", "말산염", "타르타르산염", "주석산염", "시트르산염", "구연산염", "숙신산염", "아세트산염", "초산염", "락트산염", "젖산염", "글루콘산염", "글루쿠론산염", "아스파르트산염", "팔미트산염", "스테아르산염", "벤조산염", "살리실산염", "옥살산염", "인산염", "황산염", "질산염", "염산염", "브롬산염", "요오드화물", "브롬화물", "염화물", ) HYDRATE_SUFFIXES: tuple[str, ...] = ( "일수화물", "이수화물", "삼수화물", "사수화물", "오수화물", "육수화물", "칠수화물", "팔수화물", "반수화물", "수화물", "무수물", ) CATION_SUFFIXES: tuple[str, ...] = ( "나트륨", "칼륨", "칼슘", "마그네슘", "아연", "리튬", "암모늄", "메글루민", "트로메타민", "디에탄올아민", "에탄올아민", "베타덱스", ) _STRIPPABLE_SUFFIXES: tuple[str, ...] = tuple( sorted(SALT_SUFFIXES + HYDRATE_SUFFIXES + CATION_SUFFIXES, key=len, reverse=True) ) # -------------------------------------------------------------------------- # 국가명 별칭 — 표시값 단계에서 표준명으로 접는다. # -------------------------------------------------------------------------- COUNTRY_ALIASES: Mapping[str, str] = { "대한민국": "한국", "korea": "한국", "republicofkorea": "한국", "kr": "한국", "미합중국": "미국", "usa": "미국", "us": "미국", "unitedstates": "미국", "중화인민공화국": "중국", "china": "중국", "cn": "중국", "인도": "인도", "india": "인도", "in": "인도", "일본": "일본", "japan": "일본", "jp": "일본", "이탈리아": "이탈리아", "italy": "이탈리아", "it": "이탈리아", "스위스": "스위스", "switzerland": "스위스", "ch": "스위스", "독일": "독일", "germany": "독일", "de": "독일", "스페인": "스페인", "spain": "스페인", "es": "스페인", "프랑스": "프랑스", "france": "프랑스", "fr": "프랑스", "영국": "영국", "unitedkingdom": "영국", "uk": "영국", "gb": "영국", "아일랜드": "아일랜드", "ireland": "아일랜드", "ie": "아일랜드", "이스라엘": "이스라엘", "israel": "이스라엘", "il": "이스라엘", "대만": "대만", "taiwan": "대만", "tw": "대만", "슬로베니아": "슬로베니아", "헝가리": "헝가리", "폴란드": "폴란드", "오스트리아": "오스트리아", "네덜란드": "네덜란드", "벨기에": "벨기에", "덴마크": "덴마크", "스웨덴": "스웨덴", "핀란드": "핀란드", "노르웨이": "노르웨이", "포르투갈": "포르투갈", "그리스": "그리스", "루마니아": "루마니아", "체코": "체코", "슬로바키아": "슬로바키아", "크로아티아": "크로아티아", "캐나다": "캐나다", "canada": "캐나다", "멕시코": "멕시코", "브라질": "브라질", "아르헨티나": "아르헨티나", "호주": "호주", "australia": "호주", "뉴질랜드": "뉴질랜드", "싱가포르": "싱가포르", "말레이시아": "말레이시아", "인도네시아": "인도네시아", "베트남": "베트남", "태국": "태국", "터키": "튀르키예", "튀르키예": "튀르키예", "turkey": "튀르키예", } _COUNTRY_SPLIT_RE = re.compile(r"[,/·;|]+") _SITE_ROLE_RE = re.compile(r"^\[(?P[^\]]{1,40})\]\s*") _DATE_RE = re.compile(r"(?P(?:19|20)\d{2})\D?(?P\d{1,2})\D?(?P\d{1,2})") def normalize_org_display(s: object) -> str: """업체명·제조소명의 표시값. 법인격 기호만 통일한다.""" text = clean_text(s) for src, dst in _LEGAL_DISPLAY_FOLD.items(): text = text.replace(src, dst) return text def normalize_org_key(display: str) -> str: """업체명 매칭키. 법인격 어휘를 제거한 뒤 공백·구두점을 지운다.""" text = display for token in LEGAL_FORMS_KO: text = text.replace(token, " ") key = matching_key(text) changed = True while changed: changed = False for token in LEGAL_FORMS_EN: if len(key) > len(token) and key.endswith(token): key = key[: -len(token)] changed = True if len(key) > len(token) and key.startswith(token): key = key[len(token):] changed = True return key def normalize_ingredient(raw: object) -> tuple[str, str, str, bool]: """성분명 정규화. 반환: (표시값, 매칭키, 기본명(집계 전용), 미분화 여부) 기본명은 염/수화물/양이온 접미와 미분화 접두를 반복 제거한 결과다. 동일성 판정에는 절대 쓰지 않는다. """ display = clean_text(raw) key = matching_key(display) micronized = False base = key for prefix in MICRONIZE_PREFIXES: pkey = matching_key(prefix) if base.startswith(pkey): micronized = True base = base[len(pkey):] break changed = True while changed and base: changed = False for suffix in _STRIPPABLE_SUFFIXES: skey = matching_key(suffix) if len(base) > len(skey) and base.endswith(skey): base = base[: -len(skey)] changed = True break return display, key, (base or key), micronized def split_countries(raw: object) -> tuple[str, ...]: """제조국가 다중값을 분해·별칭 정규화·중복 제거·정렬한다. 정렬하는 이유: 원본이 '이탈리아,스위스' 와 '스위스,이탈리아' 를 오가면 순서만 바뀐 것이 CHANGED 로 오탐된다(요구 R2.3). """ text = clean_text(raw) if not text: return () out: set[str] = set() for token in _COUNTRY_SPLIT_RE.split(text): name = token.strip() if not name: continue canonical = COUNTRY_ALIASES.get(matching_key(name), name) out.add(canonical) return tuple(sorted(out)) def split_sites(name_raw: object, place_raw: object, countries: tuple[str, ...]) -> tuple[Site, ...]: """제조소 다중값 분해. 집계 전용이며 diff 에는 쓰지 않는다. 주소 안의 콤마('Rho(MI) - Via Terrazzano, 77, Italy')와 충돌하지 않도록 '공백+콤마' 를 1순위 구분자로 쓴다(리서치 §10.6-3). """ names = _split_multi(clean_text(name_raw)) places = _split_multi(clean_text(place_raw)) n = max(len(names), len(places), 1) sites: list[Site] = [] for i in range(n): raw_name = names[i] if i < len(names) else "" role = "" m = _SITE_ROLE_RE.match(raw_name) if m: role = m.group("role").strip() raw_name = raw_name[m.end():].strip() sites.append(Site( seq=i, name=raw_name, address=places[i] if i < len(places) else "", country=countries[i] if i < len(countries) else (countries[0] if countries else ""), role=role, )) return tuple(s for s in sites if s.name or s.address) def _split_multi(text: str) -> list[str]: """'공백+콤마' 우선 분해. 그 구분자가 없으면 분해하지 않는다.""" if not text: return [] if " ," in text: return [part.strip() for part in text.split(" ,") if part.strip()] return [text] def normalize_date(raw: object) -> str: """날짜를 ISO YYYY-MM-DD 로. 실패하면 빈 문자열(널 방어).""" text = clean_text(raw) if not text: return "" m = _DATE_RE.search(text) if not m: return "" iso = f"{m.group('y')}-{int(m.group('m')):02d}-{int(m.group('d')):02d}" return iso if _is_valid_iso_date(iso) else "" # -------------------------------------------------------------------------- # 해시와 키 # -------------------------------------------------------------------------- _UNIT_SEP = "\x1f" # 필드 구분자. 데이터에 절대 등장하지 않는 제어문자. def _hash16(parts: Iterable[str]) -> str: joined = _UNIT_SEP.join(parts) return hashlib.sha256(joined.encode("utf-8")).hexdigest()[:16] def compute_content_hash(rec: DmfRecord) -> str: """COMPARE_FIELDS 표시값 기준. 표기 변화까지 잡는다(버전 식별용).""" view = rec.compare_view() return _hash16(view[f] for f in COMPARE_FIELDS) def compute_identity_hash(rec: DmfRecord) -> str: """COMPARE_KEY_FIELDS 정규화키 기준. 실질 변화만 잡는다(등급 판정용).""" view = rec.key_view() return _hash16(view[f] for f in COMPARE_KEY_FIELDS) def make_dmf_key(permit_no: str, dup_index: int = 0, fallback_seed: Optional[str] = None) -> str: """레코드 키 생성. - 정상 : permit_no - 중복 : permit_no#2, permit_no#3 ... - 부재 : SYN- (합성키. 리포트에서 식별 가능해야 한다) """ if not permit_no: seed = fallback_seed or "" digest = hashlib.sha1(seed.encode("utf-8")).hexdigest()[:12] base = f"SYN-{digest}" else: base = permit_no return base if dup_index == 0 else f"{base}#{dup_index + 1}" ``` ### 4.4 레코드 조립과 중복 처리 ```python # src/dmf_crawler/normalize.py (계속 — 조립 진입점) def normalize_one(raw: Mapping[str, object], dup_index: int = 0) -> DmfRecord: """원본 1건을 DmfRecord 로. 예외를 던지지 않는다.""" raw_str = {col: ("" if raw.get(col) is None else str(raw.get(col))) for col in SOURCE_COLUMNS} permit_no_raw = raw_str[COL_PERMIT_NO] permit = parse_permit_no(permit_no_raw) ingredient_name, ingredient_key, ingredient_base, micronized = \ normalize_ingredient(raw_str[COL_INGREDIENT]) applicant = normalize_org_display(raw_str[COL_APPLICANT]) manufacturer = normalize_org_display(raw_str[COL_MANUFACTURER]) place = clean_text(raw_str[COL_PLACE]) countries = split_countries(raw_str[COL_COUNTRY]) sites = split_sites(manufacturer, place, countries) # 제조소명에 [미분화공정 제조소] 역할 표기가 있으면 그것도 미분화 신호다. if any(m in manufacturer for m in ("미분화공정", "미분화 공정")): micronized = True permit_no = permit.normalized fallback_seed = _UNIT_SEP.join(( ingredient_key, normalize_org_key(applicant), normalize_org_key(manufacturer), normalize_date(raw_str[COL_PERMIT_DATE]), )) rec = DmfRecord( dmf_key=make_dmf_key(permit_no, dup_index, fallback_seed), permit_no=permit_no, permit_no_raw=permit_no_raw, ingredient_name=ingredient_name, ingredient_key=ingredient_key, ingredient_base=ingredient_base, is_micronized=micronized, applicant=applicant, applicant_key=normalize_org_key(applicant), manufacturer=manufacturer, manufacturer_key=normalize_org_key(manufacturer), manufacture_place=place, manufacture_place_key=matching_key(place), countries=countries, sites=sites, permit_date=normalize_date(raw_str[COL_PERMIT_DATE]), permit=permit, raw=raw_str, dup_index=dup_index, ) # frozen dataclass 이므로 해시는 재생성으로 채운다. from dataclasses import replace return replace(rec, content_hash=compute_content_hash(rec), identity_hash=compute_identity_hash(rec)) def _dup_sort_key(rec: DmfRecord) -> tuple[str, str, str, str, str]: """중복 그룹 내 결정론적 정렬 키. 아키텍처는 '제조소명 사전순' 이라고 했으나 동률 시 순서가 흔들리므로 5요소 전체를 쓴다. 그래도 흔들릴 수 있어 diff 가 §4.5 로 이중 방어한다. """ return (rec.manufacturer_key, rec.manufacture_place_key, "|".join(rec.countries), rec.ingredient_key, rec.permit_date) def normalize_all(raws: list[Mapping[str, object]]) -> tuple[list[DmfRecord], NormalizeStats]: """원본 전량을 표준화한다. 중복 등록번호에 결정론적 접미사를 부여한다.""" staged: list[DmfRecord] = [] rejected: list[dict[str, str]] = [] replacement_chars = 0 oversize: dict[str, int] = {col: 0 for col in SOURCE_COLUMNS} for raw in raws: try: rec = normalize_one(raw, dup_index=0) except Exception as exc: # noqa: BLE001 — 개별 실패가 전체를 멈추지 않는다 rejected.append({"error": f"{type(exc).__name__}: {exc}", "raw": repr(raw)[:500]}) continue for col in SOURCE_COLUMNS: value = rec.raw.get(col, "") if "�" in value: replacement_chars += 1 if len(value) > SOURCE_MAX_LEN[col]: oversize[col] += 1 staged.append(rec) # 등록번호별로 묶어 중복에 접미사를 부여한다. groups: dict[str, list[DmfRecord]] = {} for rec in staged: groups.setdefault(rec.permit_no, []).append(rec) out: list[DmfRecord] = [] duplicate_rows = 0 duplicate_groups = 0 for permit_no, members in groups.items(): if len(members) == 1: out.append(members[0]) continue duplicate_groups += 1 duplicate_rows += len(members) - 1 for idx, rec in enumerate(sorted(members, key=_dup_sort_key)): out.append(normalize_one(rec.raw, dup_index=idx)) null_counts = { "permit_no": sum(1 for r in out if not r.permit_no), "ingredient_name": sum(1 for r in out if not r.ingredient_name), "applicant": sum(1 for r in out if not r.applicant), } stats = NormalizeStats( total_in=len(raws), total_out=len(out), rejected=tuple(rejected), null_counts=null_counts, duplicate_permit_no=duplicate_rows, duplicate_groups=duplicate_groups, unparsed_permit_no=sum(1 for r in out if r.permit.fmt == "unknown"), new_substance_count=sum(1 for r in out if r.permit.fmt == "new_substance"), synthetic_key_count=sum(1 for r in out if r.dmf_key.startswith("SYN-")), invalid_permit_date=sum(1 for r in out if not r.permit_date), oversize_fields=oversize, replacement_char_count=replacement_chars, ) out.sort(key=lambda r: r.dmf_key) return out, stats ``` ### 4.5 중복 등록번호 그룹의 대응 문제 접미사(`#2`, `#3`)는 **정렬 순서에 의존**한다. 정렬 키의 어느 한 요소가 바뀌면 어제의 `#2` 가 오늘의 `#1` 이 될 수 있고, 그러면 diff 가 **두 건의 CHANGED** 를 만든다. 실제 변경은 한 건인데 두 건이 뜨고, 그 내용도 뒤섞인다. **해법: 접미사를 짝짓지 말고 그룹을 짝짓는다.** ``` 같은 base(접미사 제거한 permit_no)를 가진 어제 멤버 P = {p1, p2} 오늘 멤버 C = {c1, c2} 1. content_hash 가 완전히 같은 쌍을 먼저 확정한다 (변경 없음). 2. 남은 쌍에 대해 COMPARE_FIELDS 6개 중 몇 개가 같은지로 점수를 매긴다. 가중치: manufacturer_key 3, manufacture_place_key 2, ingredient_key 3, applicant_key 2, countries 1, permit_date 1 (합 12) 3. 점수 내림차순 탐욕 매칭. 점수가 임계(6, 즉 절반) 미만이면 매칭하지 않는다. 4. 매칭된 쌍 → CHANGED(또는 동일). 남은 오늘 멤버 → NEW. 남은 어제 멤버 → WITHDRAWN. ``` 그룹 크기가 1인 절대다수 케이스에서는 이 알고리즘이 **단순 키 비교와 정확히 같은 결과**를 낸다. 비용은 O(n²)이지만 n은 사실상 2~3이다. ```python # src/dmf_crawler/diff.py (1/3 — 중복 그룹 매칭) MATCH_WEIGHTS: Mapping[str, int] = { "ingredient_key": 3, "manufacturer_key": 3, "applicant_key": 2, "manufacture_place_key": 2, "countries": 1, "permit_date": 1, } MATCH_TOTAL = sum(MATCH_WEIGHTS.values()) # 12 MATCH_THRESHOLD = MATCH_TOTAL // 2 # 6 def group_base(dmf_key: str) -> str: """'20121228-168-I-169-04#2' -> '20121228-168-I-169-04'""" return dmf_key.split("#", 1)[0] def match_score(a: DmfRecord, b: DmfRecord) -> int: va, vb = a.key_view(), b.key_view() return sum(w for f, w in MATCH_WEIGHTS.items() if va[f] == vb[f]) def pair_group(previous: list[DmfRecord], current: list[DmfRecord] ) -> tuple[list[tuple[DmfRecord, DmfRecord]], list[DmfRecord], list[DmfRecord]]: """(짝지어진 쌍, 짝 없는 오늘 = NEW 후보, 짝 없는 어제 = WITHDRAWN 후보)""" pairs: list[tuple[DmfRecord, DmfRecord]] = [] left = list(previous) right = list(current) # 1단계: content_hash 완전 일치부터 확정한다. by_hash: dict[str, list[DmfRecord]] = {} for rec in left: by_hash.setdefault(rec.content_hash, []).append(rec) remaining_right: list[DmfRecord] = [] for rec in right: bucket = by_hash.get(rec.content_hash) if bucket: pairs.append((bucket.pop(0), rec)) else: remaining_right.append(rec) remaining_left = [r for bucket in by_hash.values() for r in bucket] # 2단계: 남은 것끼리 점수 탐욕 매칭. scored = sorted( ((match_score(p, c), i, j) for i, p in enumerate(remaining_left) for j, c in enumerate(remaining_right)), key=lambda t: (-t[0], t[1], t[2]), ) used_left: set[int] = set() used_right: set[int] = set() for score, i, j in scored: if score < MATCH_THRESHOLD or i in used_left or j in used_right: continue used_left.add(i) used_right.add(j) pairs.append((remaining_left[i], remaining_right[j])) unmatched_new = [c for j, c in enumerate(remaining_right) if j not in used_right] unmatched_gone = [p for i, p in enumerate(remaining_left) if i not in used_left] return pairs, unmatched_new, unmatched_gone ``` --- ## 5. 변경 탐지 알고리즘 ### 5.1 판정 규칙과 의사코드 | 판정 | 규칙 | 이벤트 | |---|---|---| | **신규** | 오늘 그룹에만 존재하고 짝을 못 찾음 | `NEW` | | **재등장** | `NEW` 인데 `records.status = 'WITHDRAWN'` 인 이력이 있음 | `NEW` + `is_reappearance=1`, 등급 `HIGH` | | **변경** | 짝지어졌고 `content_hash` 가 다름 | `CHANGED` + 필드별 변경 목록 | | **취하** | 어제 그룹에만 존재하고 짝을 못 찾음 | `WITHDRAWN`, 등급 `CRITICAL` | | **동일** | 짝지어졌고 `content_hash` 가 같음 | 없음(`unchanged_count` 만 증가) | ``` FUNCTION compute_diff(previous, current): IF previous 가 비어 있음: # 기준선 수립일. 전량을 NEW 로 만들면 첫날 리포트가 1만 건 신규가 된다. RETURN DiffResult(new=[], changed=[], withdrawn=[], unchanged_count=len(current), baseline=True) groups ← previous 와 current 를 group_base(dmf_key) 로 묶어 합집합 순회 FOR EACH base IN groups: pairs, only_current, only_previous ← pair_group(prev[base], curr[base]) FOR EACH (before, after) IN pairs: IF before.content_hash == after.content_hash: unchanged_count += 1 ELSE: changes ← diff_fields(before, after) IF changes 가 비어 있음: # 방어: 해시는 다른데 필드는 같다 unchanged_count += 1 # (해시 알고리즘 변경 등) ELSE: changed.append(CHANGED 이벤트) FOR EACH rec IN only_current: new.append(NEW 이벤트) FOR EACH rec IN only_previous: withdrawn.append(WITHDRAWN 이벤트) RETURN DiffResult(정렬된 세 목록, unchanged_count) ``` **diff 는 DB 를 모른다.** `previous` 는 `repo.load_snapshot(마지막 성공 run)` 이 만들어 인자로 넘긴다. 순수 함수라서 픽스처 두 개만 있으면 전 경로를 테스트할 수 있다(아키텍처 §3.7). **`is_reappearance` 는 diff 가 판정하지 않는다.** 순수 함수는 과거 이력을 모른다. `repo.insert_events()` 가 `records.status = 'WITHDRAWN'` 을 조회해 플래그를 세우고 등급을 `INFO` → `HIGH` 로 올린다. 판정 책임을 아는 계층에 둔다. ### 5.2 필드별 변경 중요도 등급 **등급 결정은 2단계다.** ① 표시값이 다른가 → 이벤트가 생긴다. ② 매칭키도 다른가 → 실질 변경이다. 키까지 같으면 `COSMETIC`. | 필드 | 한글명 | 키도 다를 때 등급 | 근거 | 표시값만 다를 때 | |---|---|---|---|---| | `manufacturer` | 제조소명 | **CRITICAL** | 제조원 변경은 공급 리스크의 1순위 신호. 「의약품 등의 안전에 관한 규칙」 제17조제1항제1호의 중요 변경 축 | `COSMETIC` | | `manufacture_place` | 제조소 소재지 | **CRITICAL** | 제조소 이전·주소 변경은 GMP 실사 대상 변경 | `COSMETIC` | | `countries` | 제조국가 | **CRITICAL** | 제조국 변경 = 지정학 리스크·통관·실사 체계 전부 변경 | `COSMETIC`(순서 변화는 정렬로 이미 흡수) | | `applicant` | 신청인 | **HIGH** | 양도양수 가능성. 처리기한 25일 조항이 붙는 절차 | `COSMETIC` | | `ingredient_name` | 성분명 | **HIGH** | 성분이 바뀌는 것은 정상이 아니다. 원본 정정이거나 등록 재정의 | `COSMETIC` | | `permit_date` | 발급일자 | **MEDIUM** | 재공고·정정. 등록 자체는 유지된다 | (키=표시값이라 해당 없음) | **이벤트 등급 = 그 이벤트에 속한 필드 변경 등급의 최댓값.** 단, 모든 필드가 `COSMETIC` 이면 이벤트 등급도 `COSMETIC` 이고, 리포트 '오늘 변경분' 시트에서 **기본 제외**된다(메타 시트에 건수만 표시). | 이벤트 타입 | 기본 등급 | 승격 조건 | |---|---|---| | `WITHDRAWN` | **CRITICAL** | 항상. 워치리스트 무관 | | `NEW` | `INFO` | 워치리스트 매칭 시 `HIGH`, 재등장이면 `HIGH` | | `CHANGED` | 필드 최댓값 | 워치리스트 매칭 시 최소 `HIGH` | > 워치리스트 매칭에 의한 승격은 diff 가 아니라 `repo.insert_events()` → `watchlist_hits` 계산 뒤에 수행한다. diff 의 순수성을 지키기 위해서다. ### 5.3 diff 코드 전문 ```python # src/dmf_crawler/diff.py (2/3 — 판정) """전일 스냅샷과 금일 레코드를 비교해 도메인 이벤트를 만든다. DB·네트워크·파일에 접근하지 않는 순수 함수 모듈이다(아키텍처 §3.7). """ from __future__ import annotations from dataclasses import dataclass from typing import Literal, Mapping from dmf_crawler.normalize import ( COMPARE_FIELDS, FIELD_KEY_MAP, FIELD_LABELS_KO, DmfRecord, ) EventType = Literal["NEW", "CHANGED", "WITHDRAWN"] ChangeClass = Literal["CRITICAL", "HIGH", "MEDIUM", "INFO", "COSMETIC"] SEVERITY_ORDER: Mapping[str, int] = { "CRITICAL": 0, "HIGH": 1, "MEDIUM": 2, "INFO": 3, "COSMETIC": 4, } # 실질 변경(매칭키까지 다름)일 때의 필드별 등급. FIELD_SEVERITY: Mapping[str, ChangeClass] = { "manufacturer": "CRITICAL", "manufacture_place": "CRITICAL", "countries": "CRITICAL", "applicant": "HIGH", "ingredient_name": "HIGH", "permit_date": "MEDIUM", } @dataclass(frozen=True, slots=True) class FieldChange: field: str label: str before: str after: str change_class: ChangeClass @dataclass(frozen=True, slots=True) class DiffEvent: dmf_key: str permit_no: str event_type: EventType severity: ChangeClass changes: tuple[FieldChange, ...] before: dict[str, str] | None after: dict[str, str] | None before_hash: str | None after_hash: str | None @property def changed_fields(self) -> str: return ",".join(sorted(c.field for c in self.changes)) @dataclass(frozen=True, slots=True) class DiffResult: new: tuple[DiffEvent, ...] changed: tuple[DiffEvent, ...] withdrawn: tuple[DiffEvent, ...] unchanged_count: int baseline: bool = False @property def is_empty(self) -> bool: return not (self.new or self.changed or self.withdrawn) @property def total_events(self) -> int: return len(self.new) + len(self.changed) + len(self.withdrawn) def churn_ratio(self, population: int) -> float: """전체 대비 변동 비율. §5.4 게이트 6 의 입력.""" return (self.total_events / population) if population else 0.0 def cosmetic_count(self) -> int: return sum(1 for e in self.changed if e.severity == "COSMETIC") def classify(field: str, before: DmfRecord, after: DmfRecord) -> ChangeClass: """표시값은 다르지만 매칭키가 같으면 COSMETIC(원본 표기 정정).""" key_field = FIELD_KEY_MAP[field] if before.key_view()[key_field] == after.key_view()[key_field]: return "COSMETIC" return FIELD_SEVERITY[field] def diff_fields(before: DmfRecord, after: DmfRecord) -> tuple[FieldChange, ...]: """COMPARE_FIELDS 를 순서대로 비교해 달라진 것만 돌려준다.""" bv, av = before.compare_view(), after.compare_view() out: list[FieldChange] = [] for field in COMPARE_FIELDS: if bv[field] == av[field]: continue out.append(FieldChange( field=field, label=FIELD_LABELS_KO[field], before=bv[field], after=av[field], change_class=classify(field, before, after), )) return tuple(out) def _worst(classes: tuple[ChangeClass, ...], default: ChangeClass) -> ChangeClass: if not classes: return default return min(classes, key=lambda c: SEVERITY_ORDER[c]) def compute_diff(previous: Mapping[str, DmfRecord], current: Mapping[str, DmfRecord]) -> DiffResult: """전일 스냅샷 vs 금일 레코드. 순수 함수.""" if not previous: # 기준선 수립일 — 전량을 NEW 로 만들지 않는다(아키텍처 §3.7). return DiffResult(new=(), changed=(), withdrawn=(), unchanged_count=len(current), baseline=True) prev_groups: dict[str, list[DmfRecord]] = {} for key, rec in previous.items(): prev_groups.setdefault(group_base(key), []).append(rec) curr_groups: dict[str, list[DmfRecord]] = {} for key, rec in current.items(): curr_groups.setdefault(group_base(key), []).append(rec) new_events: list[DiffEvent] = [] changed_events: list[DiffEvent] = [] withdrawn_events: list[DiffEvent] = [] unchanged = 0 for base in sorted(set(prev_groups) | set(curr_groups)): pairs, only_current, only_previous = pair_group( prev_groups.get(base, []), curr_groups.get(base, [])) for before, after in pairs: if before.content_hash == after.content_hash: unchanged += 1 continue changes = diff_fields(before, after) if not changes: # 해시는 다른데 필드 비교로는 같다. 해시 정의가 바뀐 경우의 방어. unchanged += 1 continue changed_events.append(DiffEvent( dmf_key=after.dmf_key, permit_no=after.permit_no, event_type="CHANGED", severity=_worst(tuple(c.change_class for c in changes), "INFO"), changes=changes, before=before.compare_view(), after=after.compare_view(), before_hash=before.content_hash, after_hash=after.content_hash, )) for rec in only_current: new_events.append(DiffEvent( dmf_key=rec.dmf_key, permit_no=rec.permit_no, event_type="NEW", severity="INFO", changes=(), before=None, after=rec.compare_view(), before_hash=None, after_hash=rec.content_hash, )) for rec in only_previous: withdrawn_events.append(DiffEvent( dmf_key=rec.dmf_key, permit_no=rec.permit_no, event_type="WITHDRAWN", severity="CRITICAL", changes=(), before=rec.compare_view(), after=None, before_hash=rec.content_hash, after_hash=None, )) def sort_key(e: DiffEvent) -> tuple[int, str]: return (SEVERITY_ORDER[e.severity], e.dmf_key) return DiffResult( new=tuple(sorted(new_events, key=sort_key)), changed=tuple(sorted(changed_events, key=sort_key)), withdrawn=tuple(sorted(withdrawn_events, key=sort_key)), unchanged_count=unchanged, baseline=False, ) ``` ### 5.4 안전장치 — 소스가 일부만 반환했을 때 취하로 오판하지 않기 이것이 이 프로젝트에서 **가장 위험한 실패 모드**다. API 가 절반만 돌려주면 나머지 절반이 전부 `WITHDRAWN` 이 되고, 리포트는 "오늘 5,000건 취하" 라는 거짓말을 인쇄한다. 그리고 그 잘못된 스냅샷이 내일의 기준선이 되어 **모레는 5,000건 신규**가 뜬다. 한 번의 부분 응답이 사흘을 오염시킨다. **게이트 5종(차단) + 게이트 6(신설, 조건부 차단).** 하나라도 실패하면 **diff 를 수행하지 않고, 스냅샷 INSERT 도 하지 않는다.** | # | 게이트 | 임계값(설정 키) | 실패 시 | |---|---|---|---| | 1 | 모든 페이지 HTTP 200 · `resultCode == "00"` · 본문 시그니처 통과 | 고정 | **차단** | | 2 | `len(records) == total_count_reported` | 고정(완전 일치) | **차단** | | 3 | `total_count` 가 직전 성공 대비 급감하지 않음 | `integrity.max_drop_ratio` = 0.05 | **차단** | | 4 | 필수 3필드 널 비율 ≤ 임계 | `integrity.max_null_ratio` = 0.01 | **차단** | | 5 | 중복 등록번호 비율 ≤ 임계 | `integrity.max_duplicate_ratio` = 0.02 | **차단** | | 6 | **변동률(churn) ≤ 임계** — 신설 | `integrity.max_churn_ratio` = 0.10 | **차단** | **게이트 6 을 신설하는 이유**: 게이트 2·3 은 **API 가 자기 `totalCount` 와 일관되게 거짓말할 때** 뚫린다. 예를 들어 백엔드가 특정 조건의 레코드를 통째로 누락한 채 `totalCount` 도 함께 줄여 응답하면 게이트 2 는 통과하고, 감소폭이 4%면 게이트 3 도 통과한다. 그런데 그 4%가 400건 취하로 리포트에 실린다. **diff 결과 자체를 보고 판단하는 마지막 관문**이 필요하다. - 게이트 6 은 **diff 계산 뒤에 판정**한다. 다른 게이트와 순서가 다르므로 `integrity.evaluate_post_diff()` 라는 별도 진입점을 둔다. - 실패해도 **이미 계산된 diff 는 버리고 스냅샷은 저장하지 않는다.** 실행 상태는 `PARTIAL`, 알림은 `CRITICAL`. - 제도 변화로 진짜 대량 변동이 일어날 수 있다(리서치 §2.9: 연도별 등록이 237→653건으로 튄 사례). 그래서 차단이 **영구 봉쇄가 아니다.** 알림 문구에 "정상 변동이라면 `config.local.toml` 의 `integrity.max_churn_ratio` 를 올리고 `python -m dmf_crawler run --force` 로 재실행하라" 는 다음 행동을 넣는다(요구 R7 의 알림 4요소). ```python # src/dmf_crawler/integrity.py """요구 R2.2 안전장치. diff 로 가는 유일한 관문.""" from __future__ import annotations from dataclasses import dataclass from typing import Optional from dmf_crawler.diff import DiffResult from dmf_crawler.normalize import DmfRecord, NormalizeStats @dataclass(frozen=True, slots=True) class Gate: name: str passed: bool blocking: bool observed: Optional[float] threshold: Optional[float] detail: str @dataclass(frozen=True, slots=True) class IntegrityVerdict: ok: bool gates: tuple[Gate, ...] blocking_reason: Optional[str] @property def failed(self) -> tuple[Gate, ...]: return tuple(g for g in self.gates if not g.passed) @dataclass(frozen=True, slots=True) class PrevSnapshotStats: run_id: str run_date: str total_count_reported: int records_normalized: int payload_sha256: str def evaluate(fetch, records: list[DmfRecord], stats: NormalizeStats, prev: Optional[PrevSnapshotStats], cfg) -> IntegrityVerdict: """게이트 1~5. diff 이전에 판정한다.""" n = max(len(records), 1) gates: list[Gate] = [] # 게이트 1 — 전송 계층 건전성 transport_ok = bool(fetch.body_signature_ok) and fetch.result_code == "00" gates.append(Gate( name="transport_ok", passed=transport_ok, blocking=True, observed=None, threshold=None, detail=(f"resultCode={fetch.result_code!r} " f"body_signature_ok={fetch.body_signature_ok} " f"pages={fetch.pages_fetched}/{fetch.pages_expected}"), )) # 게이트 2 — 완결성 reported = fetch.total_count_reported complete = reported is not None and len(fetch.records) == reported gates.append(Gate( name="total_count_match", passed=complete, blocking=True, observed=float(len(fetch.records)), threshold=float(reported or 0), detail=f"수집 {len(fetch.records)}건 vs totalCount {reported}", )) # 게이트 3 — 전일 대비 급감 if prev is None or prev.total_count_reported <= 0: gates.append(Gate("drop_ratio", True, True, None, None, "직전 성공 스냅샷이 없다(기준선 수립). 판정 생략")) else: drop = (prev.total_count_reported - (reported or 0)) / prev.total_count_reported gates.append(Gate( name="drop_ratio", passed=drop <= cfg.max_drop_ratio, blocking=True, observed=round(drop, 6), threshold=cfg.max_drop_ratio, detail=(f"{prev.run_date} {prev.total_count_reported}건 → " f"오늘 {reported}건 (감소율 {drop:.2%})"), )) # 게이트 4 — 필수 필드 널 비율 worst_field, worst_ratio = "", 0.0 for field, count in stats.null_counts.items(): ratio = count / n if ratio > worst_ratio: worst_field, worst_ratio = field, ratio gates.append(Gate( name="null_ratio", passed=worst_ratio <= cfg.max_null_ratio, blocking=True, observed=round(worst_ratio, 6), threshold=cfg.max_null_ratio, detail=(f"최악 필드 {worst_field or '-'} 널 {worst_ratio:.3%} " f"(전체 {dict(stats.null_counts)})"), )) # 게이트 5 — 중복 등록번호 비율 dup_ratio = stats.duplicate_permit_no / n gates.append(Gate( name="duplicate_ratio", passed=dup_ratio <= cfg.max_duplicate_ratio, blocking=True, observed=round(dup_ratio, 6), threshold=cfg.max_duplicate_ratio, detail=(f"중복 {stats.duplicate_permit_no}건 / {stats.duplicate_groups}종 " f"({dup_ratio:.3%})"), )) failed = [g for g in gates if not g.passed and g.blocking] reason = None if not failed else " / ".join(f"{g.name}: {g.detail}" for g in failed) return IntegrityVerdict(ok=not failed, gates=tuple(gates), blocking_reason=reason) def evaluate_post_diff(diff: DiffResult, population: int, cfg) -> IntegrityVerdict: """게이트 6 — diff 결과를 보고 내리는 마지막 판정. API 가 자기 totalCount 와 일관되게 축소 응답하면 게이트 2·3 은 통과한다. 변동률 자체를 보는 관문이 필요하다. """ if diff.baseline: return IntegrityVerdict( ok=True, gates=(Gate("churn_ratio", True, True, 0.0, cfg.max_churn_ratio, "기준선 수립일. 판정 생략"),), blocking_reason=None) churn = diff.churn_ratio(population) withdrawn_ratio = (len(diff.withdrawn) / population) if population else 0.0 gates = ( Gate(name="churn_ratio", passed=churn <= cfg.max_churn_ratio, blocking=True, observed=round(churn, 6), threshold=cfg.max_churn_ratio, detail=(f"신규 {len(diff.new)} / 변경 {len(diff.changed)} / " f"취하 {len(diff.withdrawn)} = 모집단 {population} 대비 {churn:.2%}")), Gate(name="withdrawn_ratio", passed=withdrawn_ratio <= cfg.max_withdrawn_ratio, blocking=False, observed=round(withdrawn_ratio, 6), threshold=cfg.max_withdrawn_ratio, detail=f"취하 {len(diff.withdrawn)}건 ({withdrawn_ratio:.2%})"), ) failed = [g for g in gates if not g.passed and g.blocking] reason = None if not failed else " / ".join(f"{g.name}: {g.detail}" for g in failed) return IntegrityVerdict(ok=not failed, gates=gates, blocking_reason=reason) ``` **설정 키 2개 추가** (아키텍처 §6.1 의 `[integrity]` 절에 이어 붙인다). | 키 | 타입 | 기본값 | 설명 | |---|---|---|---| | `integrity.max_churn_ratio` | float | `0.10` | (신규+변경+취하)/모집단 상한. 초과 시 diff 폐기 + CRITICAL | | `integrity.max_withdrawn_ratio` | float | `0.02` | 취하 비율 경고선. **비차단** — 초과 시 리포트 배너 + WARN | ### 5.5 기준선·스테일·강제 실행 | 상황 | 판정 | 저장 | 리포트 | 알림 | |---|---|---|---|---| | **첫 실행** (`previous` 비어 있음) | `baseline=True`, 이벤트 0건 | 스냅샷 저장 ○ | "기준선 수립일 — 내일부터 변경이 표시됩니다" 배너 | 없음 | | 게이트 1~5 실패 | `integrity_ok=0`, diff 미수행 | 스냅샷 저장 **✕** | **마지막 성공 스냅샷**으로 생성 + "오늘 수집 실패 — 마지막 성공: YYYY-MM-DD" 배너 | `CRITICAL` | | 게이트 6 실패 | diff 계산됨 → **폐기** | 스냅샷 저장 **✕** | 위와 동일 + 변동률 수치 명시 | `CRITICAL` | | 소스 서킷 OPEN | fetch 스킵 | 저장 없음 | 마지막 성공 스냅샷 + 스테일 배너 | 상태 전환 시 1회 | | `--force` 재실행 | 게이트 6 임계를 설정으로 올린 뒤 재실행 | 정상 | 정상 | 없음 | | 소스 미갱신(`payload_sha256` 동일) | 정상 실행, 이벤트 0건 | 정상 | "소스 데이터가 어제와 동일합니다" 표기 | 없음 | **연속 스테일의 처리**: `payload_sha256` 이 **7일 연속** 동일하면 소스가 갱신을 멈춘 것일 수도, 우리가 캐시된 응답을 받는 것일 수도 있다. `INFO` 알림 1회를 기록하고 리포트 메타 시트에 "소스 최종 변동: N일 전" 을 표시한다. 요구 SSOT 7.4 의 "갱신 주기 실측" 이 여기서 자동으로 이뤄진다. ### 5.6 채택하지 않은 이벤트 타입 — 기록 리서치 §10.3 은 9종을 제안했다. 공식 API 확정으로 6종의 근거 필드가 사라졌다. **폐기가 아니라 보류**이며, 소스가 확장되면 되살릴 지점을 남긴다. | 리서치 이벤트 | 필요 필드 | 상태 | 되살리는 조건 | |---|---|---|---| | `CHANGED_REG` | `최종변경일자` | **보류** | 화면 컬럼을 주는 소스가 생기면 | | `ANNUAL_REPORT` | `최종연차보고년도` | **보류** | 동상 | | `DOC_VERSION_CHANGED` | `문서번호` | **보류** | 동상 | | `LINKED_REVIEW_SET` | `연계심사문서번호` | **보류** | 동상 | | `SITE_CHANGED` | 제조소 3필드 | **흡수됨** | `CHANGED` + `change_class=CRITICAL` 로 이미 표현된다 | | `APPLICANT_CHANGED` | `신청인` | **흡수됨** | `CHANGED` + `applicant` 필드 등급 `HIGH` | | `DISAPPEARED` | — | **통합됨** | `WITHDRAWN` 과 구분할 근거가 API 에 없다. 하나로 합치고 §5.4 게이트로 방어 | | `REAPPEARED` | — | **채택(변형)** | `NEW` + `is_reappearance=1` | | `NEW` / `MODIFIED` / `WITHDRAWN` | — | **채택** | — | --- ## 6. 이벤트 스키마 ### 6.1 이벤트 레코드 구조 이벤트는 **세 곳에 같은 내용이 다른 형태로** 남는다. 용도가 다르기 때문이다. | 저장소 | 형태 | 용도 | 보존 | |---|---|---|---| | `events` + `event_changes` 테이블 | 관계형 | 리포트 SQL, 이력 조회, 집계 | **영구** | | `logs/run_*/events.jsonl` | JSON Lines | 기계 판독 실행 로그, 장애 분석 | 90일 | | `agy` 프롬프트 입력 | 절단·요약된 JSON | AI 브리핑 | 미보존(호출 원문은 `agy.stdout.json`) | **필드 정의 (`events` 테이블 기준).** | 필드 | 타입 | 널 | 의미 | |---|---|---|---| | `event_id` | INTEGER | ✕ | 자동 증가 PK | | `run_seq` | INTEGER | ✕ | 이 이벤트를 만든 실행 | | `event_date` | TEXT | ✕ | `YYYY-MM-DD` (KST 실행일자) | | `occurred_at` | TEXT | ✕ | ISO 8601 `+09:00`. **관측 시각**이지 실제 등록 변경 시각이 아니다 | | `record_id` | INTEGER | ✕ | `records` 대리키 | | `dmf_key` | TEXT | ✕ | 비정규화(리포트 조인 절약) | | `permit_no` | TEXT | ✕ | 비정규화. nedrug 검색 링크 생성에 쓴다 | | `event_type` | TEXT | ✕ | `NEW` / `CHANGED` / `WITHDRAWN` | | `severity` | TEXT | ✕ | `CRITICAL` / `HIGH` / `MEDIUM` / `INFO` / `COSMETIC` | | `is_reappearance` | INTEGER | ✕ | `NEW` 이면서 과거 취하 이력이 있으면 1 | | `change_count` | INTEGER | ✕ | 변경된 필드 수. `NEW`·`WITHDRAWN` 은 0 | | `changed_fields` | TEXT | ✕ | 정렬된 CSV. 필터·인덱스용 | | `before_version_id` / `after_version_id` | INTEGER | ○ | 버전 테이블 포인터 | | `changes_json` | TEXT | ✕ | 필드별 변경 배열 | | `before_json` / `after_json` | TEXT | ○ | `COMPARE_FIELDS` 스냅샷. 버전 행이 정리된 뒤에도 이벤트만으로 읽히게 | **`occurred_at` 이 "관측 시각" 임을 문서에 못 박는 이유**: API 는 변경 시점을 주지 않는다. 우리가 아는 것은 "06:00 에 봤더니 달랐다" 뿐이다. 이 필드를 실제 변경 시각으로 오해하면 심사 소요일 같은 파생 지표가 전부 틀어진다. 리포트에도 "관측일" 로 표기한다. ### 6.2 예시 JSON **`CHANGED` — 제조소 변경 (CRITICAL)** ```json { "event_id": 48213, "run_id": "run_20260902_060013", "event_date": "2026-09-02", "occurred_at": "2026-09-02T06:01:44+09:00", "dmf_key": "20121228-168-I-169-04", "permit_no": "20121228-168-I-169-04", "event_type": "CHANGED", "severity": "CRITICAL", "is_reappearance": 0, "change_count": 2, "changed_fields": "countries,manufacturer", "before_hash": "9f2c1ab4e7d05631", "after_hash": "c07be1d3aa914f28", "changes": [ { "field": "manufacturer", "label": "제조소명", "before": "SICOR SOCIETA'ITALIANA CORTICOSTEROIDO S.R.L. , [미분화공정 제조소] Micro-Macinazione SA.", "after": "SICOR SOCIETA'ITALIANA CORTICOSTEROIDO S.R.L.", "class": "CRITICAL" }, { "field": "countries", "label": "제조국가", "before": "스위스, 이탈리아", "after": "이탈리아", "class": "CRITICAL" } ], "before": { "ingredient_name": "포르모테롤푸마르산염수화물", "applicant": "(주)대웅제약", "manufacturer": "SICOR SOCIETA'ITALIANA CORTICOSTEROIDO S.R.L. , [미분화공정 제조소] Micro-Macinazione SA.", "manufacture_place": "Rho(MI) - Via Terrazzano, 77, Italy , 6995 Madonna del Piano, Switzerland", "countries": "스위스, 이탈리아", "permit_date": "2015-02-26" }, "after": { "ingredient_name": "포르모테롤푸마르산염수화물", "applicant": "(주)대웅제약", "manufacturer": "SICOR SOCIETA'ITALIANA CORTICOSTEROIDO S.R.L.", "manufacture_place": "Rho(MI) - Via Terrazzano, 77, Italy", "countries": "이탈리아", "permit_date": "2015-02-26" } } ``` **`CHANGED` — 원본 오탈자 정정 (COSMETIC, 리포트 기본 제외)** ```json { "event_id": 48214, "run_id": "run_20260902_060013", "event_date": "2026-09-02", "occurred_at": "2026-09-02T06:01:44+09:00", "dmf_key": "20121228-168-I-169-05", "permit_no": "20121228-168-I-169-05", "event_type": "CHANGED", "severity": "COSMETIC", "is_reappearance": 0, "change_count": 1, "changed_fields": "manufacturer", "before_hash": "1d4a7f90cc23b805", "after_hash": "5b8e02c1df647a3e", "changes": [ { "field": "manufacturer", "label": "제조소명", "before": "SICOR SOCIETA'ITALIANA CORTICOSTER OIDO S.R.L.", "after": "SICOR SOCIETA'ITALIANA CORTICOSTEROIDO S.R.L.", "class": "COSMETIC" } ], "note": "manufacturer_key 가 동일하므로 표기 정정으로 판정. 대시보드 KPI 에 포함되지 않는다." } ``` **`NEW` — 신규 등록 (재등장 아님)** ```json { "event_id": 48215, "run_id": "run_20260902_060013", "event_date": "2026-09-02", "occurred_at": "2026-09-02T06:01:44+09:00", "dmf_key": "20260901-209-J-2270", "permit_no": "20260901-209-J-2270", "event_type": "NEW", "severity": "INFO", "is_reappearance": 0, "change_count": 0, "changed_fields": "", "before_hash": null, "after_hash": "8ac3d5e0117b94f2", "changes": [], "before": null, "after": { "ingredient_name": "탐스로신염산염", "applicant": "주식회사지맥스파마켐", "manufacturer": "Hema Pharmaceuticals Pvt. Ltd.", "manufacture_place": "Plot No. 6201/A & B, G.I.D.C., Gujarat State, India", "countries": "인도", "permit_date": "2026-09-01" }, "permit_parts": { "fmt": "standard", "accept_date": "2026-09-01", "ingr_no": 209, "group": "J", "group_table": "별표1", "group_effective_from": "2017-12-25", "serial": 2270, "sub": null, "grant": null, "ingr_group_key": "209-J" } } ``` **`WITHDRAWN` — 목록 이탈 (항상 CRITICAL)** ```json { "event_id": 48216, "run_id": "run_20260902_060013", "event_date": "2026-09-02", "occurred_at": "2026-09-02T06:01:44+09:00", "dmf_key": "20110531-71-B-317-05", "permit_no": "20110531-71-B-317-05", "event_type": "WITHDRAWN", "severity": "CRITICAL", "is_reappearance": 0, "change_count": 0, "changed_fields": "", "before_hash": "44f1c8b09e2d7a61", "after_hash": null, "changes": [], "before": { "ingredient_name": "부데소니드", "applicant": "(주)하이플", "manufacturer": "Avik Pharmaceutical Limited", "manufacture_place": "Plot No. 21, Gujarat, India", "countries": "인도", "permit_date": "2011-05-31" }, "after": null, "note": "API 는 취하 사유를 제공하지 않는다. 리포트는 이 행에 nedrug 검색 링크를 붙여 사람이 확인하게 한다." } ``` ### 6.3 `events.jsonl` 라인 포맷 파이프라인 실행 로그. 도메인 이벤트뿐 아니라 **스테이지 전이·게이트 판정·알림 기록도 같은 파일에** 한 줄씩 쌓는다. 장애 분석 시 시간순 단일 뷰가 필요하기 때문이다. ```jsonl {"ts":"2026-09-02T06:00:13+09:00","run_id":"run_20260902_060013","kind":"stage","stage":"fetch","status":"START"} {"ts":"2026-09-02T06:00:51+09:00","run_id":"run_20260902_060013","kind":"fetch","pages":99,"records":9841,"total_count":9841,"elapsed_s":37.8,"payload_sha256":"3b1f…"} {"ts":"2026-09-02T06:00:52+09:00","run_id":"run_20260902_060013","kind":"stage","stage":"fetch","status":"SUCCESS","duration_ms":38912} {"ts":"2026-09-02T06:01:02+09:00","run_id":"run_20260902_060013","kind":"gate","name":"total_count_match","passed":true,"observed":9841,"threshold":9841} {"ts":"2026-09-02T06:01:02+09:00","run_id":"run_20260902_060013","kind":"gate","name":"drop_ratio","passed":true,"observed":0.0002,"threshold":0.05} {"ts":"2026-09-02T06:01:44+09:00","run_id":"run_20260902_060013","kind":"event","event_type":"CHANGED","severity":"CRITICAL","dmf_key":"20121228-168-I-169-04","changed_fields":"countries,manufacturer"} {"ts":"2026-09-02T06:01:44+09:00","run_id":"run_20260902_060013","kind":"event","event_type":"NEW","severity":"INFO","dmf_key":"20260901-209-J-2270","changed_fields":""} {"ts":"2026-09-02T06:01:45+09:00","run_id":"run_20260902_060013","kind":"diff","new":12,"changed":7,"withdrawn":0,"cosmetic":3,"unchanged":9822,"churn_ratio":0.0019} {"ts":"2026-09-02T06:02:31+09:00","run_id":"run_20260902_060013","kind":"alert","level":"WARN","code":"BACKUP_SKIPPED_LOW_DISK","dedup_key":"backup:lowdisk"} ``` **규칙 4가지.** 1. 한 줄 = 하나의 JSON 객체. 개행 없음. `ensure_ascii=False` (한글을 그대로 읽을 수 있어야 한다). 2. 모든 줄에 `ts`(ISO 8601 `+09:00`) · `run_id` · `kind` 세 필드가 반드시 있다. 3. `kind` 는 `stage` / `fetch` / `gate` / `event` / `diff` / `agy` / `alert` / `error` 8종. 4. **비밀 값은 절대 쓰지 않는다.** `logging.mask_patterns`(`serviceKey`, `access_token`)에 걸리는 키는 기록 직전에 `***` 로 치환한다. --- ## 7. 데이터 품질 검증 ### 7.1 검증의 두 층 | 층 | 이름 | 시점 | 실패 시 | |---|---|---|---| | **차단** | 게이트 (§5.4) | `integrity` 스테이지 / `diff` 직후 | diff 미수행 · 스냅샷 미저장 · `PARTIAL` · `CRITICAL` 알림 | | **비차단** | 품질 지표 | `integrity` 스테이지에서 함께 계산 | 저장·diff·리포트 전부 정상 진행. `quality_checks` 기록 + 등급별 알림 + 리포트 메타 시트 표시 | **모든 결과는 통과·실패 무관하게 `quality_checks` 에 기록한다.** 통과 기록이 있어야 "임계값이 적절한가" 를 나중에 실측으로 조정할 수 있다. ### 7.2 비차단 품질 지표 13종 | # | `check_name` | 관측 대상 | 기본 임계 | 위반 등급 | 위반 시 동작 | 근거 | |---|---|---|---|---|---|---| | 1 | `row_count_range` | 정규화 레코드 수 | 6,000 ≤ N ≤ 30,000 | `WARN` | 리포트 배너 + 알림 | 관측 9,840건 기준 상하 여유. 부트스트랩 후 실측으로 조정 | | 2 | `permit_parse_rate` | `fmt != 'unknown'` 비율 | ≥ 0.97 | `WARN` | 미파싱 목록을 로그에 덤프 | 급락 = 원본 번호 체계 변경 신호 | | 3 | `new_substance_ratio` | `fmt == 'new_substance'` 비율 | ≤ 0.05 | `INFO` | 메타 시트 표기만 | 실측 표본이 1건뿐이라 관측 목적 | | 4 | `unknown_group_count` | `GROUP_TABLE` 에 없는 군 건수 | = 0 | `INFO` | 새 군 값을 알림 본문에 명시 | 제도 확대로 `L` 군이 생길 수 있다 | | 5 | `synthetic_key_count` | `SYN-` 합성키 건수 | = 0 | `WARN` | 메타 시트 경고 + 해당 목록 로그 | 합성키는 오탐에 취약하다(§2.3) | | 6 | `permit_date_valid_rate` | `permit_date != ''` 비율 | ≥ 0.99 | `WARN` | 실패 샘플 20건 로그 | 날짜 표기 변경 탐지 | | 7 | `permit_date_future_count` | 오늘보다 미래인 발급일자 | = 0 | `WARN` | 해당 목록 로그 | 원본 입력 오류 | | 8 | `permit_date_ancient_count` | 1990-01-01 이전 발급일자 | = 0 | `INFO` | 로그만 | DMF 제도 도입이 2002년이다 | | 9 | `date_order_violation_rate` | `permit_date < accept_date` 비율 | ≤ 0.05 | `INFO` | 메타 시트 수치 | 재공고 건은 정상적으로 역전될 수 있다 | | 10 | `country_unresolved_rate` | `COUNTRY_ALIASES` 에 없는 국가명 비율 | ≤ 0.05 | `INFO` | 미등록 국가명 상위 10개 로그 | 별칭 사전 보강 근거 수집 | | 11 | `field_oversize_count` | 명세 크기 초과 필드 건수 | = 0 | `WARN` | 필드별 건수 알림 | **스키마 드리프트 1순위 신호** | | 12 | `encoding_replacement_count` | `U+FFFD` 포함 레코드 수 | = 0 | `WARN` | 해당 원문 아카이브 경로 안내 | 인코딩 사고. 재파싱 필요 | | 13 | `payload_unchanged_days` | `payload_sha256` 연속 동일 일수 | ≤ 7 | `INFO` | 메타 시트 "소스 최종 변동: N일 전" | 요구 SSOT 7.4 갱신 주기 실측 | > `rejected`(정규화 자체가 실패한 원본)는 지표가 아니라 **무조건 로그와 `runs.notes` 에 남긴다.** 1건이라도 있으면 `WARN` 이다. 정상 데이터에서는 절대 발생하지 않아야 한다. ### 7.3 검증 코드 ```python # src/dmf_crawler/integrity.py (계속 — 비차단 품질 지표) from datetime import date ANCIENT_CUTOFF = "1990-01-01" def quality_metrics(records: list[DmfRecord], stats: NormalizeStats, fetch, prev: Optional[PrevSnapshotStats], payload_unchanged_days: int, cfg) -> tuple[Gate, ...]: """비차단 지표 13종. 실패해도 파이프라인을 멈추지 않는다.""" n = max(len(records), 1) today = date.today().isoformat() out: list[Gate] = [] def add(name: str, passed: bool, observed, threshold, detail: str) -> None: out.append(Gate(name=name, passed=passed, blocking=False, observed=None if observed is None else float(observed), threshold=None if threshold is None else float(threshold), detail=detail)) # 1. 건수 범위 add("row_count_range", cfg.min_expected_rows <= len(records) <= cfg.max_expected_rows, len(records), cfg.max_expected_rows, f"{len(records)}건 (기대 {cfg.min_expected_rows}~{cfg.max_expected_rows})") # 2. 등록번호 파싱률 parse_rate = 1.0 - (stats.unparsed_permit_no / n) add("permit_parse_rate", parse_rate >= 0.97, parse_rate, 0.97, f"파싱 실패 {stats.unparsed_permit_no}건 ({1 - parse_rate:.3%})") # 3. 신물질 포맷 비율 ns_ratio = stats.new_substance_count / n add("new_substance_ratio", ns_ratio <= 0.05, ns_ratio, 0.05, f"신물질 포맷 {stats.new_substance_count}건") # 4. 미지 알파벳군 unknown_groups = sorted({r.permit.group for r in records if r.permit.group and r.permit.group_table is None}) add("unknown_group_count", not unknown_groups, len(unknown_groups), 0, f"미등록 군: {unknown_groups or '없음'}") # 5. 합성키 add("synthetic_key_count", stats.synthetic_key_count == 0, stats.synthetic_key_count, 0, f"등록번호 부재로 합성키를 쓴 레코드 {stats.synthetic_key_count}건") # 6. 발급일자 유효율 valid_rate = 1.0 - (stats.invalid_permit_date / n) add("permit_date_valid_rate", valid_rate >= 0.99, valid_rate, 0.99, f"발급일자 파싱 실패 {stats.invalid_permit_date}건") # 7. 미래 일자 future = [r.dmf_key for r in records if r.permit_date and r.permit_date > today] add("permit_date_future_count", not future, len(future), 0, f"미래 발급일자 {len(future)}건 (예: {future[:3]})") # 8. 과거 일자 ancient = [r.dmf_key for r in records if r.permit_date and r.permit_date < ANCIENT_CUTOFF] add("permit_date_ancient_count", not ancient, len(ancient), 0, f"{ANCIENT_CUTOFF} 이전 발급일자 {len(ancient)}건") # 9. 날짜 역전 (발급일자 < 등록수리일자) inverted = [r for r in records if r.permit_date and r.accepted_date and r.permit_date < r.accepted_date] inv_rate = len(inverted) / n add("date_order_violation_rate", inv_rate <= 0.05, inv_rate, 0.05, f"발급일자가 등록수리일자보다 이른 건 {len(inverted)}건 ({inv_rate:.2%})") # 10. 미등록 국가명 from dmf_crawler.normalize import COUNTRY_ALIASES, matching_key unresolved: dict[str, int] = {} country_rows = 0 for rec in records: for country in rec.countries: country_rows += 1 if matching_key(country) not in COUNTRY_ALIASES: unresolved[country] = unresolved.get(country, 0) + 1 unresolved_rate = (sum(unresolved.values()) / country_rows) if country_rows else 0.0 top = sorted(unresolved.items(), key=lambda kv: -kv[1])[:10] add("country_unresolved_rate", unresolved_rate <= 0.05, unresolved_rate, 0.05, f"별칭 미등록 국가 {len(unresolved)}종 상위: {top}") # 11. 필드 길이 초과 (스키마 드리프트) oversize_total = sum(stats.oversize_fields.values()) add("field_oversize_count", oversize_total == 0, oversize_total, 0, f"명세 크기 초과: {dict(stats.oversize_fields)}") # 12. 인코딩 사고 add("encoding_replacement_count", stats.replacement_char_count == 0, stats.replacement_char_count, 0, f"U+FFFD 포함 레코드 {stats.replacement_char_count}건. " f"원문: {fetch.archive_dir}") # 13. 소스 미갱신 연속 일수 add("payload_unchanged_days", payload_unchanged_days <= 7, payload_unchanged_days, 7, f"소스 데이터가 {payload_unchanged_days}일 연속 동일하다") return tuple(out) QUALITY_SEVERITY: Mapping[str, str] = { "row_count_range": "WARN", "permit_parse_rate": "WARN", "new_substance_ratio": "INFO", "unknown_group_count": "INFO", "synthetic_key_count": "WARN", "permit_date_valid_rate": "WARN", "permit_date_future_count": "WARN", "permit_date_ancient_count": "INFO", "date_order_violation_rate": "INFO", "country_unresolved_rate": "INFO", "field_oversize_count": "WARN", "encoding_replacement_count": "WARN", "payload_unchanged_days": "INFO", } ``` ### 7.4 위반 시 동작 매트릭스 | 위반 조합 | 실행 상태 | 스냅샷 저장 | diff | 리포트 | 알림 | 다음 실행에 미치는 영향 | |---|---|---|---|---|---|---| | 없음 | `SUCCESS` | ○ | ○ | 정상 | 없음 | 없음 | | 비차단 `INFO` 만 | `SUCCESS` | ○ | ○ | 메타 시트에 표기 | 없음 | 없음 | | 비차단 `WARN` 1개 이상 | `SUCCESS` | ○ | ○ | 대시보드 주의 배너 + 메타 시트 상세 | `WARN` 1건(쿨다운 4시간) | 없음 | | 차단 게이트 1~5 실패 | `PARTIAL` | **✕** | **✕** | 마지막 성공 스냅샷 + 실패 배너 | `CRITICAL` | **기준선 불변** — 다음 실행은 여전히 마지막 성공과 비교 | | 차단 게이트 6 실패 | `PARTIAL` | **✕** | 계산 후 폐기 | 위와 동일 + 변동률 명시 | `CRITICAL` | 동일 | | 정규화 `rejected` ≥ 1 | 다른 조건 따름 | 다른 조건 따름 | 다른 조건 따름 | 메타 시트에 건수·원문 경로 | `WARN` | 없음 | | 3회 연속 차단 | `PARTIAL` | ✕ | ✕ | 스테일 배너 강조 | `CRITICAL` **강제 모달** | 소스 서킷 `OPEN` 전이 | **핵심 불변식**: 차단이 걸린 날은 **기준선이 움직이지 않는다.** 그래서 소스가 사흘 뒤 정상으로 돌아오면 그날의 diff 는 "마지막 성공일 대비 사흘치 변화" 를 정확히 보여준다. 변화를 놓치는 것이 아니라 **미루는 것**이다. 리포트는 그 사실을 "비교 기준: 2026-08-30 (3일 전)" 로 명시해야 한다. --- ## 8. 보존 정책 ### 8.1 자산별 보존 기간 | 자산 | 경로 / 테이블 | 기본 보존 | 설정 키 | 정리 방법 | 정리 주체 | |---|---|---|---|---|---| | 도메인 이벤트 | `events`, `event_changes` | **영구** | 없음(고정) | 삭제하지 않는다 | — | | 레코드 버전 | `record_versions` | **영구** | 없음(고정) | 삭제하지 않는다 | — | | 현재 상태 | `records` | **영구** | 없음(고정) | 삭제하지 않는다 | — | | 스냅샷 멤버십 | `snapshots` | 영구(=`0`) | `storage.snapshot_retain_days` | `run_seq` 기준 배치 DELETE + `VACUUM` | `repo.prune_snapshots()` | | 실행 기록 | `runs`, `stage_status`, `fetch_stats`, `quality_checks` | 영구 | 없음 | 삭제하지 않는다(행이 작다) | — | | AI 산출물 | `agy_calls`, `enrichment*` | 영구 | 없음 | 삭제하지 않는다 | — | | 알림 | `alerts` | 영구 | 없음 | 삭제하지 않는다 | — | | **API 원문** | `data/raw/YYYY-MM-DD/` | **180일** | `source.archive_retain_days` | 날짜 디렉터리 통째 삭제 | `backup.prune_raw_archive()` | | **실행 로그** | `logs/run_*/` | **90일** | `logging.retain_days` | 디렉터리 통째 삭제 | `backup.prune_logs()` | | **리포트** | `reports/DMF_리포트_*.xlsx` | **365일** | `report.retain_days` | 파일 삭제. `*_최신.xlsx` 는 제외 | `backup.prune_reports()` | | **DB 백업** | `backup/dmf_*.sqlite3` | **최근 30개** | `backup.keep_count` | 오래된 것부터 삭제 | `backup.prune_backups()` | | 마이그레이션 전 백업 | `backup/premigrate_*.sqlite3` | **최근 5개** | 고정 | 오래된 것부터 삭제 | `backup.prune_backups()` | | AI 제안 | `state/proposals/` | 사람이 처리할 때까지 | 없음 | **자동 삭제하지 않는다** | — | **보존 기간의 근거.** - `data/raw/` 180일 — 재파싱·회귀 픽스처의 원자료다. 반년이면 정규화 규칙을 고쳤을 때 과거 데이터를 다시 만들어볼 수 있는 충분한 창이다. 하루 ~99페이지 × ~40KB ≈ 4MB/일 → 180일 = **720MB**. 이게 이 프로젝트에서 두 번째로 큰 디스크 소비처다. - `logs/` 90일 — 분기 단위 장애 회고에 필요한 최소. 하루 ~2MB → 180MB. - `reports/` 365일 — 전년 동기 리포트를 열어볼 수 있어야 한다. 파일당 ~2MB → 730MB. - `backup/` 30개 — 한 달치 일일 백업. `VACUUM INTO` 산출물은 원본과 비슷한 크기(~150MB 가정) → **4.5GB**. **가장 큰 소비처이며 다른 드라이브를 권장하는 이유다.** **총 디스크 예산 (1년 운영 시)**: DB 150MB + raw 720MB + logs 180MB + reports 730MB + backup 4.5GB ≈ **6.3GB**. 온보딩 GUI 가 이 숫자를 보여주고 백업 드라이브 선택을 유도한다. ### 8.2 정리 코드 ```python # src/dmf_crawler/backup.py (보존 정리 함수군) """VACUUM INTO 백업과 자산별 보존 정리. 정리는 finalize 스테이지에서만 수행한다. required=False 취급이라 실패해도 실행 상태를 FAILED 로 끌어내리지 않는다. """ from __future__ import annotations import re import shutil from dataclasses import dataclass from datetime import date, datetime, timedelta from pathlib import Path RUN_DIR_RE = re.compile(r"^run_(?P\d{8}_\d{6})$") DATE_DIR_RE = re.compile(r"^(?P\d{4}-\d{2}-\d{2})$") @dataclass(frozen=True, slots=True) class PruneResult: asset: str removed: int freed_bytes: int errors: tuple[str, ...] def _dir_size(path: Path) -> int: total = 0 for item in path.rglob("*"): if item.is_file(): try: total += item.stat().st_size except OSError: pass return total def prune_raw_archive(raw_root: Path, retain_days: int, today: date | None = None) -> PruneResult: """data/raw/YYYY-MM-DD/ 를 날짜 디렉터리 단위로 지운다.""" if retain_days <= 0 or not raw_root.exists(): return PruneResult("raw", 0, 0, ()) cutoff = (today or date.today()) - timedelta(days=retain_days) removed = freed = 0 errors: list[str] = [] for child in sorted(raw_root.iterdir()): if not child.is_dir(): continue m = DATE_DIR_RE.match(child.name) if not m or m.group("d") >= cutoff.isoformat(): continue size = _dir_size(child) try: shutil.rmtree(child) except OSError as exc: errors.append(f"{child}: {exc}") continue removed += 1 freed += size return PruneResult("raw", removed, freed, tuple(errors)) def prune_logs(logs_root: Path, retain_days: int, today: date | None = None) -> PruneResult: """logs/run_YYYYMMDD_HHMMSS/ 를 디렉터리 단위로 지운다.""" if retain_days <= 0 or not logs_root.exists(): return PruneResult("logs", 0, 0, ()) cutoff = (today or date.today()) - timedelta(days=retain_days) removed = freed = 0 errors: list[str] = [] for child in sorted(logs_root.iterdir()): if not child.is_dir(): continue m = RUN_DIR_RE.match(child.name) if not m: continue try: stamp = datetime.strptime(m.group("stamp"), "%Y%m%d_%H%M%S").date() except ValueError: continue if stamp >= cutoff: continue size = _dir_size(child) try: shutil.rmtree(child) except OSError as exc: errors.append(f"{child}: {exc}") continue removed += 1 freed += size return PruneResult("logs", removed, freed, tuple(errors)) def prune_reports(reports_dir: Path, retain_days: int, latest_name: str, today: date | None = None) -> PruneResult: """일자별 리포트를 지운다. 최신본 고정 링크는 절대 지우지 않는다.""" if retain_days <= 0 or not reports_dir.exists(): return PruneResult("reports", 0, 0, ()) cutoff = (today or date.today()) - timedelta(days=retain_days) removed = freed = 0 errors: list[str] = [] for item in sorted(reports_dir.glob("*.xlsx")): if latest_name and item.name == latest_name: continue m = re.search(r"(\d{4}-\d{2}-\d{2})", item.name) if not m or m.group(1) >= cutoff.isoformat(): continue try: size = item.stat().st_size item.unlink() except OSError as exc: # 사용자가 열어둔 파일은 잠겨 있다. 다음에 지운다. errors.append(f"{item.name}: {exc}") continue removed += 1 freed += size return PruneResult("reports", removed, freed, tuple(errors)) def prune_backups(backup_dir: Path, keep_daily: int, keep_premigrate: int = 5) -> PruneResult: """일일 백업과 마이그레이션 전 백업을 각각 개수 기준으로 정리한다.""" if not backup_dir.exists(): return PruneResult("backup", 0, 0, ()) removed = freed = 0 errors: list[str] = [] for pattern, keep in (("dmf_*.sqlite3", keep_daily), ("premigrate_*.sqlite3", keep_premigrate)): files = sorted(backup_dir.glob(pattern), key=lambda p: p.name, reverse=True) for item in files[max(keep, 0):]: try: size = item.stat().st_size item.unlink() except OSError as exc: errors.append(f"{item.name}: {exc}") continue removed += 1 freed += size return PruneResult("backup", removed, freed, tuple(errors)) ``` ```python # src/dmf_crawler/storage/repo.py (보존 정리 — 스냅샷 멤버십) import sqlite3 from datetime import date def prune_snapshots(conn: sqlite3.Connection, retain_days: int, today: str | None = None) -> int: """오래된 스냅샷 멤버십 행을 지운다. retain_days <= 0 이면 영구 보존. events + records + record_versions 가 남아 있으므로 이력은 손실되지 않는다. 사라지는 것은 '그날 어떤 레코드가 있었는가' 의 빠른 조회 인덱스뿐이다. """ if retain_days <= 0: return 0 cutoff = today or date.today().isoformat() rows = conn.execute( "SELECT run_seq FROM runs" " WHERE run_date < date(?, '-' || ? || ' days')" " AND run_seq IN (SELECT DISTINCT run_seq FROM snapshots)" " ORDER BY run_seq", (cutoff, int(retain_days)), ).fetchall() if not rows: return 0 # 항상 가장 최근 성공 실행 하나는 남긴다(diff 기준선 보호). keep = conn.execute( "SELECT run_seq FROM runs WHERE status='SUCCESS' ORDER BY run_seq DESC LIMIT 1" ).fetchone() keep_seq = int(keep["run_seq"]) if keep else -1 deleted = 0 for row in rows: run_seq = int(row["run_seq"]) if run_seq == keep_seq: continue cur = conn.execute("DELETE FROM snapshots WHERE run_seq = ?", (run_seq,)) deleted += cur.rowcount return deleted ``` ### 8.3 정리 실행 규칙 5가지 1. **`finalize` 스테이지에서만** 실행한다. 수집·diff·리포트가 전부 끝난 뒤다. 2. **트랜잭션 밖에서** 실행한다. 파일 삭제는 롤백되지 않으므로 DB 트랜잭션과 섞으면 안 된다. 3. **실패해도 무시한다.** 사용자가 xlsx 를 열어두면 그 파일은 잠겨 있다. `errors` 에 담아 로그에만 남기고 다음 실행에 다시 시도한다. 4. **`snapshots` 정리 뒤에만 `VACUUM` 을 고려한다.** `VACUUM` 은 DB 크기만큼의 임시 공간을 쓰므로 `backup.min_free_gb` 확인 후에만 실행하고, 실패해도 무시한다. 5. **가장 최근 성공 실행의 스냅샷은 절대 지우지 않는다.** 그것이 내일 diff 의 기준선이다. 코드에 `keep_seq` 가드로 박아 뒀다. --- ## 부록. 미해결 / 실측 필요 serviceKey 발급 직후 또는 첫 정상 실행 직후에 확인하고, 이 문서를 갱신한다. **A. 원본 데이터 실측 (serviceKey 발급 즉시)** - [ ] **전체 등록 건수(`totalCount`) 실측값** — 이 문서는 9,840건을 가정했다. §7 지표 1의 범위(6,000~30,000)를 실측 기준으로 재설정할 것 - [ ] **`numOfRows` 최대값** — 명세 크기가 3자리다. 999 가 실제로 통하는지 확인. 통하면 `source.page_size` 를 999로 올려 호출 수를 1/10로 줄인다 - [ ] **중복 `DMF_PERMIT_NO` 가 실재하는가** — §4.5 그룹 매칭 로직의 존재 이유다. 중복이 0건이면 무해한 no-op 으로 남는다 - [ ] **`DMF_PERMIT_NO` 가 빈 값인 레코드가 있는가** — 있으면 `SYN-` 합성키 경로가 실제로 동작한다. 건수를 기록할 것 - [ ] **7필드 각각의 실제 널 비율** — 게이트 4 임계 0.01 이 현실적인지 검증 - [ ] **`MANUF_COUNTRY_CODE_NM` 의 실제 구분자** — 콤마만인지, 슬래시·중점도 쓰이는지. `_COUNTRY_SPLIT_RE` 보강 근거 - [ ] **`MNFCTR_NAME` 다중값의 실제 구분자** — ` ,`(공백+콤마) 가정이 맞는지. 아니면 `_split_multi()` 를 고쳐야 한다 - [ ] **`DMF_PERMIT_DATE` 의 실제 표기** — `2015-02-26` 외 `20150226`·`2015.02.26` 이 섞이는지 **B. 등록번호 체계** - [ ] **신물질 포맷 전수 조사** — `수6580-16-ND(20)` 외에 `제`·`허` 등 다른 접두어가 있는가. `RE_PERMIT_NEW_SUBSTANCE` 커버리지 확인 - [ ] **관측되는 알파벳군 집합** — A~K 전부 등장하는가. `GROUP_TABLE` 에 없는 군(`L` 등)이 있는가 - [ ] **`sub` 토큰이 J군에서 정말 없는가** — FAQ 는 "J는 제외" 라고 했다. 반례가 있으면 파서 주석을 고칠 것 - [ ] **허여서 괄호가 3자리(`(10)` 이상)로 가는가** — 현재 정규식은 `\d{1,3}` 로 여유를 뒀다 - [ ] **[별표1] 성분 목록 전문(1~211호) 확보** — `ingr_no` 를 성분명으로 역매핑하려면 필수. 현행 고시 별표1 파싱 **C. 변경 탐지 임계값 (2~4주 관측 후 확정)** - [ ] **일일 실제 변동 규모** — 신규/변경/취하 각각의 중앙값·최댓값. `integrity.max_churn_ratio = 0.10` 이 너무 빡빡하거나 너무 느슨한지 판정 - [ ] **`COSMETIC` 이벤트의 실제 발생 빈도** — 원본 표기 정정이 얼마나 잦은가. 잦으면 리포트 기본 제외가 정당하고, 0에 가까우면 `identity_hash` 계층의 비용 대비 효용을 재검토 - [ ] **소스 갱신 주기** — `payload_sha256` 변화 간격의 실측 분포. 06:00 스케줄이 적절한지, 갱신이 주 1회면 매일 호출이 낭비인지 - [ ] **`WITHDRAWN` 이 실제로 발생하는가** — API 가 취하 건을 목록에서 빼는지, 아니면 계속 유지하는지. **후자라면 취하 탐지가 원리적으로 불가능**하고 이 문서 §5.4 전체의 전제가 바뀐다. **가장 중요한 미검증 항목** **D. 스키마·저장** - [ ] **`snapshots` 를 `WITHOUT ROWID` 로 둔 것의 실제 크기 이득** — 1만 행 × 30일 적재 후 `dbstat` 로 측정 - [ ] **1년 운영 후 DB 실제 크기** — §3.2 산정(연 128MB)의 검증. 2GB 임계 경고가 언제 걸리는지 - [ ] **`VACUUM INTO` 백업 1회 소요 시간** — 30분 실행 제한(`schedule.execution_time_limit_minutes`) 안에 드는지 - [ ] **부분 유니크 인덱스 `ux_runs_success_per_day` 가 backfill 과 충돌하는가** — 과거 날짜를 채우는 `backfill` 서브커맨드가 같은 날짜에 두 번 SUCCESS 를 넣으려 하면 막힌다. 의도된 동작인지 확인 **E. 도메인 규칙** - [ ] **염·수화물 접미 사전의 커버리지** — `_STRIPPABLE_SUFFIXES` 로 실제 성분명 전량을 처리해 `ingredient_base` 가 얼마나 잘 묶이는지 측정. 과다 제거(예: `벤조산`이 성분 본체인 경우) 사례 확인 - [ ] **`COUNTRY_ALIASES` 미등록 국가명 목록** — §7 지표 10 이 수집한다. 1주 관측 후 사전 보강 - [ ] **워치리스트 초기 목록** — 요구 00-REQUIREMENTS.md 의 미해결 항목. 사용자에게 관심 성분·업체를 받아야 `watchlist` 테이블이 의미를 갖는다