- 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 문서 지도 갱신
2894 lines
148 KiB
Markdown
2894 lines
148 KiB
Markdown
# 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).
|
||
|
||
**흔들림 유형별 처리표.**
|
||
|
||
| 흔들림 유형 | 실례 | 처리 단계 | 처리 방법 |
|
||
|---|---|---|---|
|
||
| 전각/반각 | `SICOR` vs `SICOR` | 1단 | `unicodedata.normalize("NFKC")` |
|
||
| 전각 하이픈·괄호 | `20121228-168` / `(1)` | 1단 | `_PUNCT_FOLD` 치환표 |
|
||
| 비단절 공백·제로폭 | U+00A0, U+200B, U+FEFF | 1단 | `_PUNCT_FOLD` + `Cc` 카테고리 제거 |
|
||
| 연속 공백·앞뒤 공백 | `(주) 대웅제약 ` | 1단 | `\s+` → 한 칸, `strip()` |
|
||
| 법인격 표기 | `㈜하이플` / `(주)하이플` / `주식회사 하이플` | 1단(표시 통일) + 2단(어휘 제거) | `㈜`→`(주)` 통일 후, 키에서는 법인격 어휘 삭제 |
|
||
| 원본 오탈자 공백 | `CORTICOSTER OIDO` | 2단 | 공백 전부 삭제 → 정정본과 동일 키 |
|
||
| 대소문자 | `S.R.L.` vs `s.r.l.` | 2단 | `lower()` |
|
||
| 구두점 | `Co., Ltd.` vs `Co.,Ltd` | 2단 | 영숫자·한글 외 전부 삭제 |
|
||
| 국가 다중값 순서 | `이탈리아,스위스` vs `스위스,이탈리아` | 1단 | 분해 후 **정렬** |
|
||
| 국가 별칭 | `미국` / `미합중국` / `USA` | 1단 | `COUNTRY_ALIASES` 치환표 |
|
||
| 날짜 표기 | `2015-02-26` / `20150226` / `2015.02.26` | 1단 | `normalize_date()` → ISO |
|
||
| 다중 제조소 구분자 | `A , B` vs `A, B` | 1단 | ` ,` 우선 분해 후 ` , ` 로 재조립 |
|
||
|
||
### 4.3 정규화 코드 전문
|
||
|
||
```python
|
||
# src/dmf_crawler/normalize.py (3/3 — 필드 정규화 · 레코드 조립)
|
||
|
||
# --------------------------------------------------------------------------
|
||
# 법인격 어휘 — 매칭키에서 제거한다. 긴 것부터 지워야 부분 삭제가 안 생긴다.
|
||
# --------------------------------------------------------------------------
|
||
LEGAL_FORMS_KO: tuple[str, ...] = (
|
||
"주식회사", "유한책임회사", "유한회사", "합자회사", "합명회사",
|
||
"재단법인", "사단법인", "의료법인", "학교법인", "(주)", "(유)", "(재)", "(사)",
|
||
)
|
||
LEGAL_FORMS_EN: tuple[str, ...] = (
|
||
"coltd", "companylimited", "limited", "ltd", "llc", "llp", "inc", "incorporated",
|
||
"corporation", "corp", "gmbh", "ag", "sa", "srl", "spa", "bv", "nv", "as",
|
||
"pvtltd", "pvt", "plc", "kg", "oy", "ab", "sas", "sl", "pte", "sdnbhd",
|
||
)
|
||
|
||
# 표시값 단계의 법인격 표기 통일
|
||
_LEGAL_DISPLAY_FOLD = {
|
||
"㈜": "(주)", "㈔": "(사)", "㈖": "(재)", "㈕": "(유)",
|
||
}
|
||
|
||
# --------------------------------------------------------------------------
|
||
# 성분명 접두/접미 — ingredient_base(집계 전용) 산출에만 쓴다.
|
||
# 동일성 판정에는 절대 쓰지 않는다(염 형태가 다르면 다른 등록이다).
|
||
# --------------------------------------------------------------------------
|
||
MICRONIZE_PREFIXES: tuple[str, ...] = ("초미분화", "미분화")
|
||
|
||
SALT_SUFFIXES: tuple[str, ...] = (
|
||
"브롬화수소산염", "메탄술폰산염", "메탄설폰산염", "메실산염", "베실산염",
|
||
"토실산염", "에실산염", "이세티온산염", "파모산염", "팜산염",
|
||
"푸마르산염", "말레산염", "말산염", "타르타르산염", "주석산염",
|
||
"시트르산염", "구연산염", "숙신산염", "아세트산염", "초산염",
|
||
"락트산염", "젖산염", "글루콘산염", "글루쿠론산염", "아스파르트산염",
|
||
"팔미트산염", "스테아르산염", "벤조산염", "살리실산염", "옥살산염",
|
||
"인산염", "황산염", "질산염", "염산염", "브롬산염",
|
||
"요오드화물", "브롬화물", "염화물",
|
||
)
|
||
HYDRATE_SUFFIXES: tuple[str, ...] = (
|
||
"일수화물", "이수화물", "삼수화물", "사수화물", "오수화물",
|
||
"육수화물", "칠수화물", "팔수화물", "반수화물", "수화물", "무수물",
|
||
)
|
||
CATION_SUFFIXES: tuple[str, ...] = (
|
||
"나트륨", "칼륨", "칼슘", "마그네슘", "아연", "리튬", "암모늄",
|
||
"메글루민", "트로메타민", "디에탄올아민", "에탄올아민", "베타덱스",
|
||
)
|
||
_STRIPPABLE_SUFFIXES: tuple[str, ...] = tuple(
|
||
sorted(SALT_SUFFIXES + HYDRATE_SUFFIXES + CATION_SUFFIXES, key=len, reverse=True)
|
||
)
|
||
|
||
# --------------------------------------------------------------------------
|
||
# 국가명 별칭 — 표시값 단계에서 표준명으로 접는다.
|
||
# --------------------------------------------------------------------------
|
||
COUNTRY_ALIASES: Mapping[str, str] = {
|
||
"대한민국": "한국", "korea": "한국", "republicofkorea": "한국", "kr": "한국",
|
||
"미합중국": "미국", "usa": "미국", "us": "미국", "unitedstates": "미국",
|
||
"중화인민공화국": "중국", "china": "중국", "cn": "중국",
|
||
"인도": "인도", "india": "인도", "in": "인도",
|
||
"일본": "일본", "japan": "일본", "jp": "일본",
|
||
"이탈리아": "이탈리아", "italy": "이탈리아", "it": "이탈리아",
|
||
"스위스": "스위스", "switzerland": "스위스", "ch": "스위스",
|
||
"독일": "독일", "germany": "독일", "de": "독일",
|
||
"스페인": "스페인", "spain": "스페인", "es": "스페인",
|
||
"프랑스": "프랑스", "france": "프랑스", "fr": "프랑스",
|
||
"영국": "영국", "unitedkingdom": "영국", "uk": "영국", "gb": "영국",
|
||
"아일랜드": "아일랜드", "ireland": "아일랜드", "ie": "아일랜드",
|
||
"이스라엘": "이스라엘", "israel": "이스라엘", "il": "이스라엘",
|
||
"대만": "대만", "taiwan": "대만", "tw": "대만",
|
||
"슬로베니아": "슬로베니아", "헝가리": "헝가리", "폴란드": "폴란드",
|
||
"오스트리아": "오스트리아", "네덜란드": "네덜란드", "벨기에": "벨기에",
|
||
"덴마크": "덴마크", "스웨덴": "스웨덴", "핀란드": "핀란드", "노르웨이": "노르웨이",
|
||
"포르투갈": "포르투갈", "그리스": "그리스", "루마니아": "루마니아",
|
||
"체코": "체코", "슬로바키아": "슬로바키아", "크로아티아": "크로아티아",
|
||
"캐나다": "캐나다", "canada": "캐나다", "멕시코": "멕시코", "브라질": "브라질",
|
||
"아르헨티나": "아르헨티나", "호주": "호주", "australia": "호주",
|
||
"뉴질랜드": "뉴질랜드", "싱가포르": "싱가포르", "말레이시아": "말레이시아",
|
||
"인도네시아": "인도네시아", "베트남": "베트남", "태국": "태국",
|
||
"터키": "튀르키예", "튀르키예": "튀르키예", "turkey": "튀르키예",
|
||
}
|
||
|
||
_COUNTRY_SPLIT_RE = re.compile(r"[,/·;|]+")
|
||
_SITE_ROLE_RE = re.compile(r"^\[(?P<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` 테이블이 의미를 갖는다
|