DMF_Crawler/docs/design/02-data-model.md
Yun Chan 56a6e2da93 chore: 저장소 구조 정리 및 문서화, 첫 커밋
- src/dist 산출물 분리 원칙 정리(.gitignore, .gitattributes)
- 루트 및 주요 폴더(config/scripts/prompts/tests/src, 런타임 폴더 5종)에
  안내용 README.md 추가
- CHANGELOG.md, LICENSE, docs/ops/05-release-and-versioning.md 추가
- docs/README.md 문서 지도 갱신
2026-09-04 09:25:44 +09:00

2894 lines
148 KiB
Markdown
Raw Permalink Blame History

This file contains invisible Unicode characters

This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# 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<accept_date>(?:19|20)\d{2}(?:0[1-9]|1[0-2])(?:0[1-9]|[12]\d|3[01]))"
r"-(?P<ingr_no>\d{1,4})"
r"-(?P<group>[A-Z])"
r"-(?P<serial>\d{1,6})"
r"(?:-(?P<sub>\d{1,3}))?"
r"(?:\((?P<grant>\d{1,3})\))?$"
)
# --------------------------------------------------------------------------
# 포맷 B: 신물질(신약 원료) 계열 — 예: 수6580-16-ND(20)
# 토큰 의미는 미검증이므로 구조 인식만 하고 값 해석은 하지 않는다.
# --------------------------------------------------------------------------
RE_PERMIT_NEW_SUBSTANCE = re.compile(
r"^(?P<prefix>[가-힣]{1,2})(?P<doc_no>\d{3,7})"
r"-(?P<year_seq>\d{1,3})"
r"-(?P<kind>ND)"
r"(?:\((?P<seq>\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<num>\d{4})_(?P<name>[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).
**흔들림 유형별 처리표.**
| 흔들림 유형 | 실례 | 처리 단계 | 처리 방법 |
|---|---|---|---|
| 전각/반각 | `` vs `SICOR` | 1단 | `unicodedata.normalize("NFKC")` |
| 전각 하이픈·괄호 | `20121228168` / `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<role>[^\]]{1,40})\]\s*")
_DATE_RE = re.compile(r"(?P<y>(?:19|20)\d{2})\D?(?P<m>\d{1,2})\D?(?P<d>\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-<sha1 앞 12자> (합성키. 리포트에서 식별 가능해야 한다)
"""
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 "<22>" 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<stamp>\d{8}_\d{6})$")
DATE_DIR_RE = re.compile(r"^(?P<d>\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` 테이블이 의미를 갖는다