diff --git a/ER_DOSE_ERROR.md b/ER_DOSE_ERROR.md index 64ec37c..58fb076 100644 --- a/ER_DOSE_ERROR.md +++ b/ER_DOSE_ERROR.md @@ -8,7 +8,7 @@ mbeat.er_data_raw -> er_dose batch -> prism_common.er_dose_raw_parsed - -> prism_common.de_trend_die_yield_daily (summary) + -> Airflow 후속 작업: prism_common.de_trend_die_yield_daily (summary) ``` Root cause는 이 배치와 별도 흐름이다. @@ -17,7 +17,7 @@ Root cause는 이 배치와 별도 흐름이다. mbeat.er_data_raw_euv -> contents root cause 파싱 -> prism_common.er_dose_euv_parsed - -> prism_common.de_trend_root_cause_daily (summary) + -> Airflow 후속 작업: prism_common.de_trend_root_cause_daily (summary) ``` `prism_common.er_dose_raw_parsed`와 `prism_common.er_dose_euv_parsed`는 서로 조인하거나 매칭하지 않는다. @@ -30,6 +30,7 @@ mbeat.er_data_raw_euv - `mbeat.er_data_raw_euv`: Root cause source description 후보 RAW. `contents`에 `dose error detected in file`, `root cause`, `exposure id`, 각종 EUV 지표가 들어온다. `er_date`, `er_index`가 없다. - `prism_common.er_dose_euv_parsed`: FE 조회용 root cause 결과 테이블. `er_data_raw_euv.contents`를 파싱한 구조화 컬럼과 원문을 저장하며, `er_dose_raw_parsed`와 무관하다. - `prism_common.de_trend_root_cause_daily`: EUV Root Cause 일별 발생 빈도 요약 서머리 테이블 (`occur_date`, `eq_name`, `root_cause`, `frequency`). +- `mbeat.batch_event_log`: 여러 배치가 공통으로 사용하는 이벤트 로그 테이블. `batch_name`, `target_date`, `event_type`, `message`와 가변 상세 데이터인 `data jsonb`를 저장한다. DDL: @@ -184,22 +185,23 @@ Mermaid ERD는 렌더링 호환성을 위해 타입 표기를 단순화했다. - RAW는 이전 chunk의 `COPY`를 적재 worker 1개에서 실행하는 동안 다음 chunk를 조회·파싱한다. 동시에 대기하는 적재 작업은 1개로 제한해 처리 순서와 메모리 사용량을 유지한다. 4. 각 `chunk`를 `prism_common.er_dose_raw_parsed` 일별 파티션에 `COPY` append insert - 파티션 적재는 공통 `copy_insert_to_partition_table`을 사용하며, DataFrame 컬럼을 테이블 물리 컬럼 순서와 동일하게 정렬한 뒤 `COPY ... FROM STDIN WITH CSV HEADER`를 실행한다. -5. 적재된 파티션 날짜를 기준으로 DIE Yield 서머리 테이블(`prism_common.de_trend_die_yield_daily`) 및 EUV Root Cause 서머리 테이블(`prism_common.de_trend_root_cause_daily`)에 `UPSERT` 집계 업데이트 실행 +5. 적재한 파티션을 ANALYZE하고 파서 작업 종료 -환경변수 기반 기본 실행에서 target date와 `ER_DOSE_START_TIME`, `ER_DOSE_END_TIME`가 모두 없으면 raw/euv 배치는 최근 2일 lookback 모드로 동작한다. +최근 날짜별 판단은 Airflow DAG에서 수행한다. -1. 실행일 기준 `오늘 포함 최근 2일`을 날짜 오름차순으로 순회 -2. 각 날짜에 대해 원천 raw 전체 건수와 타겟 parsed 전체 건수를 비교 -3. 건수가 같으면 해당 날짜는 스킵 -4. RAW/EUV 건수가 다르고 parsed 건수가 0보다 크면 원천의 `eq_name`, `code`, `code_occur_time` 기준 중복 제거 건수를 추가 계산 -5. 중복 제거 건수가 parsed 건수와 같으면 스킵 -6. 중복 제거 건수도 다르거나 parsed 건수가 0이면 해당 날짜의 parsed 파티션을 `TRUNCATE` -7. 원천 raw를 해당 날짜 처음부터 다시 조회해 chunk 단위로 파싱 후 insert -8. 적재 완료 후 해당 날짜의 서머리 테이블 2종을 `UPSERT` 업데이트 +1. 기준일 포함 최근 10일을 날짜 오름차순으로 확인 +2. 각 날짜의 원천 raw 전체 건수와 parsed 전체 건수를 비교 +3. 건수가 같으면 해당 파서의 처리 대상에서 제외 +4. 전체 건수가 다르고 parsed 건수가 0보다 크면 원천의 `eq_name`, `code`, `code_occur_time` 기준 DISTINCT 건수를 추가 비교 +5. DISTINCT 건수도 다르거나 parsed 건수가 0이면 재적재 대상으로 확정 +6. 전체 날짜 계획을 XCom에 저장한 후, 날짜별 EUV → RAW 파서 작업 실행 +7. 대상 파서는 지정된 날짜 파티션을 TRUNCATE하고 재적재. 재시도에도 계획을 유지하며 건수를 다시 비교하지 않음 +8. 해당 날짜 적재 성공 후 서머리 두 종류와 RAW/EUV 통계 두 종류를 독립 Airflow 작업으로 실행 +9. 실패한 서머리·통계는 같은 기간으로 해당 작업만 재시도 -`ER_DOSE_EUV_TARGET_DATE`가 있으면 해당 날짜 1일만 같은 방식으로 count 비교 후 필요 시 재적재한다. `ER_DOSE_START_TIME`/`ER_DOSE_END_TIME`으로 시간 범위를 직접 지정하면 count 비교 없이 해당 범위를 처리한다. +환경변수/CLI 파서 실행에는 날짜 또는 시간 범위를 반드시 지정한다. target date가 있으면 해당 날짜만 재적재하며, `ER_DOSE_START_TIME`/`ER_DOSE_END_TIME`으로 시간 범위를 직접 지정하면 기존 방식으로 해당 범위를 처리한다. 파서 단독 실행은 서머리·통계를 갱신하지 않는다. -`ER_DOSE_EUV` 배치는 `mbeat.er_data_raw_euv`를 기간 조건으로 `chunk` 조회하고, root cause 형식의 `contents`만 파싱해 `prism_common.er_dose_euv_parsed`에 적재한다. EUV source count도 parsed count와 맞추기 위해 `contents`에 `dose error detected in file:`과 `root cause`가 있는 row만 계산하며, 전체 건수가 다를 때는 RAW와 동일하게 고유 이벤트 count를 추가 비교한다. EUV parsed 결과에는 `eq_name`, `er_type`, `code`, `code_occur_time`, `title`, `contents`, `reason_code`, `task`, `compile_script`와 root cause 파싱 컬럼만 저장하며, 적재 완료 후 `de_trend_root_cause_daily` 서머리 테이블을 `UPSERT` 업데이트한다. +`ER_DOSE_EUV` 배치는 `mbeat.er_data_raw_euv`를 기간 조건으로 `chunk` 조회하고, root cause 형식의 `contents`만 파싱해 `prism_common.er_dose_euv_parsed`에 적재한다. EUV source count도 parsed count와 맞추기 위해 `contents`에 `dose error detected in file:`과 `root cause`가 있는 row만 계산하며, 전체 건수가 다를 때는 RAW와 동일하게 고유 이벤트 count를 추가 비교한다. EUV parsed 결과에는 `eq_name`, `er_type`, `code`, `code_occur_time`, `title`, `contents`, `reason_code`, `task`, `compile_script`와 root cause 파싱 컬럼만 저장하며, Root Cause 집계와 건수 로그는 이후 RAW 마지막 단계에서 저장한다. RAW와 EUV 모두 대용량 처리를 위해 전체 결과를 한 번에 메모리로 올리지 않고 `read chunk -> parse -> insert` 방식으로 반복 처리한다. 또한, 데이터베이스 드라이버 단의 메모리 팽창을 방지하기 위해 SQLAlchemy 서버사이드 커서(`stream_results=True`, `max_row_buffer=chunk_size`)를 활성화하여 스트리밍 조회를 수행한다. 다만 실제 메모리 사용량은 `chunk` 크기와 raw `contents` 크기에 영향을 받기 때문에 운영 환경에서 조정이 필요할 수 있다. RAW와 EUV 모두 조회 SQL에서 `prism_dev.photo_eqp_info`의 `use_yn = 'Y'`이고 `eqp_model_name like 'NXE%'`인 `eqp_id`를 서브쿼리로 조회해 `eq_name` 필터로 사용한다. RAW의 이전 `lot_seq`, `wafer_seq` 상태 조회에도 같은 조건을 적용한다. @@ -258,10 +260,12 @@ RAW parsed 저장 필드: ## LO-0050/LO-0051/LO-0052 파싱 규칙 -- `LO-0050`의 lot 시작, `LO-0051`의 lot 종료, `LO-0052`의 lot 중단 메시지에 동일한 파싱 규칙을 적용한다. -- `lot_id`: 원문 텍스트의 `lot '([^']+)'` 정규식 패턴에서 추출한다. +- `LO-0050`의 lot 시작, `LO-0051`의 lot 종료, `LO-0052`의 lot 중단 메시지에서 동일한 lot 필드를 저장한다. +- `LO-0050/LO-0051`은 기존 `lot '...' (id=...)` 형식을 사용한다. +- `LO-0052`는 별도 코드 분기에서 `lot '...' (id=...)`와 `lot (name='...', id=...)` 형식을 모두 처리한다. +- `lot_id`: 각 메시지의 lot 문자열에서 추출한다. - `lot_name`: 추출된 `lot_id`에서 첫 번째 `.`(점) 문자를 기준으로 이전 텍스트를 추출(`lot_id.split('.', maxsplit=1)[0]`)한다. -- `lot_seq`: `(id=\s*\d+)` 정규식 패턴에서 우선 추출하며, 미매칭 시 기존 `_LOT_SEQ_PATTERNS` 패턴으로 폴백한다. +- `lot_seq`: 각 메시지의 `id` 값에서 추출하며, 미매칭 시 기존 `_LOT_SEQ_PATTERNS` 패턴으로 폴백한다. ## 실행 @@ -271,10 +275,11 @@ EUV 날짜 변수는 `ER_DOSE_EUV_TARGET_DATE` 를 사용한다. DB 접속은 `--dsn`, 프로젝트 루트 `er_dose.properties`, `ER_DOSE_DB_DSN`, `DATABASE_URL` 순서로 사용한다. 기본 `chunk` 크기는 `ER_DOSE_RAW` 및 `ER_DOSE_EUV` 배치 모두 `30000`이며 `--chunk-size`로 조정할 수 있다. -RAW 기본 실행은 최근 2일 lookback 모드이며, `--lookback-days` 또는 환경변수 기반 실행의 `ER_DOSE_LOOKBACK_DAYS`로 일수를 바꿀 수 있다. +최근 10일 대상 결정은 Airflow DAG에서 수행한다. CLI에서는 `--date` 또는 `--start-time/--end-time`을 반드시 지정한다. ```bash python -m er_dose.run_er_dose_batch \ + --date 2026-04-13 \ --parser ER_DOSE_RAW \ --chunk-size 30000 \ --dsn 'postgresql://user:password@host:5432/dbname' @@ -305,3 +310,17 @@ python -m er_dose.run_er_dose_batch \ --chunk-size 30000 \ --dsn 'postgresql://user:password@host:5432/dbname' ``` + +## ER Dose 날짜별 Airflow 실행과 재시도 + +`dags/er_dose_daily_dag.py`가 최근 10일의 건수를 비교해 날짜별 RAW/EUV 처리 여부를 결정합니다. 원천 건수와 parsed 건수가 같으면 제외하고, parsed 건수가 있을 때에는 기존과 동일하게 원천 DISTINCT 건수도 비교합니다. 계획은 `plan_dates` 작업의 XCom에 저장한 뒤 파싱을 시작합니다. 실행·건너뛰기는 DAG의 `BranchPythonOperator`에서 결정하며, 선택된 파서·서머리 함수는 별도 실행 여부 판단 없이 지정 기간을 처리합니다. + +파서는 `target_date` 또는 `start_time/end_time`을 명시적으로 받아 처리합니다. 파서 내부의 lookback과 서머리·통계 호출은 제거했습니다. `--date` 실행은 해당 날짜 파티션을 비우고 다시 적재하며, 시간 범위 직접 실행은 기존 적재 동작을 유지합니다. CLI 파서 실행만으로는 서머리·통계를 갱신하지 않습니다. + +DAG는 과거 날짜부터 EUV → RAW 파서를 SSH로 Linux 서버에서 실행합니다. 해당 날짜 적재 성공 후 수율 서머리·원인 서머리·RAW 통계·EUV 통계는 Airflow worker가 기존 repository로 직접 조회·저장합니다. SSH Connection과 실행 명령 Variable 설정은 운영 안내를 참고하세요. 서머리 실패는 다른 통계나 다음 날짜 파싱의 선행 조건이 아닙니다. 실제 병렬 실행 수는 Airflow executor/pool 설정에 따릅니다. + +서머리 작업은 기존 DELETE 메서드와 INSERT 메서드를 직접 호출하며, 두 실행을 하나의 트랜잭션으로 묶지 않습니다. 0건 집계도 DELETE 후 처리합니다. 통계는 SELECT 후 동일 배치·날짜·이벤트·시간 범위 메시지에 해당하는 이전 스냅샷을 DELETE하고 INSERT합니다. 재시도 시 중복 행이 남지 않도록 하며, 같은 기간의 과거 스냅샷은 교체됩니다. + +RAW와 EUV 로그는 별도 행이며, `data`에는 `equipment_counts` 배열만 저장합니다. 각 원소는 `eq_name`, `source_count`, `target_count`입니다. 시간 범위는 `message`, 시작일은 `target_date`, 저장 시각은 `created_at`에 기록합니다. 메시지와 JSON 구성은 `airflow_modules/er_dose_jobs.py`에 있고, 공통 로그 repository는 전달받은 조건과 데이터만 처리합니다. + +DAG 배포·운영 및 장애 복구 절차는 [Airflow 재처리 안내](docs/er_dose_airflow.md)를 참고하세요. 이전 스레드 구성의 측정치는 [과거 검증 기록](docs/pr/er_dose_final_statistics_validation.md)에 있으며 새 DAG의 성능 측정치는 아닙니다. diff --git a/README.md b/README.md index 4674cf4..d44e20d 100644 --- a/README.md +++ b/README.md @@ -37,15 +37,17 @@ RUBI 텍스트와 RUIP 이미지를 수집 및 매칭하여 reticle backside 오 ## ER Dose Error 배치 -`ER_DOSE_RAW` 배치는 `mbeat.er_data_raw`의 dose warning 로그를 파싱해 `prism_common.er_dose_raw_parsed`에 적재하며, 수집 완료 후 `prism_common.de_trend_die_yield_daily` 일별 DIE Yield 서머리 테이블을 갱신합니다. parsed 테이블에는 `eq_name`, `code`, `code_occur_time`, `title`, `contents`와 `contents`에서 실제로 필요한 `exposure_handle`, `action_handle`, `lot_id`, `lot_name`, `lot_seq`, `wafer_seq`, `de_err`, `n_slit`, `use_yn`을 저장합니다. 조회 대상 `code`는 `DW-3411`, `DW-3425`, `DW-343A`, `DW-343B`, `LO-0050`, `LO-0051`, `LO-0052`, `LO-0061`, `LO-8166`, `LO-8167`, `KE-9103`, `KE-9104`이며, 코드 값은 DB에 저장된 원본 형식 그대로 비교합니다. +`ER_DOSE_RAW` 배치는 `mbeat.er_data_raw`의 dose warning 로그를 파싱해 `prism_common.er_dose_raw_parsed`에 적재하며, Airflow의 별도 후속 작업이 `prism_common.de_trend_die_yield_daily` 일별 DIE Yield 서머리를 갱신합니다. parsed 테이블에는 `eq_name`, `code`, `code_occur_time`, `title`, `contents`와 `contents`에서 실제로 필요한 `exposure_handle`, `action_handle`, `lot_id`, `lot_name`, `lot_seq`, `wafer_seq`, `de_err`, `n_slit`, `use_yn`을 저장합니다. 조회 대상 `code`는 `DW-3411`, `DW-3425`, `DW-343A`, `DW-343B`, `LO-0050`, `LO-0051`, `LO-0052`, `LO-0061`, `LO-8166`, `LO-8167`, `KE-9103`, `KE-9104`이며, 코드 값은 DB에 저장된 원본 형식 그대로 비교합니다. 배치는 `code_occur_time` 기간 조건으로 조회한 후보를 한 번에 메모리로 올리지 않고, `chunk` 단위로 읽어서 파싱 후 바로 `COPY` 적재합니다. 파티션 적재는 공통 `copy_insert_to_partition_table`을 사용하며, 적재 전에 DataFrame 컬럼을 테이블의 물리 컬럼 순서와 동일하게 정렬한 뒤 `COPY ... FROM STDIN WITH CSV HEADER`를 실행합니다. 현재 기본 `chunk` 크기는 `ER_DOSE_RAW` 및 `ER_DOSE_EUV` 배치 모두 `30000`이며 실행 시 조정할 수 있습니다. 조회는 SQLAlchemy 서버사이드 커서(`stream_results=True`, `max_row_buffer=chunk_size`) 기반 스트리밍으로 수행되지만, 실제 메모리 사용량은 `chunk` 크기와 raw `contents` 크기에 영향을 받으므로 운영 환경에 맞게 조정해야 합니다. 청크 단위로 처리되더라도 설비(`eq_name`)별로 이전에 파싱한 `lot_seq`와 `wafer_seq`를 기억하여 지속 적용합니다. `ER_DOSE_RAW`은 적재 worker 1개를 사용해 이전 청크의 `COPY`와 다음 청크의 조회·파싱을 겹쳐 실행하며, 대기 중인 DataFrame은 최대 1개로 제한합니다. 이 변경의 운영 실측 결과는 [DB 스트리밍 및 RAW 성능 개선 문서](docs/db_streaming_optimization.md#5-er-dose-raw-파싱적재-파이프라인-실측)에 기록합니다. -`ER_DOSE_RAW`와 `ER_DOSE_EUV`의 processor 기본 실행은 최근 2일 lookback 모드입니다. 실행일 기준 `오늘 포함 최근 2일`을 날짜별로 검사하고, 먼저 원천 raw와 parsed의 전체 건수를 비교합니다. 두 배치 모두 전체 건수가 다르고 parsed 건수가 0보다 클 때만 원천에서 `(eq_name, code, code_occur_time)`이 같은 행을 중복 제거한 건수를 한 번 더 계산합니다. 이 건수가 parsed 건수와 같으면 중복으로 인한 차이이므로 스킵하고, 여전히 다르거나 parsed 건수가 0이면 해당 날짜 parsed 파티션을 `TRUNCATE`한 뒤 원천 raw를 처음부터 다시 파싱해 적재합니다. `ER_DOSE_EUV_TARGET_DATE`를 명시하면 해당 날짜 1일만 같은 방식으로 검사하고, `ER_DOSE_START_TIME`/`ER_DOSE_END_TIME`를 명시하면 count 비교 없이 지정한 시간 범위를 처리합니다. EUV source count는 root cause 파싱 대상인 `contents`만 세어 parsed count와 비교합니다. +최근 날짜 판단은 `er_dose_daily` DAG가 담당합니다. 기준일 포함 최근 10일을 검사하며, 원천/parsed 전체 건수가 다르고 parsed 건수가 있을 때 원천 DISTINCT 건수를 추가 비교하는 기존 규칙을 유지합니다. 확정된 날짜별 처리 계획을 XCom에 저장한 후 파서는 지정된 날짜만 재적재합니다. EUV source count는 root cause 파싱 대상인 `contents`만 셉니다. 파서 단독 실행에는 날짜 또는 시간 범위가 필요합니다. + +날짜별 적재 완료 후 Airflow 통계 작업에서 RAW/EUV 통계를 각각 조회해 `mbeat.batch_event_log`에 배치별 한 행씩 저장합니다. 각 행의 `data.equipment_counts` 배열에 설비별 원천/parsed 건수를 담습니다. 로그 테이블이 없으면 [create_batch_event_log.sql](er_dose/sql/create_batch_event_log.sql)을 먼저 적용해야 합니다. DW 로그에서 `exposure_handle`이 같은 설비의 이전 값보다 `1000` 이상 커지면 테스트샷성 row로 보고 저장은 하되 `use_yn='N'`으로 표시합니다. 일반 분석에서는 `use_yn='Y'` 조건을 사용하면 되고, row 자체는 저장되므로 raw count와 parsed count 비교가 계속 어긋나는 문제를 피할 수 있습니다. -`ER_DOSE_EUV`는 `mbeat.er_data_raw_euv` 기반 root cause 결과용 실행입니다. 대상 결과는 `prism_common.er_dose_euv_parsed`에 저장하며, `er_line`, `belong`, `type`은 저장하지 않습니다. `contents`에서 `dose_error_detected_in_file`, `exposure_id`, `time`, `root_cause`와 각종 EUV metric 컬럼을 파싱해 적재합니다. 컬럼명은 소문자 snake_case 기준으로 공백, `.`, `-`, `<`, `=`를 `_`로 치환하며, 파생 컬럼은 `root_cause_code`만 저장합니다. 파싱 및 적재가 완료되면 `prism_common.de_trend_root_cause_daily` 서머리 테이블에 일별/설비별/원인별 발생 빈도(`frequency`)를 자동 `UPSERT` 합니다. +`ER_DOSE_EUV`는 `mbeat.er_data_raw_euv` 기반 root cause 결과용 실행입니다. 대상 결과는 `prism_common.er_dose_euv_parsed`에 저장하며, `er_line`, `belong`, `type`은 저장하지 않습니다. `contents`에서 `dose_error_detected_in_file`, `exposure_id`, `time`, `root_cause`와 각종 EUV metric 컬럼을 파싱해 적재합니다. 컬럼명은 소문자 snake_case 기준으로 공백, `.`, `-`, `<`, `=`를 `_`로 치환하며, 파생 컬럼은 `root_cause_code`만 저장합니다. EUV 적재 결과의 일별/설비별/원인별 발생 빈도(`frequency`)는 이후 RAW 마지막 단계에서 DELETE 후 INSERT합니다. 상세 스키마와 파싱 규칙은 [ER_DOSE_ERROR.md](ER_DOSE_ERROR.md)를 기준으로 관리합니다. @@ -219,7 +221,6 @@ ER_DOSE_DB_DSN=postgresql://user:password@host:5432/dbname - `ER_DOSE_RAW_TARGET_DATE`: ER Dose raw 대상 날짜, `YYYY-MM-DD` - `ER_DOSE_EUV_TARGET_DATE`: ER Dose EUV 대상 날짜, `YYYY-MM-DD` - `ER_DOSE_CHUNK_SIZE`: ER Dose raw/euv fetch chunk 크기 -- `ER_DOSE_LOOKBACK_DAYS`: `ER_DOSE_RAW` 기본 lookback 일수, 기본값 `2` - `INPUT_DATE`: `RBI_INPUT_DATE` 대체값 - `ER_DOSE_TARGET_DATE`: raw 레거시 대상 날짜 이름 - `TARGET_DATE`: `ER_DOSE_TARGET_DATE` 레거시 대체값 @@ -244,7 +245,7 @@ pip3 install --target .vendor SQLAlchemy psycopg2-binary - 내부적으로 `ER_DOSE_RAW`는 `ER_DOSE_RAW_TARGET_DATE`, `ER_DOSE_EUV`는 `ER_DOSE_EUV_TARGET_DATE`를 사용합니다. - raw/euv processor 모두 대상 날짜 기준으로 하루 범위를 계산합니다. - `BATCH_TARGET=ER_DOSE_RAW`가 현재 raw 배치의 기본 이름입니다. 레거시 `ER_DOSE`도 계속 지원합니다. -- `BATCH_TARGET=ER_DOSE_RAW` 또는 `BATCH_TARGET=ER_DOSE_EUV`를 환경변수만으로 실행하고 날짜 인자를 주지 않으면 최근 2일 lookback + 날짜별 count 비교 기반 재적재 전략이 적용됩니다. 두 배치 모두 전체 count가 다를 때만 고유 이벤트 count를 추가로 비교합니다. +- `BATCH_TARGET=ER_DOSE_RAW` 또는 `BATCH_TARGET=ER_DOSE_EUV`는 날짜 또는 시간 범위를 반드시 지정합니다. 최근 날짜 대상 결정은 Airflow DAG에서 수행합니다. - 인자 없이 `main.py`를 실행하면 기존과 동일하게 환경변수 기반 실행입니다. ### DB 초기화 @@ -312,13 +313,7 @@ ER_DOSE_DB_DSN='postgresql://user:password@host:5432/dbname' \ 직접 ER Dose 실행 스크립트를 사용할 수도 있습니다. -`ER_DOSE_RAW`를 날짜 없이 실행하면 최근 2일 lookback 모드로 동작합니다. - -```bash -python3 -m er_dose.run_er_dose_batch \ - --parser ER_DOSE_RAW \ - --dsn 'postgresql://user:password@host:5432/dbname' -``` +`--date` 또는 `--start-time/--end-time`을 지정합니다. CLI 파서는 적재만 수행합니다. ```bash python3 -m er_dose.run_er_dose_batch \ @@ -385,3 +380,17 @@ python3 -m er_dose.run_er_dose_batch \ - PostgreSQL 전환 자세한 매칭 규칙은 [/Users/parkjunho/PycharmProjects/PythonStudy/IMAGE_TEXT_MATCHING.md](/Users/parkjunho/PycharmProjects/PythonStudy/IMAGE_TEXT_MATCHING.md) 를 참고하면 됩니다. + +## ER Dose 날짜별 Airflow 실행과 재시도 + +`dags/er_dose_daily_dag.py`가 최근 10일의 건수를 비교해 날짜별 RAW/EUV 처리 여부를 결정합니다. 원천 건수와 parsed 건수가 같으면 제외하고, parsed 건수가 있을 때에는 기존과 동일하게 원천 DISTINCT 건수도 비교합니다. 계획은 `plan_dates` 작업의 XCom에 저장한 뒤 파싱을 시작합니다. 실행·건너뛰기는 DAG의 `BranchPythonOperator`에서 결정하며, 선택된 파서·서머리 함수는 별도 실행 여부 판단 없이 지정 기간을 처리합니다. + +파서는 `target_date` 또는 `start_time/end_time`을 명시적으로 받아 처리합니다. 파서 내부의 lookback과 서머리·통계 호출은 제거했습니다. `--date` 실행은 해당 날짜 파티션을 비우고 다시 적재하며, 시간 범위 직접 실행은 기존 적재 동작을 유지합니다. CLI 파서 실행만으로는 서머리·통계를 갱신하지 않습니다. + +DAG는 과거 날짜부터 EUV → RAW 파서를 SSH로 Linux 서버에서 실행합니다. 해당 날짜 적재 성공 후 수율 서머리·원인 서머리·RAW 통계·EUV 통계는 Airflow worker가 기존 repository로 직접 조회·저장합니다. SSH Connection과 실행 명령 Variable 설정은 운영 안내를 참고하세요. 서머리 실패는 다른 통계나 다음 날짜 파싱의 선행 조건이 아닙니다. 실제 병렬 실행 수는 Airflow executor/pool 설정에 따릅니다. + +서머리 작업은 기존 DELETE 메서드와 INSERT 메서드를 직접 호출하며, 두 실행을 하나의 트랜잭션으로 묶지 않습니다. 0건 집계도 DELETE 후 처리합니다. 통계는 SELECT 후 동일 배치·날짜·이벤트·시간 범위 메시지에 해당하는 이전 스냅샷을 DELETE하고 INSERT합니다. 재시도 시 중복 행이 남지 않도록 하며, 같은 기간의 과거 스냅샷은 교체됩니다. + +RAW와 EUV 로그는 별도 행이며, `data`에는 `equipment_counts` 배열만 저장합니다. 각 원소는 `eq_name`, `source_count`, `target_count`입니다. 시간 범위는 `message`, 시작일은 `target_date`, 저장 시각은 `created_at`에 기록합니다. 메시지와 JSON 구성은 `airflow_modules/er_dose_jobs.py`에 있고, 공통 로그 repository는 전달받은 조건과 데이터만 처리합니다. + +DAG 배포·운영 및 장애 복구 절차는 [Airflow 재처리 안내](docs/er_dose_airflow.md)를 참고하세요. 이전 스레드 구성의 측정치는 [과거 검증 기록](docs/pr/er_dose_final_statistics_validation.md)에 있으며 새 DAG의 성능 측정치는 아닙니다. diff --git a/airflow_modules/er_dose_jobs.py b/airflow_modules/er_dose_jobs.py new file mode 100644 index 0000000..fb67540 --- /dev/null +++ b/airflow_modules/er_dose_jobs.py @@ -0,0 +1,99 @@ +from __future__ import annotations + +from datetime import date, datetime, timedelta + +from er_dose.common.batch_log_repository import delete_batch_log +from er_dose.euv.euv_repository import ERDoseEUVRepository +from er_dose.infra.postgres_db import PostgresDB +from er_dose.raw.raw_repository import ERDoseRepository + + +PLAN_TASK_ID = "plan_dates" + + +def needs_reload(repository, target_date: date) -> bool: + source_count = repository.fetch_source_count(target_date) + target_count = repository.fetch_target_count(target_date) + if source_count == target_count: + return False + if target_count > 0: + return repository.fetch_source_count(target_date, distinct=True) != target_count + return True + + +def plan_dates(reference_date: str, lookback_days: int = 10) -> list[dict]: + """Return a JSON-safe plan; Airflow persists it before any parsing starts.""" + if lookback_days <= 0: + raise ValueError("lookback_days must be greater than 0") + end_date = date.fromisoformat(reference_date) + db = PostgresDB() + raw_repository = ERDoseRepository(db) + euv_repository = ERDoseEUVRepository(db) + plan = [] + for offset in range(lookback_days - 1, -1, -1): + target_date = end_date - timedelta(days=offset) + start_time = datetime.combine(target_date, datetime.min.time()) + plan.append({ + "target_date": target_date.isoformat(), + "start_time": start_time.isoformat(), + "end_time": (start_time + timedelta(days=1)).isoformat(), + "raw": needs_reload(raw_repository, target_date), + "euv": needs_reload(euv_repository, target_date), + }) + return plan + + +def _planned_day(ti, day_index: int) -> dict: + plan = ti.xcom_pull(task_ids=PLAN_TASK_ID) + if plan is None: + raise ValueError("saved ER Dose plan is missing; do not recompute dates during retry") + return plan[day_index] + + +def select_parse_task(parser_name: str, day_index: int, ti) -> str: + window = _planned_day(ti, day_index) + if window[parser_name]: + return f"day_{day_index:02d}_parse_{parser_name}" + return f"day_{day_index:02d}_{parser_name}_done" + + +def select_report_tasks(day_index: int, ti) -> list[str]: + window = _planned_day(ti, day_index) + if window["raw"] or window["euv"]: + return [f"day_{day_index:02d}_{job}" for job in + ("die_yield", "root_cause", "raw_count", "euv_count")] + return [] + + +def write_equipment_count_log(repository, batch_name: str, start_time: datetime, end_time: datetime) -> None: + counts = repository.fetch_equipment_counts(start_time, end_time) + message = f"equipment counts start_time={start_time.isoformat()} end_time={end_time.isoformat()}" + # A successful INSERT followed by a worker failure must not duplicate the snapshot on retry. + delete_batch_log(repository.db, batch_name, start_time.date(), "EQUIPMENT_COUNT", message) + repository.insert_batch_log( + batch_name=batch_name, + target_date=start_time.date(), + event_type="EQUIPMENT_COUNT", + message=message, + data={"equipment_counts": counts["equipment_counts"]}, + ) + + +def report_day(job_type: str, day_index: int, ti) -> None: + if job_type not in ("die_yield", "root_cause", "raw_count", "euv_count"): + raise ValueError("unknown ER Dose report job") + window = _planned_day(ti, day_index) + start_time = datetime.fromisoformat(window["start_time"]) + end_time = datetime.fromisoformat(window["end_time"]) + db = PostgresDB() + repository = ERDoseRepository(db) + if job_type == "die_yield": + repository.delete_die_yield_daily_summary(start_time, end_time) + repository.insert_die_yield_daily_summary(start_time, end_time) + elif job_type == "root_cause": + repository.delete_root_cause_daily_summary(start_time, end_time) + repository.insert_root_cause_daily_summary(start_time, end_time) + elif job_type == "raw_count": + write_equipment_count_log(repository, "ER_DOSE_RAW", start_time, end_time) + else: + write_equipment_count_log(ERDoseEUVRepository(db), "ER_DOSE_EUV", start_time, end_time) diff --git a/dags/er_dose_daily_dag.py b/dags/er_dose_daily_dag.py new file mode 100644 index 0000000..915f27a --- /dev/null +++ b/dags/er_dose_daily_dag.py @@ -0,0 +1,86 @@ +from __future__ import annotations + +from datetime import timedelta + +import pendulum +from airflow import DAG +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.providers.ssh.operators.ssh import SSHOperator +from airflow.utils.trigger_rule import TriggerRule + +from airflow_modules.er_dose_jobs import ( + PLAN_TASK_ID, plan_dates, report_day, select_parse_task, select_report_tasks, +) + + +LOOKBACK_DAYS = 10 +SSH_CONN_ID = "er_dose_parser" +LOCAL_TZ = pendulum.timezone("Asia/Seoul") + + +with DAG( + dag_id="er_dose_daily", + description="Plan ER Dose dates, parse, and retry reports independently", + schedule=None, + start_date=pendulum.datetime(2026, 4, 11, tz=LOCAL_TZ), + catchup=False, + max_active_runs=1, + default_args={"retries": 3, "retry_delay": timedelta(minutes=5)}, + tags=["er_dose", "batch"], +) as dag: + plan_task = PythonOperator( + task_id=PLAN_TASK_ID, + python_callable=plan_dates, + op_kwargs={ + "reference_date": "{{ dag_run.conf.get('reference_date', data_interval_end.in_timezone('Asia/Seoul').strftime('%Y-%m-%d')) }}", + "lookback_days": LOOKBACK_DAYS, + }, + ) + previous_parse = plan_task + for day_index in range(LOOKBACK_DAYS): + # Oldest day first; preserve RAW's prior-day lot state dependency. + for parser_name in ("euv", "raw"): + branch = BranchPythonOperator( + task_id=f"day_{day_index:02d}_choose_{parser_name}", + python_callable=select_parse_task, + op_kwargs={"parser_name": parser_name, "day_index": day_index}, + ) + parse_task = SSHOperator( + task_id=f"day_{day_index:02d}_parse_{parser_name}", + ssh_conn_id=SSH_CONN_ID, + command=( + "{{ var.value.er_dose_parser_command }}" + f" --parser ER_DOSE_{parser_name.upper()}" + " --date {{ ti.xcom_pull(task_ids='plan_dates')[" + f"{day_index}" + "]['target_date'] }} --chunk-size 30000" + ), + cmd_timeout=None, + get_pty=True, + do_xcom_push=False, + ) + done = EmptyOperator( + task_id=f"day_{day_index:02d}_{parser_name}_done", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + previous_parse >> branch + branch >> [parse_task, done] + parse_task >> done + previous_parse = done + + report_branch = BranchPythonOperator( + task_id=f"day_{day_index:02d}_choose_reports", + python_callable=select_report_tasks, + op_kwargs={"day_index": day_index}, + do_xcom_push=False, + ) + previous_parse >> report_branch + for job_type in ("die_yield", "root_cause", "raw_count", "euv_count"): + report_task = PythonOperator( + task_id=f"day_{day_index:02d}_{job_type}", + python_callable=report_day, + op_kwargs={"job_type": job_type, "day_index": day_index}, + do_xcom_push=False, + ) + report_branch >> report_task diff --git a/docs/er_dose_airflow.md b/docs/er_dose_airflow.md new file mode 100644 index 0000000..2c57afa --- /dev/null +++ b/docs/er_dose_airflow.md @@ -0,0 +1,70 @@ +# ER Dose 날짜별 실행과 장애 복구 + +## 구성 + +- DAG: `dags/er_dose_daily_dag.py`, ID `er_dose_daily` +- 작업 함수: `airflow_modules/er_dose_jobs.py` +- 최근 10일(기준일 포함)을 과거 날짜부터 확인합니다. `LOOKBACK_DAYS`로 범위를 변경합니다. +- `plan_dates`가 기존 원천/parsed/DISTINCT 건수 비교로 RAW·EUV 처리 여부를 확정합니다. 날짜와 자정 기준 `[start_time, end_time)`은 문자열로 XCom에 저장합니다. DB 연결이나 쿼리 결과 데이터는 XCom에 넣지 않습니다. +- 계획 작업은 조회만 수행합니다. 계획이 성공해 XCom이 저장되기 전에는 적재를 시작하지 않습니다. +- 날짜별 EUV → RAW 파싱은 `SSHOperator`로 Linux 파서 서버에서 순차 실행합니다. RAW의 이전 날짜 lot 상태 조회 순서를 보존합니다. +- 각 날짜 파싱 성공 후 수율·원인 서머리와 RAW·EUV 통계 네 작업은 Airflow worker에서 직접 DB를 조회·저장합니다. SSH로 통계를 실행하지 않습니다. 기존 repository 쿼리와 PostgreSQL 9.4 호환 Python JSON 구성을 재사용합니다. 어느 한쪽 파서만 재적재해도 네 후속 작업을 실행합니다. +- 실행 여부는 DAG의 `BranchPythonOperator`만 결정합니다. `choose_euv`/`choose_raw`는 파서 또는 합류 작업을 선택하고, `choose_reports`는 후속 작업 네 개 또는 빈 목록을 선택합니다. +- 처리 대상이 아니면 파서·서머리는 Airflow에서 `skipped` 처리되어 함수 자체가 호출되지 않습니다. 파서·서머리 함수에는 실행 여부 조건이 없으며, 호출되면 XCom의 날짜·기간대로 무조건 처리합니다. +- 파서 합류 작업은 `none_failed_min_one_success`를 사용합니다. 정상 건너뛰기는 다음 날짜로 이어지지만 파서 실패는 합류 이후 작업을 차단합니다. 서머리 실패는 다음 날짜 파싱을 차단하지 않습니다. + +예를 들어 기준일이 2026-05-14라면 `day_07_*`는 12일, `day_08_*`는 13일, `day_09_*`는 14일입니다. 모든 작업은 해당 실행의 `plan_dates`를 읽습니다. 재시도 시 오늘 날짜나 건수를 이용해 기간을 다시 결정하지 않습니다. + +## 배포 + +저장소에는 기존 ER Dose DAG가 없으므로 신규 DAG를 추가했습니다. Airflow 2.6.3과 `apache-airflow-providers-ssh==3.7.1`을 검증 기준으로 삼습니다. 운영 버전에 맞는 SSH provider를 설치해야 합니다. + +1. Airflow worker에는 이 저장소의 작업 함수·repository·DB 모듈과 pandas, psycopg2, SQLAlchemy를 제공합니다. scheduler에서도 DAG import가 가능해야 합니다. worker에서 직접 통계 대상 DB에 접속할 수 있어야 합니다. +2. Airflow worker의 기존 `er_dose.properties` 또는 `ER_DOSE_DB_DSN`/`DATABASE_URL` 설정으로 건수 비교와 통계를 실행합니다. DB 연결을 작업 인자로 전달하지 않습니다. +3. Linux 파서 서버에도 수정된 파서 코드를 배포하고, 서버 자체의 기존 DB 설정을 유지합니다. 원격 CLI는 받은 `--date`만 재적재하며 서머리·통계를 실행하지 않습니다. +4. Airflow SSH Connection `er_dose_parser`에 실제 Linux 접속 정보를 등록합니다. 기존 Connection을 사용할 경우 DAG의 `SSH_CONN_ID`를 해당 ID로 변경합니다. 호스트와 자격 증명은 코드에 넣지 않습니다. +5. Airflow Variable `er_dose_parser_command`에 실행 접두 명령을 설정합니다. Linux 예시(실제 경로로 변경): + + ```sh + cd /opt/tempus-parser && exec /opt/venv/bin/python -m er_dose.run_er_dose_batch + ``` + + DAG는 저장된 계획에서 날짜를 읽어 `--parser ER_DOSE_RAW --date 2026-05-12 --chunk-size 30000` 같은 인자를 덧붙입니다. 공백 있는 경로는 Variable에서 셸 인용해야 합니다. Variable은 운영자가 관리하는 명령 설정이며, 실행 요청의 사용자 입력을 명령 접두어로 사용하지 않습니다. `nohup`, `&` 등으로 백그라운드 실행하지 않고 원격 파서가 끝날 때까지 기다려야 종료 코드를 확인할 수 있습니다. +6. DAG를 배포하고 수동 실행합니다. 예: `{"reference_date": "2026-05-14"}`. 생략하면 해당 실행의 `data_interval_end`를 서울 날짜로 변환합니다. +7. 원천 적재 완료 시각이 확인되지 않았으므로 신규 DAG는 `schedule=None`입니다. 운영 주기에 맞게 연결·설정한 후 기존 직접 실행 스케줄과 중복되지 않도록 전환합니다. + +Airflow 메타데이터 DB는 파서 대상 PostgreSQL과 별도 장애 영역에 두는 것이 전제입니다. DB 연결 종료나 SSH 명령의 0이 아닌 종료 코드는 작업 실패로 전파되고, Airflow가 5분 간격으로 최대 3회 재시도합니다. 건수 비교와 통계 작업은 시도마다 새로운 `PostgresDB`를 생성합니다. SSH 출력은 대량 로그가 XCom에 저장되지 않도록 `do_xcom_push=False`입니다. + +긴 파싱 작업이 SSH 기본 명령 제한 시간으로 끊기지 않도록 `cmd_timeout=None`을 사용합니다. Linux 원격 프로세스는 `get_pty=True`로 실행합니다. SSH 접속 단절 시 원격 프로세스 종료 여부를 실제 서버에서 확인해야 하며, 살아 있는 동일 날짜 작업과 재시도가 겹치면 안 됩니다. 전체 작업 및 락 대기의 제한 시간은 운영 소요 시간을 기준으로 별도 설정합니다. + +## 장애 복구 + +- 파싱 실패: 실패한 날짜의 파서 작업을 재실행합니다. 계획의 처리 여부를 다시 조회하지 않고 해당 날짜 파티션을 다시 적재합니다. 이후 선행 실패로 막힌 작업도 함께 Clear하여 실행합니다. +- 서머리 실패: 해당 DAG Run에서 실패한 `day_XX_die_yield` 또는 `day_XX_root_cause`만 Clear합니다. 성공한 파싱과 `plan_dates`는 Clear하지 않습니다. +- 통계 실패: 실패한 `day_XX_raw_count` 또는 `day_XX_euv_count`만 Clear합니다. +- 서머리는 DELETE → INSERT를 별도 메서드·별도 커밋으로 수행합니다. DELETE 후 실패하면 집계가 비어 있을 수 있으며, 재시도는 DELETE부터 시작합니다. +- 통계는 조회 성공 후 동일 배치·날짜·이벤트·기간 메시지의 기존 행을 DELETE하고 INSERT합니다. INSERT 커밋 후 worker가 종료되어도 재시도 시 동일 기간 스냅샷을 중복 저장하지 않습니다. 이전 스냅샷을 매 실행 이력으로 보존하는 방식은 아닙니다. +- `plan_dates`가 없는 상태에서는 임의로 날짜를 재계산하지 않고 오류로 종료합니다. 계획을 실수로 삭제했다면 필요한 기간을 명시한 별도 복구가 필요합니다. +- 자동 재시도를 소진한 실패 작업은 다음 정기 실행의 건수 비교로 복구되지 않을 수 있습니다. 기존 실패 DAG Run의 작업을 Clear해야 하며, 실패 알림을 운영에 연결해야 합니다. +- `max_active_runs=1`로 정상 DAG 실행 중첩을 제한합니다. 과거 Run 복구 시 같은 날짜를 쓰는 다른 DAG·CLI 작업은 중지해야 합니다. 별도 시스템 간 동시 실행 잠금은 추가하지 않았습니다. + +## 기존 진입점 변경 + +`processor.run()`은 날짜나 기간 없이 호출할 수 없습니다. `lookback_days`, `reference_date`, `run_recent_days()`는 더 이상 파서 API가 아닙니다. `--date`는 건수와 관계없이 지정된 날짜만 재적재합니다. 직접 시간 범위를 지정하는 호출은 기존 파싱/적재 동작을 유지합니다. 파서 CLI만 실행하면 서머리·통계는 실행되지 않으므로, 운영 스케줄을 신규 DAG와 함께 전환해야 합니다. + +## 검증 범위 + +기본 테스트 108개가 통과했고, Airflow 미설치 환경에서는 DAG 테스트 모듈 1개를 건너뜁니다. 임시 Airflow 2.6.3 + SSH provider 3.7.1 + Python 3.10 환경에서는 DAG·작업 테스트 25개가 통과했습니다. + +111개 작업의 DAG import와 의존성을 확인했으며, 실제 `dag.test()`와 격리된 SQLite 메타데이터 DB를 사용해 다음을 검증했습니다. + +- 변경 없는 날짜는 SSH 파서를 호출하지 않고 다음 날짜로 진행 +- RAW만/EUV만/둘 다 필요한 경우 각각 저장된 날짜로 SSH 명령 실행 +- SSH 실패 시 해당 날짜 서머리·통계와 다음 날짜 파싱 차단 +- 서머리 실패 시 다른 통계·다음 날짜 파싱은 계속 진행 +- 원격 명령의 비정상 종료 코드가 Airflow 실패로 전파되고, 재시도 명령은 동일 날짜 유지 +- 통계는 Airflow 작업이 repository를 직접 호출해 조회·저장 + +SSH 전송과 대상 DB 호출은 mock입니다. 실제 Linux 서버 SSH 접속, 운영 DB의 관리자 연결 종료, 운영 executor 배포는 별도 검증이 필요합니다. + +참고: [Airflow 2.6.3 XCom](https://airflow.apache.org/docs/apache-airflow/2.6.3/core-concepts/xcoms.html). 실패한 작업 자신의 XCom은 재시도 때 정리될 수 있으므로, 계획은 별도의 성공한 `plan_dates` 작업에 보관합니다. diff --git a/docs/er_dose_flow.drawio b/docs/er_dose_flow.drawio index a687b67..791e981 100644 --- a/docs/er_dose_flow.drawio +++ b/docs/er_dose_flow.drawio @@ -187,39 +187,39 @@ - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + diff --git a/docs/er_dose_learning_roadmap.md b/docs/er_dose_learning_roadmap.md new file mode 100644 index 0000000..4d1133e --- /dev/null +++ b/docs/er_dose_learning_roadmap.md @@ -0,0 +1,161 @@ +# ER Dose 프로젝트 기반 기술 학습 로드맵 + +## 목적 + +이 문서는 ER Dose 대용량 배치에서 실제로 발생한 메모리 증가, 긴 처리 시간, 부분 적재, 재실행 판단 문제를 학습 과제로 연결한다. + +현재 프로젝트에서는 Kafka나 Redis 같은 별도 인프라를 먼저 도입하기보다 다음 역량을 우선해서 학습하는 편이 효과적이다. + +1. PostgreSQL 대량 조회·적재와 실행계획 분석 +2. Airflow 배치의 멱등성과 실패 복구 +3. 실제 PostgreSQL을 사용하는 통합 테스트 +4. Python 및 Linux 성능 분석 +5. Repository, 트랜잭션 경계와 기술 의사결정 기록 + +기존 성능 측정 결과와 처리 구조는 [Database Streaming Optimization](./db_streaming_optimization.md)을 기준 자료로 사용한다. + +## 1. PostgreSQL 대량 처리 + +### 추천 자료 + +- [PostgreSQL EXPLAIN 사용법](https://www.postgresql.org/docs/current/using-explain.html) +- [PostgreSQL 대량 데이터 적재](https://www.postgresql.org/docs/current/populate.html) +- [PostgreSQL COPY](https://www.postgresql.org/docs/current/sql-copy.html) +- [PostgreSQL 테이블 파티셔닝](https://www.postgresql.org/docs/current/ddl-partitioning.html) +- [PostgreSQL 트랜잭션 격리](https://www.postgresql.org/docs/current/transaction-iso.html) +- [PostgreSQL MVCC](https://www.postgresql.org/docs/current/mvcc-intro.html) +- [PostgreSQL WAL](https://www.postgresql.org/docs/current/wal-intro.html) +- [pg_stat_statements](https://www.postgresql.org/docs/current/pgstatstatements.html) +- [PostgreSQL 명령 진행률 확인](https://www.postgresql.org/docs/current/progress-reporting.html) + +### 집중해서 볼 내용 + +- `EXPLAIN ANALYZE`의 예상 행 수와 실제 행 수 차이 +- `Buffers`, `temp read/written`, `Sort Method`, `loops` +- 날짜 조건이 원천 파티션을 제대로 제거하는지 나타내는 partition pruning +- `COPY`와 일반 `INSERT`의 처리 방식 차이 +- 청크별 커밋과 날짜 전체 트랜잭션의 성능·복구 범위 차이 +- 인덱스를 유지한 채 대량 적재할 때 발생하는 쓰기 비용 +- 임시 테이블에 적재·검증한 뒤 파티션을 교체하는 방식 +- `pg_stat_statements`의 실행시간, 블록 읽기, 임시 블록, WAL 발생량 +- 실행 중인 `COPY`를 `pg_stat_progress_copy`로 확인하는 방법 + +### 프로젝트 실습 + +1. `docs/db_streaming_optimization.md`의 RAW 조회 SQL에 `EXPLAIN (ANALYZE, BUFFERS, TIMING OFF, SUMMARY ON)`을 적용한다. +2. source count, target count, 설비별 count, 일별 summary 쿼리의 실행계획을 각각 저장한다. +3. 대상 날짜 파티션만 읽는지, 정렬이 디스크로 내려가는지, 예상 행 수가 실제 값과 크게 다른지 비교한다. +4. 동일한 테스트 데이터로 청크별 커밋과 날짜 단위 커밋의 시간 및 실패 결과를 비교한다. +5. 가능하면 임시 테이블 적재 후 검증·파티션 교체 방식을 별도 실험한다. + +`EXPLAIN ANALYZE`는 쿼리를 실제로 실행한다. 운영 환경에서는 부하가 낮은 시간에 수행하고, `pg_stat_statements` 활성화는 DB 관리자와 먼저 협의한다. + +## 2. Airflow 멱등성과 실패 복구 + +### 추천 자료 + +- [Airflow Best Practices](https://airflow.apache.org/docs/apache-airflow/stable/best-practices.html) +- [Airflow DAG Run과 Data Interval](https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/dag-run.html) +- [Airflow Backfill](https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/backfill.html) +- [Airflow Metrics](https://airflow.apache.org/docs/apache-airflow/stable/administration-and-deployment/logging-monitoring/metrics.html) + +### 집중해서 볼 내용 + +- Airflow task를 하나의 트랜잭션처럼 취급하는 이유 +- 같은 날짜를 여러 번 실행해도 최종 결과가 같아야 하는 멱등성 +- 현재 시간이 아니라 명시적인 data interval을 기준으로 읽고 쓰는 방식 +- retry, clear, backfill 실행 시 데이터가 중복되거나 일부만 남지 않도록 하는 방법 +- 태스크 실행시간, 실패, 재시도, 지연을 지표로 남기는 방법 + +### 프로젝트 실습 + +다음 실패 시나리오를 같은 날짜 데이터로 반복 검증한다. + +| 시나리오 | 검증할 결과 | +| --- | --- | +| 정상 실행 후 같은 날짜 재실행 | 최종 건수와 내용이 동일함 | +| 세 번째 청크 적재 전에 실패 | 불완전한 결과가 사용자에게 노출되는지 확인 | +| 세 번째 청크 적재 후 실패 | 재실행 후 중복과 누락이 없는지 확인 | +| summary 갱신 전 실패 | parsed와 summary의 불일치 복구 여부 확인 | +| 동일 날짜 동시 실행 | 중복 재적재 또는 truncate 충돌 여부 확인 | + +Airflow 문서는 최신 버전을 가리킨다. 실제 코드에 적용할 때는 운영 서버의 `airflow version`과 같은 버전의 문서를 선택한다. + +## 3. 실제 PostgreSQL 통합 테스트 + +### 추천 자료 + +- [Testcontainers Python](https://testcontainers-python.readthedocs.io/en/latest/) +- [Docker 공식 Testcontainers Python 실습](https://docs.docker.com/guides/testcontainers-python-getting-started/) +- [pytest fixture](https://docs.pytest.org/en/stable/how-to/fixtures.html) +- [pytest parameterize](https://docs.pytest.org/en/stable/how-to/parametrize.html) +- [GitHub Actions Python 테스트](https://docs.github.com/en/actions/tutorials/build-and-test-code/python) +- [Python pyproject.toml 작성법](https://packaging.python.org/en/latest/guides/writing-pyproject-toml/) + +### 먼저 만들 테스트 + +1. PostgreSQL 컨테이너에 실제 parent table과 날짜 파티션을 생성한다. +2. 작은 RAW 원천 데이터를 넣고 parser와 `COPY`를 거쳐 parsed 결과를 검증한다. +3. 지정한 청크에서 예외를 발생시킨 뒤 트랜잭션이 어디까지 반영됐는지 확인한다. +4. 같은 날짜를 재실행하고 중복·누락 없이 결과가 동일한지 확인한다. +5. RAW와 EUV에 동일한 재적재 판단 조건을 매개변수화해 검증한다. + +FakeDB 단위 테스트는 분기 검증에 유지하고, SQL 문법, 파티션, `COPY`, commit/rollback 동작은 실DB 통합 테스트가 담당하도록 나눈다. + +## 4. 성능 분석 + +### 추천 자료 + +- [py-spy](https://github.com/benfred/py-spy) +- [Python cProfile](https://docs.python.org/3/library/profile.html) +- [Python tracemalloc](https://docs.python.org/3/library/tracemalloc.html) +- [Python concurrent.futures](https://docs.python.org/3/library/concurrent.futures.html) +- [Brendan Gregg의 USE Method](https://www.brendangregg.com/Articles/The_USE_Method.pdf) +- [Brendan Gregg의 Linux Performance 자료](https://www.brendangregg.com/linuxperf.html) + +### 분석 순서 + +1. 전체 시간을 fetch, parse, insert, summary로 나눠 기록한다. +2. `py-spy`로 Python CPU 사용 구간과 DB 응답 대기 구간을 구분한다. +3. `cProfile`로 parser 함수별 누적시간과 호출 횟수를 확인한다. +4. 메모리가 계속 증가하면 `tracemalloc` snapshot을 청크 전후로 비교한다. +5. 서버에서는 CPU 사용률, run queue, 메모리·swap, 디스크 대기시간, 네트워크 에러를 확인한다. +6. worker 수와 chunk size는 한 번에 하나씩만 변경하고 같은 데이터로 비교한다. + +멀티스레드나 멀티프로세스는 병목이 Python CPU인지 DB I/O인지 확인한 뒤 적용한다. worker 수만 늘리고 DB가 포화되면 처리시간은 줄지 않고 경합과 메모리 사용량만 증가할 수 있다. + +## 5. 코드 구조와 의사결정 + +### 추천 자료 + +- [Cosmic Python Repository Pattern](https://www.cosmicpython.com/book/chapter_02_repository) +- [Cosmic Python Unit of Work Pattern](https://www.cosmicpython.com/book/chapter_06_uow.html) +- [Architecture Decision Records](https://adr.github.io/) +- [Martin Fowler의 Architecture Decision Record](https://martinfowler.com/bliki/ArchitectureDecisionRecord.html) + +### 프로젝트에 적용할 기준 + +- Repository는 SQL과 영속성 처리를 담당한다. +- Processor는 날짜, 청크, 파싱·적재 순서와 같은 처리 흐름을 담당한다. +- 여러 Repository 작업을 하나의 원자적 작업으로 묶어야 할 때만 Unit of Work를 검토한다. +- 한 줄을 줄이기 위한 추상화보다 실패 범위와 테스트 가능성을 명확히 만드는 추상화를 우선한다. +- 성능 구조를 변경할 때는 선택 이유, 대안, 실측 결과와 되돌릴 조건을 ADR로 남긴다. + +우선 기록할 ADR 후보는 다음과 같다. + +1. 적재 worker를 1개로 제한한 이유 +2. 원천·대상 count가 다를 때만 `DISTINCT`를 실행하는 이유 +3. 청크별 commit과 날짜 단위 commit 중 어떤 방식을 선택했는지 +4. 부분 적재 노출을 막기 위해 임시 파티션 교체를 도입할지 + +## 권장 학습 순서 + +| 단계 | 학습 및 실습 | 결과물 | +| --- | --- | --- | +| 1 | `EXPLAIN`, partition pruning, `pg_stat_statements` | 실제 쿼리 실행계획 분석 문서 | +| 2 | Airflow 멱등성 및 실패 복구 | 실패 지점별 재실행 테스트 표 | +| 3 | Testcontainers와 pytest | 실제 PostgreSQL 통합 테스트 | +| 4 | `py-spy`, `cProfile`, USE Method | flame graph와 병목 분석 결과 | +| 5 | Repository, Unit of Work, ADR | 트랜잭션 경계 ADR | + +가장 먼저 수행할 과제는 RAW 하루치 조회 실행계획 분석과 강제 실패 통합 테스트다. 이 두 결과가 있어야 다음 최적화가 parser, DB query, `COPY`, transaction 중 어디를 대상으로 해야 하는지 근거를 갖고 결정할 수 있다. diff --git a/docs/pr/er_dose_final_statistics_validation.md b/docs/pr/er_dose_final_statistics_validation.md new file mode 100644 index 0000000..5bb57bb --- /dev/null +++ b/docs/pr/er_dose_final_statistics_validation.md @@ -0,0 +1,36 @@ +# RAW 마지막 서머리·통계 검증 + +검증일: 2026-09-14. 격리된 로컬 PostgreSQL 16 컨테이너를 사용했으며 운영 DB에는 접근하지 않았다. + +## DB 실행 및 기능 + +- 기존 커밋의 `execute("select :value", ...)`는 실제 psycopg2에서 SyntaxError를 재현했다. 기존 SQL 표기를 유지하고 execute 경계에서 바인딩 형식만 변환했다. +- 합성 RAW parsed/원천 각 60,000행(20개 설비, 20,000개 lot), EUV parsed/원천 각 100,000행을 사용했다. +- DIE Yield 합계 20,000, Root Cause 빈도 합계 100,000을 확인했다. +- 로그 24행의 SQL 집계 결과, 건수 일치, STATISTICS_RECORDED, created_at을 확인했다. +- INSERT를 의도적으로 실패시켜 앞서 실행한 DELETE가 되돌아가지 않음을 확인했다. 두 메서드는 별도 커밋이다. + +## 시간 비교 + +아래 수치는 DELETE와 INSERT를 각 서머리 작업 안에서 연속 실행하던 수정 전 구조의 측정값이다. 현재 코드는 processor에서 DELETE 두 개를 먼저 직접 실행한 뒤 INSERT 두 개와 통계 두 개를 병렬 제출한다. 따라서 아래 수치를 현재 구조의 성능 측정값으로 간주하지 않는다. + +동일한 DB에서 순차·병렬을 각각 한 번 예열하고 5회씩 교차 실행했다. 병렬 작업자는 4개이며 서머리 2개와 통계 조회/로그 INSERT 2개를 함께 실행했다. RAW 파싱·적재 시간은 포함하지 않는다. + +| 실행 | 순차(초) | 병렬(초) | +|---|---:|---:| +| 1 | 0.5261 | 0.4030 | +| 2 | 0.5255 | 0.4058 | +| 3 | 0.5312 | 0.4227 | +| 4 | 0.6430 | 0.3871 | +| 5 | 0.6240 | 0.4628 | +| 중앙값 | 0.5312 | 0.4058 | + +이 로컬 표본에서는 중앙값이 23.6% 감소했다. 데이터 규모, 인덱스, 캐시 및 DB 부하가 운영과 다르므로 운영 속도 향상을 보장하는 결과가 아니다. 과거 전체 파이프라인 실행 시간과 합산하거나 비교하지 않는다. + +## 반영 전 조건 + +- 로그 테이블 DDL `er_dose/sql/create_batch_event_log.sql`을 대상 환경에 먼저 적용해야 한다. 이번에는 임시 검증 DB에만 적용했다. +- 실행 진입점은 배치 하나를 선택하며 EUV 완료를 자동으로 기다리지 않는다. 운영에서 EUV 적재가 RAW 마지막 집계 전에 완료되도록 순서를 보장해야 한다. +- 통계 로그는 전체 배치 성공 표시가 아니다. 다른 병렬 작업 실패 시 이미 저장한 통계가 남을 수 있다. + +현재 로그는 RAW/EUV 별도 두 행이며 각 `data`에는 `equipment_counts` 배열만 저장한다. 원소 필드는 eq_name/source_count/target_count다. 처리 범위는 message에 기록하며 통계 로그는 전체 배치 성공 표시가 아니다. 위 시간 기록은 현재 로그 형식의 재측정 결과가 아니다. diff --git a/er_dose/common/batch_log_repository.py b/er_dose/common/batch_log_repository.py new file mode 100644 index 0000000..c290980 --- /dev/null +++ b/er_dose/common/batch_log_repository.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +import json +from datetime import date +from typing import Any + +from er_dose.infra.postgres_db import PostgresDB + + +BATCH_LOG_TABLE = "mbeat.batch_event_log" + + +def delete_batch_log( + db: PostgresDB, + batch_name: str, + target_date: date, + event_type: str, + message: str, +) -> None: + db.execute( + f""" + delete from {BATCH_LOG_TABLE} + where batch_name = :batch_name + and target_date = :target_date + and event_type = :event_type + and message = :message + """, + params={ + "batch_name": batch_name, + "target_date": target_date, + "event_type": event_type, + "message": message, + }, + ) + + +def insert_batch_log( + db: PostgresDB, + batch_name: str, + target_date: date, + event_type: str, + message: str, + data: dict[str, Any], +) -> None: + query = f""" + insert into {BATCH_LOG_TABLE} ( + batch_name, + target_date, + event_type, + message, + data + ) values ( + :batch_name, + :target_date, + :event_type, + :message, + cast(:data as jsonb) + ) + """ + db.execute( + query, + params={ + "batch_name": batch_name, + "target_date": target_date, + "event_type": event_type, + "message": message, + "data": json.dumps(data, ensure_ascii=False, default=str), + }, + ) diff --git a/er_dose/common/reload_processor.py b/er_dose/common/reload_processor.py index 27023c3..8349cf5 100644 --- a/er_dose/common/reload_processor.py +++ b/er_dose/common/reload_processor.py @@ -31,54 +31,6 @@ def _run_window( ) -> int: raise NotImplementedError - def run_recent_days( - self, - lookback_days: int = 2, - reference_date: date | None = None, - chunk_size: int = 10000, - ) -> None: - if lookback_days <= 0: - raise ValueError("lookback_days must be greater than 0") - if chunk_size <= 0: - raise ValueError("chunk_size must be greater than 0") - - end_date = reference_date or date.today() - start_date = end_date - timedelta(days=lookback_days - 1) - - checked_dates = 0 - reloaded_dates = 0 - source_rows = 0 - inserted_rows = 0 - current_date = start_date - while current_date <= end_date: - checked_dates += 1 - source_count = self.repository.fetch_source_count(current_date) - target_count = self.repository.fetch_target_count(current_date) - - if source_count == target_count: - current_date += timedelta(days=1) - continue - if target_count > 0: - distinct_source_count = self.repository.fetch_source_count(current_date, distinct=True) - if distinct_source_count == target_count: - current_date += timedelta(days=1) - continue - - reloaded_dates += 1 - source_rows += source_count - inserted_rows += self._reload_target_date(target_date=current_date, chunk_size=chunk_size) - current_date += timedelta(days=1) - - print( - f"{self.log_prefix} " - f"lookback_done start_date={start_date.isoformat()} " - f"end_date={end_date.isoformat()} " - f"checked_dates={checked_dates} " - f"reloaded_dates={reloaded_dates} " - f"source_rows={source_rows} " - f"inserted={inserted_rows}" - ) - def _reload_target_date(self, target_date: date, chunk_size: int) -> int: start_time = datetime.combine(target_date, datetime.min.time()) end_time = start_time + timedelta(days=1) diff --git a/er_dose/euv/euv_processor.py b/er_dose/euv/euv_processor.py index 31fa371..4836c64 100644 --- a/er_dose/euv/euv_processor.py +++ b/er_dose/euv/euv_processor.py @@ -23,31 +23,18 @@ def run( start_time: datetime | None = None, end_time: datetime | None = None, chunk_size: int = 10000, - lookback_days: int = 2, - reference_date: date | None = None, target_date: date | None = None, ) -> None: + if chunk_size <= 0: + raise ValueError("chunk_size must be greater than 0") if target_date is not None: - self.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=chunk_size, - ) + self._reload_target_date(target_date=target_date, chunk_size=chunk_size) return - if start_time is None and end_time is None: - self.run_recent_days( - lookback_days=lookback_days, - reference_date=reference_date, - chunk_size=chunk_size, - ) - return if start_time is None or end_time is None: raise ValueError("start_time and end_time are required") if start_time >= end_time: raise ValueError("start_time must be earlier than end_time") - if chunk_size <= 0: - raise ValueError("chunk_size must be greater than 0") self._run_window(start_time=start_time, end_time=end_time, chunk_size=chunk_size) @@ -118,11 +105,6 @@ def _run_window( for target_date in sorted(inserted_target_dates): self.repository.analyze_target_partition(target_date, connection=connection) - self.repository.upsert_root_cause_daily_summary(target_date, connection=connection) - print( - "[ER_DOSE_EUV] " - f"summary updated (root_cause) partition_date={target_date}" - ) print( "[ER_DOSE_EUV] " diff --git a/er_dose/euv/euv_repository.py b/er_dose/euv/euv_repository.py index 92fdeb4..0765053 100644 --- a/er_dose/euv/euv_repository.py +++ b/er_dose/euv/euv_repository.py @@ -1,10 +1,11 @@ from __future__ import annotations from datetime import date, datetime, timedelta -from typing import Iterator +from typing import Any, Iterator import pandas as pd +from er_dose.common.batch_log_repository import insert_batch_log from er_dose.common.sql_filters import active_nxe_eq_filter from er_dose.infra.postgres_db import PostgresDB @@ -58,6 +59,23 @@ class ERDoseEUVRepository: def __init__(self, db: PostgresDB): self.db = db + def insert_batch_log( + self, + batch_name: str, + target_date: date, + event_type: str, + message: str, + data: dict[str, Any], + ) -> None: + insert_batch_log( + self.db, + batch_name=batch_name, + target_date=target_date, + event_type=event_type, + message=message, + data=data, + ) + def fetch_source_count(self, target_date: date, distinct: bool = False) -> int: start_time = datetime.combine(target_date, datetime.min.time()) end_time = start_time + timedelta(days=1) @@ -91,6 +109,55 @@ def fetch_target_count(self, target_date: date) -> int: return 0 return int(df.iloc[0]["row_count"]) + def fetch_equipment_counts(self, start_time: datetime, end_time: datetime) -> dict[str, Any]: + query = f""" + with source_counts as ( + select + r.eq_name, + count(*) as source_count + from {EUV_RAW_TABLE} r + where r.code_occur_time >= :start_time + and r.code_occur_time < :end_time + and lower(r.contents) like '%dose error detected in file:%' + and lower(r.contents) like '%root cause%' + and {active_nxe_eq_filter("r.eq_name")} + group by r.eq_name + ), + target_counts as ( + select + p.eq_name, + count(*) as target_count + from {ROOT_CAUSE_TABLE} p + where p.code_occur_time >= :start_time + and p.code_occur_time < :end_time + and {active_nxe_eq_filter("p.eq_name")} + group by p.eq_name + ) + select + coalesce(s.eq_name, t.eq_name) as eq_name, + coalesce(s.source_count, 0) as source_count, + coalesce(t.target_count, 0) as target_count, + sum(coalesce(s.source_count, 0)) over () as total_source_count, + sum(coalesce(t.target_count, 0)) over () as total_target_count, + min(case when coalesce(s.source_count, 0) = coalesce(t.target_count, 0) + then 1 else 0 end) over () as matched + from source_counts s + full outer join target_counts t on t.eq_name = s.eq_name + order by eq_name + """ + df = self.db.select(query, params={"start_time": start_time, "end_time": end_time}) + if df is None or df.empty: + return {"source_count": 0, "target_count": 0, "matched": True, "equipment_counts": []} + return { + "source_count": int(df.iloc[0]["total_source_count"]), + "target_count": int(df.iloc[0]["total_target_count"]), + "matched": bool(df.iloc[0]["matched"]), + "equipment_counts": [ + {"eq_name": row.eq_name, "source_count": int(row.source_count), "target_count": int(row.target_count)} + for row in df.itertuples(index=False) + ], + } + def truncate_target_partition(self, target_date: date, connection=None) -> int: parsed_table = self._partition_table_name(ROOT_CAUSE_TABLE, target_date) return self.db.execute(f"truncate table {parsed_table}", connection=connection) @@ -224,64 +291,6 @@ def analyze_target_partition(self, target_date: str, connection=None) -> int: partition_table = f"{ROOT_CAUSE_TABLE}_1_prt_p{target_date.replace('-', '')}" return self.db.execute(f"ANALYZE {partition_table}", connection=connection) - def upsert_root_cause_daily_summary(self, target_date: date | str, connection=None) -> int: - if isinstance(target_date, str): - target_date = datetime.strptime(target_date, "%Y-%m-%d").date() - start_time = datetime.combine(target_date, datetime.min.time()) - end_time = start_time + timedelta(days=1) - - query = f""" - insert into prism_common.de_trend_root_cause_daily ( - occur_date, - eq_name, - root_cause, - frequency, - created_at - ) - with root_cause_data as ( - select - e.code_occur_time::date as occur_date, - e.eq_name, - replace( - trim( - split_part( - case - when e.root_cause like '%CE & MP%' then replace(e.root_cause, 'CE & MP', 'CE @ MP') - when e.root_cause like '%l2Dx & l2Dy%' then replace(e.root_cause, 'L2Dx & L2Dy', 'L2Dx @ L2Dx') - when e.root_cause like '%E&T%' then replace(e.root_cause, 'E&T', 'E@T') - else e.root_cause - end, - '&', - 1 - ) - ), - '@', - '&' - ) as root_cause - from {ROOT_CAUSE_TABLE} e - where e.code_occur_time >= :start_time - and e.code_occur_time < :end_time - and e.code = 'OSD-0200' - and e.root_cause is not null - and trim(e.root_cause) != '' - ) - select - occur_date, - eq_name, - root_cause, - count(*) as frequency, - now() as created_at - from root_cause_data - where root_cause is not null - and root_cause != '' - group by occur_date, eq_name, root_cause - on conflict (occur_date, eq_name, root_cause) - do update set - frequency = excluded.frequency, - created_at = now(); - """ - return self.db.execute(query, params={"start_time": start_time, "end_time": end_time}, connection=connection) - def transaction(self): return self.db.transaction() diff --git a/er_dose/infra/postgres_db.py b/er_dose/infra/postgres_db.py index b70149f..74c1783 100644 --- a/er_dose/infra/postgres_db.py +++ b/er_dose/infra/postgres_db.py @@ -193,6 +193,17 @@ def copy_insert_df(self, table_name: str, df, connection=None) -> int: return len(normalized_df) def execute(self, query: str, params=None, connection=None) -> int: + # Keep repository SQL in SQLAlchemy's :name style, as in select(). + # Only translate named binds; existing native driver calls stay valid. + if isinstance(params, dict): + from sqlalchemy import text + from sqlalchemy.dialects.postgresql.psycopg2 import dialect + + compiled = text(query).compile(dialect=dialect()) + if compiled.params: + query = str(compiled) + params = compiled.construct_params(params) + own_connection = connection is None conn = connection or self._connect_raw() try: diff --git a/er_dose/raw/raw_parser.py b/er_dose/raw/raw_parser.py index ef096d6..043bc6f 100644 --- a/er_dose/raw/raw_parser.py +++ b/er_dose/raw/raw_parser.py @@ -57,10 +57,14 @@ def parse_dose_error(raw: RawErLog) -> ParsedErDoseError: n_slit = extract_int(contents, rf"n_slit\s*=\s*{_INT_RE}") elif code_norm.startswith("LO-"): - if code_norm in {"LO-0050", "LO-0051", "LO-0052"}: + if code_norm in {"LO-0050", "LO-0051"}: lot_id = extract_text(contents, r"lot\s+'([^']+)'") lot_name = lot_id.split(".", maxsplit=1)[0] if lot_id is not None else None lot_seq = extract_int(contents, rf"\(id\s*=\s*{_INT_RE}\)") + elif code_norm == "LO-0052": + lot_id = extract_text(contents, r"lot\s+(?:\(\s*name\s*=\s*)?'([^']+)'") + lot_name = lot_id.split(".", maxsplit=1)[0] if lot_id is not None else None + lot_seq = extract_int(contents, rf"(?:\(|,\s*)id\s*=\s*{_INT_RE}\)") if lot_seq is None: lot_seq = extract_first_int(contents, _LOT_SEQ_PATTERNS, minimum=1) wafer_seq = extract_first_int(contents, _WAFER_SEQ_PATTERNS, minimum=1) diff --git a/er_dose/raw/raw_processor.py b/er_dose/raw/raw_processor.py index a0abb01..c418e35 100644 --- a/er_dose/raw/raw_processor.py +++ b/er_dose/raw/raw_processor.py @@ -22,6 +22,7 @@ class ERDoseProcessor(CountReloadProcessor): log_prefix = "[ER_DOSE]" + batch_name = "ER_DOSE_RAW" def __init__(self, repository: ERDoseRepository): self.repository = repository @@ -35,26 +36,17 @@ def run( end_time: datetime | None = None, chunk_size: int = 10000, target_date: date | None = None, - lookback_days: int = 2, - reference_date: date | None = None, ) -> None: + if chunk_size <= 0: + raise ValueError("chunk_size must be greater than 0") if target_date is not None: - start_time = datetime.combine(target_date, datetime.min.time()) - end_time = start_time + timedelta(days=1) - - if start_time is None and end_time is None: - self.run_recent_days( - lookback_days=lookback_days, - reference_date=reference_date, - chunk_size=chunk_size, - ) + self._reload_target_date(target_date=target_date, chunk_size=chunk_size) return + if start_time is None or end_time is None: raise ValueError("start_time and end_time are required") if start_time >= end_time: raise ValueError("start_time must be earlier than end_time") - if chunk_size <= 0: - raise ValueError("chunk_size must be greater than 0") self._run_window(start_time=start_time, end_time=end_time, chunk_size=chunk_size) @@ -81,15 +73,26 @@ def _run_window( f"preloaded_eq={len(self.lot_states)}" ) - with ThreadPoolExecutor(max_workers=1, thread_name_prefix="er-dose-insert") as insert_executor: - for chunk_index, raw_df in enumerate( - self.repository.fetch_raw_logs_in_chunks( - start_time=start_time, - end_time=end_time, - chunk_size=chunk_size, - ), - start=1, - ): + raw_chunks = iter( + self.repository.fetch_raw_logs_in_chunks( + start_time=start_time, + end_time=end_time, + chunk_size=chunk_size, + ) + ) + + with ThreadPoolExecutor(max_workers=1, thread_name_prefix="er-dose-fetch") as fetch_executor, \ + ThreadPoolExecutor(max_workers=1, thread_name_prefix="er-dose-insert") as insert_executor: + pending_fetch = fetch_executor.submit(next, raw_chunks, None) + chunk_index = 0 + + while True: + raw_df = pending_fetch.result() + if raw_df is None: + break + + chunk_index += 1 + pending_fetch = fetch_executor.submit(next, raw_chunks, None) chunk_fetched = int(len(raw_df)) fetched_count += chunk_fetched print( @@ -100,6 +103,7 @@ def _run_window( ) parsed_rows = self._parse_chunk(raw_df) + del raw_df parsed_count = len(parsed_rows) print( "[ER_DOSE] " @@ -129,6 +133,7 @@ def _run_window( continue parsed_df = pd.DataFrame(parsed_rows) + del parsed_rows chunk_occur_time = self._normalize_datetime(parsed_df.iloc[0]["code_occur_time"]) if chunk_occur_time is None: raise ValueError("code_occur_time is required") @@ -152,12 +157,6 @@ def _run_window( for target_date in sorted(inserted_target_dates): self.repository.analyze_target_partition(target_date, connection=connection) - self.repository.upsert_die_yield_daily_summary(target_date, connection=connection) - self.repository.upsert_root_cause_daily_summary(target_date, connection=connection) - print( - "[ER_DOSE] " - f"summary updated (die_yield, root_cause) partition_date={target_date}" - ) print( "[ER_DOSE] " @@ -277,31 +276,18 @@ def run( start_time: datetime | None = None, end_time: datetime | None = None, chunk_size: int = 10000, - lookback_days: int = 2, - reference_date: date | None = None, target_date: date | None = None, ) -> None: + if chunk_size <= 0: + raise ValueError("chunk_size must be greater than 0") if target_date is not None: - self.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=chunk_size, - ) + self._reload_target_date(target_date=target_date, chunk_size=chunk_size) return - if start_time is None and end_time is None: - self.run_recent_days( - lookback_days=lookback_days, - reference_date=reference_date, - chunk_size=chunk_size, - ) - return if start_time is None or end_time is None: raise ValueError("start_time and end_time are required") if start_time >= end_time: raise ValueError("start_time must be earlier than end_time") - if chunk_size <= 0: - raise ValueError("chunk_size must be greater than 0") self._run_window(start_time=start_time, end_time=end_time, chunk_size=chunk_size) @@ -372,11 +358,6 @@ def _run_window( for target_date in sorted(inserted_target_dates): self.repository.analyze_target_partition(target_date, connection=connection) - self.repository.upsert_root_cause_daily_summary(target_date, connection=connection) - print( - "[ER_DOSE_EUV] " - f"summary updated (root_cause) partition_date={target_date}" - ) print( "[ER_DOSE_EUV] " diff --git a/er_dose/raw/raw_repository.py b/er_dose/raw/raw_repository.py index 92a2758..65f2c0f 100644 --- a/er_dose/raw/raw_repository.py +++ b/er_dose/raw/raw_repository.py @@ -1,10 +1,11 @@ from __future__ import annotations from datetime import date, datetime, timedelta -from typing import Iterator +from typing import Any, Iterator import pandas as pd +from er_dose.common.batch_log_repository import insert_batch_log from er_dose.common.sql_filters import active_nxe_eq_filter from er_dose.infra.postgres_db import PostgresDB @@ -31,6 +32,23 @@ class ERDoseRepository: def __init__(self, db: PostgresDB): self.db = db + def insert_batch_log( + self, + batch_name: str, + target_date: date, + event_type: str, + message: str, + data: dict[str, Any], + ) -> None: + insert_batch_log( + self.db, + batch_name=batch_name, + target_date=target_date, + event_type=event_type, + message=message, + data=data, + ) + def fetch_raw_logs_in_chunks( self, start_time: datetime, @@ -119,6 +137,56 @@ def fetch_target_count(self, target_date: date) -> int: return 0 return int(df.iloc[0]["row_count"]) + def fetch_equipment_counts(self, start_time: datetime, end_time: datetime) -> dict[str, Any]: + target_codes_sql = ", ".join(f"'{code}'" for code in TARGET_CODES) + query = f""" + with source_counts as ( + select + r.eq_name, + count(*) as source_count + from {MAIN_RAW_TABLE} r + where r.code_occur_time >= :start_time + and r.code_occur_time < :end_time + and r.code in ({target_codes_sql}) + and {active_nxe_eq_filter("r.eq_name")} + group by r.eq_name + ), + target_counts as ( + select + p.eq_name, + count(*) as target_count + from {PARSED_TABLE} p + where p.code_occur_time >= :start_time + and p.code_occur_time < :end_time + and p.code in ({target_codes_sql}) + and {active_nxe_eq_filter("p.eq_name")} + group by p.eq_name + ) + select + coalesce(s.eq_name, t.eq_name) as eq_name, + coalesce(s.source_count, 0) as source_count, + coalesce(t.target_count, 0) as target_count, + sum(coalesce(s.source_count, 0)) over () as total_source_count, + sum(coalesce(t.target_count, 0)) over () as total_target_count, + min(case when coalesce(s.source_count, 0) = coalesce(t.target_count, 0) + then 1 else 0 end) over () as matched + from source_counts s + full outer join target_counts t on t.eq_name = s.eq_name + order by eq_name + """ + df = self.db.select(query, params={"start_time": start_time, "end_time": end_time}) + if df is None or df.empty: + return {"source_count": 0, "target_count": 0, "matched": True, "equipment_counts": []} + return { + "source_count": int(df.iloc[0]["total_source_count"]), + "target_count": int(df.iloc[0]["total_target_count"]), + "matched": bool(df.iloc[0]["matched"]), + "equipment_counts": [ + {"eq_name": row.eq_name, "source_count": int(row.source_count), "target_count": int(row.target_count)} + for row in df.itertuples(index=False) + ], + } + def truncate_target_partition(self, target_date: date, connection=None) -> int: parsed_table = self._partition_table_name(PARSED_TABLE, target_date) return self.db.execute(f"truncate table {parsed_table}", connection=connection) @@ -232,13 +300,16 @@ def analyze_target_partition(self, target_date: str, connection=None) -> int: partition_table = f"{PARSED_TABLE}_1_prt_p{target_date.replace('-', '')}" return self.db.execute(f"ANALYZE {partition_table}", connection=connection) - def upsert_die_yield_daily_summary(self, target_date: date | str, connection=None) -> int: - if isinstance(target_date, str): - target_date = datetime.strptime(target_date, "%Y-%m-%d").date() - start_time = datetime.combine(target_date, datetime.min.time()) - end_time = start_time + timedelta(days=1) + def delete_die_yield_daily_summary(self, start_time: datetime, end_time: datetime) -> None: + delete_query = """ + delete from prism_common.de_trend_die_yield_daily + where occur_date >= cast(:start_time as date) + and occur_date < :end_time + """ + self.db.execute(delete_query, params={"start_time": start_time, "end_time": end_time}) - query = f""" + def insert_die_yield_daily_summary(self, start_time: datetime, end_time: datetime) -> None: + insert_query = f""" insert into prism_common.de_trend_die_yield_daily ( occur_date, eq_name, @@ -307,26 +378,20 @@ def upsert_die_yield_daily_summary(self, target_date: date | str, connection=Non sum(case when valid_wafer_yn = 1 then reject_yn else 0 end) as reject_wafer, now() as created_at from wafer_code_data - group by occur_date, eq_name - on conflict (occur_date, eq_name) - do update set - total_die = excluded.total_die, - reject_shot = excluded.reject_shot, - to_repair_die = excluded.to_repair_die, - repair_nok = excluded.repair_nok, - total_wafer = excluded.total_wafer, - reject_wafer = excluded.reject_wafer, - created_at = now(); + group by occur_date, eq_name; """ - return self.db.execute(query, params={"start_time": start_time, "end_time": end_time}, connection=connection) + self.db.execute(insert_query, params={"start_time": start_time, "end_time": end_time}) - def upsert_root_cause_daily_summary(self, target_date: date | str, connection=None) -> int: - if isinstance(target_date, str): - target_date = datetime.strptime(target_date, "%Y-%m-%d").date() - start_time = datetime.combine(target_date, datetime.min.time()) - end_time = start_time + timedelta(days=1) + def delete_root_cause_daily_summary(self, start_time: datetime, end_time: datetime) -> None: + delete_query = """ + delete from prism_common.de_trend_root_cause_daily + where occur_date >= cast(:start_time as date) + and occur_date < :end_time + """ + self.db.execute(delete_query, params={"start_time": start_time, "end_time": end_time}) - query = """ + def insert_root_cause_daily_summary(self, start_time: datetime, end_time: datetime) -> None: + insert_query = """ insert into prism_common.de_trend_root_cause_daily ( occur_date, eq_name, @@ -345,13 +410,9 @@ def upsert_root_cause_daily_summary(self, target_date: date | str, connection=No and p.code_occur_time < :end_time and p.eq_name is not null and p.root_cause is not null - group by p.code_occur_time::date, p.eq_name, p.root_cause - on conflict (occur_date, eq_name, root_cause) - do update set - frequency = excluded.frequency, - created_at = now(); + group by p.code_occur_time::date, p.eq_name, p.root_cause; """ - return self.db.execute(query, params={"start_time": start_time, "end_time": end_time}, connection=connection) + self.db.execute(insert_query, params={"start_time": start_time, "end_time": end_time}) def transaction(self): return self.db.transaction() diff --git a/er_dose/sql/check_er_dose_counts_last_7_days.sql b/er_dose/sql/check_er_dose_counts_last_7_days.sql new file mode 100644 index 0000000..f8173b0 --- /dev/null +++ b/er_dose/sql/check_er_dose_counts_last_7_days.sql @@ -0,0 +1,128 @@ +-- [RAW] 날짜별 원천 vs 타깃 건수 비교 (오늘 포함 최근 7일) +with dates as ( + select current_date - day_offset as target_date + from generate_series(0, 6) as days(day_offset) +), +source_counts as ( + select + r.code_occur_time::date as target_date, + count(*) as source_count + from mbeat.er_data_raw r + where r.code_occur_time >= (current_date - 6)::timestamp + and r.code_occur_time < (current_date + 1)::timestamp + and r.code in ( + 'DW-3411', + 'DW-3425', + 'DW-343A', + 'DW-343B', + 'LO-0050', + 'LO-0051', + 'LO-0052', + 'LO-0061', + 'LO-8166', + 'LO-8167', + 'KE-9103', + 'KE-9104' + ) + and r.eq_name in ( + select eqp.eqp_id + from prism_dev.photo_eqp_info eqp + where eqp.use_yn = 'Y' + and eqp.eqp_model_name like 'NXE%' + ) + group by r.code_occur_time::date +), +target_counts as ( + select + p.code_occur_time::date as target_date, + count(*) as target_count + from prism_common.er_dose_raw_parsed p + where p.code_occur_time >= (current_date - 6)::timestamp + and p.code_occur_time < (current_date + 1)::timestamp + and p.code in ( + 'DW-3411', + 'DW-3425', + 'DW-343A', + 'DW-343B', + 'LO-0050', + 'LO-0051', + 'LO-0052', + 'LO-0061', + 'LO-8166', + 'LO-8167', + 'KE-9103', + 'KE-9104' + ) + and p.eq_name in ( + select eqp.eqp_id + from prism_dev.photo_eqp_info eqp + where eqp.use_yn = 'Y' + and eqp.eqp_model_name like 'NXE%' + ) + group by p.code_occur_time::date +) +select + d.target_date, + coalesce(s.source_count, 0) as source_count, + coalesce(t.target_count, 0) as target_count, + coalesce(s.source_count, 0) - coalesce(t.target_count, 0) as count_diff, + case + when coalesce(s.source_count, 0) = coalesce(t.target_count, 0) then 'Y' + else 'N' + end as matched +from dates d +left join source_counts s on s.target_date = d.target_date +left join target_counts t on t.target_date = d.target_date +order by d.target_date; + + +-- [EUV] 날짜별 원천 vs 타깃 건수 비교 (오늘 포함 최근 7일) +with dates as ( + select current_date - day_offset as target_date + from generate_series(0, 6) as days(day_offset) +), +source_counts as ( + select + r.code_occur_time::date as target_date, + count(*) as source_count + from mbeat.er_data_raw_euv r + where r.code_occur_time >= (current_date - 6)::timestamp + and r.code_occur_time < (current_date + 1)::timestamp + and lower(r.contents) like '%dose error detected in file:%' + and lower(r.contents) like '%root cause%' + and r.eq_name in ( + select eqp.eqp_id + from prism_dev.photo_eqp_info eqp + where eqp.use_yn = 'Y' + and eqp.eqp_model_name like 'NXE%' + ) + group by r.code_occur_time::date +), +target_counts as ( + select + p.code_occur_time::date as target_date, + count(*) as target_count + from prism_common.er_dose_euv_parsed p + where p.code_occur_time >= (current_date - 6)::timestamp + and p.code_occur_time < (current_date + 1)::timestamp + and p.eq_name in ( + select eqp.eqp_id + from prism_dev.photo_eqp_info eqp + where eqp.use_yn = 'Y' + and eqp.eqp_model_name like 'NXE%' + ) + group by p.code_occur_time::date +) +select + d.target_date, + coalesce(s.source_count, 0) as source_count, + coalesce(t.target_count, 0) as target_count, + coalesce(s.source_count, 0) - coalesce(t.target_count, 0) as count_diff, + case + when coalesce(s.source_count, 0) = coalesce(t.target_count, 0) then 'Y' + else 'N' + end as matched +from dates d +left join source_counts s on s.target_date = d.target_date +left join target_counts t on t.target_date = d.target_date +order by d.target_date; diff --git a/er_dose/sql/create_batch_event_log.sql b/er_dose/sql/create_batch_event_log.sql new file mode 100644 index 0000000..ffe2b79 --- /dev/null +++ b/er_dose/sql/create_batch_event_log.sql @@ -0,0 +1,24 @@ +create schema if not exists mbeat; + +create table if not exists mbeat.batch_event_log ( + id bigserial primary key, + batch_name varchar(100) not null, + target_date date, + event_type varchar(50) not null, + message text, + data jsonb not null default '{}'::jsonb, + created_at timestamp default now() not null +); + +do $$ +begin + if not exists ( + select 1 from pg_catalog.pg_indexes + where schemaname = 'mbeat' + and indexname = 'idx_batch_event_log_batch_date' + ) then + create index idx_batch_event_log_batch_date + on mbeat.batch_event_log (batch_name, target_date, created_at desc); + end if; +end +$$; diff --git a/er_dose/sql/create_er_dose_euv_parsed.sql b/er_dose/sql/create_er_dose_euv_parsed.sql index 1ed5baa..b0bd60c 100644 --- a/er_dose/sql/create_er_dose_euv_parsed.sql +++ b/er_dose/sql/create_er_dose_euv_parsed.sql @@ -15,39 +15,39 @@ create table if not exists prism_common.er_dose_euv_parsed ( dose_error_detected_in_file text, root_cause_code text, root_cause text, - exposure_length numeric(12,7), - duty_cycle numeric(12,7), - min_dose_error numeric(12,7), - max_dose_error numeric(12,7), - on_drop_euv_energy numeric(12,7), - on_drop_pp_energy numeric(12,7), - on_drop_mp_energy numeric(12,7), - on_drop_pp_dlgc_1 numeric(12,7), - on_drop_mp_dlgc_1 numeric(12,7), - bi_cell_y_3sigma numeric(12,7), - fdsc_y_error numeric(12,7), - fdsc_y_3sigma numeric(12,7), - max_cross_interval numeric(12,7), - xint_3sigma numeric(12,7), - euv_3sigma numeric(12,7), + exposure_length numeric(20,10), + duty_cycle numeric(20,10), + min_dose_error numeric(20,10), + max_dose_error numeric(20,10), + on_drop_euv_energy numeric(20,10), + on_drop_pp_energy numeric(20,10), + on_drop_mp_energy numeric(20,10), + on_drop_pp_dlgc_1 numeric(20,10), + on_drop_mp_dlgc_1 numeric(20,10), + bi_cell_y_3sigma numeric(20,10), + fdsc_y_error numeric(20,10), + fdsc_y_3sigma numeric(20,10), + max_cross_interval numeric(20,10), + xint_3sigma numeric(20,10), + euv_3sigma numeric(20,10), pulses_euv_0_6dt_tot integer, fed_pulses integer, - l2dx_maxce numeric(12,7), - l2dy_maxce numeric(12,7), - sensitivity_at_l2dx_maxce numeric(12,7), - sensitivity_at_l2dy_maxce numeric(12,7), - dose_margin numeric(12,7), - l2dx_qc_etdc_3sigma numeric(12,7), - l2dx_qc_etdc_median numeric(12,7), - l2dy_qc_etdc_3sigma numeric(12,7), - l2dy_qc_etdc_median numeric(12,7), - rbdy_peak_frequency_hf numeric(12,7), - rbdy_peak_frequency_lf numeric(12,7), - rbdy_peak_frequency_mf numeric(12,7), - rbdy_peak_power_hf numeric(12,7), - rbdy_qc_etdc_3sigma numeric(12,7), - rbdy_total_power_lf numeric(12,7), - rbdy_total_power_mf numeric(12,7), + l2dx_maxce numeric(20,10), + l2dy_maxce numeric(20,10), + sensitivity_at_l2dx_maxce numeric(20,10), + sensitivity_at_l2dy_maxce numeric(20,10), + dose_margin numeric(20,10), + l2dx_qc_etdc_3sigma numeric(20,10), + l2dx_qc_etdc_median numeric(20,10), + l2dy_qc_etdc_3sigma numeric(20,10), + l2dy_qc_etdc_median numeric(20,10), + rbdy_peak_frequency_hf numeric(20,10), + rbdy_peak_frequency_lf numeric(20,10), + rbdy_peak_frequency_mf numeric(20,10), + rbdy_peak_power_hf numeric(20,10), + rbdy_qc_etdc_3sigma numeric(20,10), + rbdy_total_power_lf numeric(20,10), + rbdy_total_power_mf numeric(20,10), software_version text, created_at timestamp default now() ) diff --git a/er_dose/sql/migrate_er_dose_euv_parsed_numeric_precision.sql b/er_dose/sql/migrate_er_dose_euv_parsed_numeric_precision.sql new file mode 100644 index 0000000..4182602 --- /dev/null +++ b/er_dose/sql/migrate_er_dose_euv_parsed_numeric_precision.sql @@ -0,0 +1,34 @@ +-- Run during a maintenance window. Increasing the scale can rewrite every +-- attached partition and hold ACCESS EXCLUSIVE locks until this statement ends. +alter table if exists prism_common.er_dose_euv_parsed + alter column exposure_length type numeric(20,10), + alter column duty_cycle type numeric(20,10), + alter column min_dose_error type numeric(20,10), + alter column max_dose_error type numeric(20,10), + alter column on_drop_euv_energy type numeric(20,10), + alter column on_drop_pp_energy type numeric(20,10), + alter column on_drop_mp_energy type numeric(20,10), + alter column on_drop_pp_dlgc_1 type numeric(20,10), + alter column on_drop_mp_dlgc_1 type numeric(20,10), + alter column bi_cell_y_3sigma type numeric(20,10), + alter column fdsc_y_error type numeric(20,10), + alter column fdsc_y_3sigma type numeric(20,10), + alter column max_cross_interval type numeric(20,10), + alter column xint_3sigma type numeric(20,10), + alter column euv_3sigma type numeric(20,10), + alter column l2dx_maxce type numeric(20,10), + alter column l2dy_maxce type numeric(20,10), + alter column sensitivity_at_l2dx_maxce type numeric(20,10), + alter column sensitivity_at_l2dy_maxce type numeric(20,10), + alter column dose_margin type numeric(20,10), + alter column l2dx_qc_etdc_3sigma type numeric(20,10), + alter column l2dx_qc_etdc_median type numeric(20,10), + alter column l2dy_qc_etdc_3sigma type numeric(20,10), + alter column l2dy_qc_etdc_median type numeric(20,10), + alter column rbdy_peak_frequency_hf type numeric(20,10), + alter column rbdy_peak_frequency_lf type numeric(20,10), + alter column rbdy_peak_frequency_mf type numeric(20,10), + alter column rbdy_peak_power_hf type numeric(20,10), + alter column rbdy_qc_etdc_3sigma type numeric(20,10), + alter column rbdy_total_power_lf type numeric(20,10), + alter column rbdy_total_power_mf type numeric(20,10); diff --git a/tests/test_er_dose_airflow.py b/tests/test_er_dose_airflow.py new file mode 100644 index 0000000..af35379 --- /dev/null +++ b/tests/test_er_dose_airflow.py @@ -0,0 +1,180 @@ +from datetime import date, datetime +from unittest.mock import Mock + +import pandas as pd +import pytest + +from airflow_modules import er_dose_jobs as jobs +from er_dose.euv.euv_processor import ERDoseEUVProcessor +from er_dose.euv.euv_repository import ERDoseEUVRepository +from er_dose.raw.raw_processor import ERDoseProcessor, ERDoseEUVProcessor as LegacyEUVProcessor +from er_dose.raw.raw_repository import ERDoseRepository +from tests.test_er_dose_processor import FakeDB +from tests.test_er_dose_euv_processor import FakeDB as EUVFakeDB + + +def saved_plan(raw=True, euv=True): + return [{"target_date": "2026-05-12", "start_time": "2026-05-12T00:00:00", + "end_time": "2026-05-13T00:00:00", "raw": raw, "euv": euv}] + + +@pytest.mark.parametrize("source,target,distinct,expected", [ + (0, 0, None, False), (2, 2, None, False), (2, 1, 1, False), + (3, 1, 2, True), (1, 0, None, True), (0, 1, 0, True), +]) +def test_planner_preserves_count_and_distinct_rules(source, target, distinct, expected): + repo = Mock() + repo.fetch_source_count.side_effect = [source, distinct] + repo.fetch_target_count.return_value = target + assert jobs.needs_reload(repo, date(2026, 5, 12)) is expected + assert repo.fetch_source_count.call_count == (2 if source != target and target > 0 else 1) + if repo.fetch_source_count.call_count == 2: + assert repo.fetch_source_count.call_args.kwargs == {"distinct": True} + + +def test_fourteenth_run_keeps_twelfth_in_json_safe_plan(monkeypatch): + monkeypatch.setattr(jobs, "PostgresDB", Mock()) + raw, euv = Mock(), Mock() + monkeypatch.setattr(jobs, "ERDoseRepository", Mock(return_value=raw)) + monkeypatch.setattr(jobs, "ERDoseEUVRepository", Mock(return_value=euv)) + raw.fetch_source_count.side_effect = lambda day, **kw: int(day == date(2026, 5, 12)) + raw.fetch_target_count.return_value = 0 + euv.fetch_source_count.return_value = euv.fetch_target_count.return_value = 0 + plan = jobs.plan_dates("2026-05-14") + assert len(plan) == 10 + assert plan[0]["target_date"] == "2026-05-05" + assert plan[-1]["target_date"] == "2026-05-14" + assert [day for day in plan if day["raw"]] == saved_plan(raw=True, euv=False) + import json + assert json.loads(json.dumps(plan)) == plan + raw.truncate_target_partition.assert_not_called() + + +@pytest.mark.parametrize("processor_type,repository_type,db_type", [ + (ERDoseProcessor, ERDoseRepository, FakeDB), + (ERDoseEUVProcessor, ERDoseEUVRepository, EUVFakeDB), + (LegacyEUVProcessor, ERDoseEUVRepository, EUVFakeDB), +]) +def test_parser_requires_explicit_period_and_date_retry_does_not_recheck_counts(processor_type, repository_type, db_type): + db = db_type(pd.DataFrame()) + repo = repository_type(db) + repo.fetch_source_count = Mock(side_effect=AssertionError("planner only")) + repo.fetch_target_count = Mock(side_effect=AssertionError("planner only")) + processor = processor_type(repo) + with pytest.raises(ValueError, match="start_time and end_time"): + processor.run() + for _ in range(2): + processor.run(target_date=date(2026, 5, 12)) + assert len([q for q, _, _ in db.executed if q.strip().lower().startswith("truncate")]) == 2 + assert not any("de_trend_" in q or "batch_event_log" in q for q, _, _ in db.executed) + + + + +@pytest.mark.parametrize("job_type,method", [("die_yield", "die_yield_daily_summary"), ("root_cause", "root_cause_daily_summary")]) +def test_report_retry_restarts_delete_without_reparsing(monkeypatch, job_type, method): + db_factory = Mock() + monkeypatch.setattr(jobs, "PostgresDB", db_factory) + repo = Mock() + monkeypatch.setattr(jobs, "ERDoseRepository", Mock(return_value=repo)) + events = [] + getattr(repo, "delete_" + method).side_effect = lambda *a: events.append(("delete", a)) + def insert(*args): + events.append(("insert", args)) + if len(events) == 2: + raise RuntimeError("admin terminated connection") + getattr(repo, "insert_" + method).side_effect = insert + ti = Mock() + ti.xcom_pull.return_value = saved_plan() + with pytest.raises(RuntimeError): + jobs.report_day(job_type, 0, ti) + jobs.report_day(job_type, 0, ti) + assert [action for action, _ in events] == ["delete", "insert", "delete", "insert"] + assert all(bounds == (datetime(2026, 5, 12), datetime(2026, 5, 13)) for _, bounds in events) + assert db_factory.call_count == 2 + repo.transaction.assert_not_called() + repo.fetch_source_count.assert_not_called() + + +def test_failed_delete_does_not_insert(monkeypatch): + monkeypatch.setattr(jobs, "PostgresDB", Mock()) + repo = Mock() + repo.delete_die_yield_daily_summary.side_effect = RuntimeError("delete failed") + monkeypatch.setattr(jobs, "ERDoseRepository", Mock(return_value=repo)) + ti = Mock() + ti.xcom_pull.return_value = saved_plan() + with pytest.raises(RuntimeError): + jobs.report_day("die_yield", 0, ti) + repo.insert_die_yield_daily_summary.assert_not_called() + + +@pytest.mark.parametrize("raw,euv", [(False, False), (True, False), (False, True), (True, True)]) +def test_only_dag_branches_decide_which_tasks_run(monkeypatch, raw, euv): + db = Mock(side_effect=AssertionError("unchanged day")) + monkeypatch.setattr(jobs, "PostgresDB", db) + ti = Mock() + ti.xcom_pull.return_value = saved_plan(raw, euv) + for parser, selected in (("raw", raw), ("euv", euv)): + expected = f"day_00_parse_{parser}" if selected else f"day_00_{parser}_done" + assert jobs.select_parse_task(parser, 0, ti) == expected + assert jobs.select_report_tasks(0, ti) == ( + [f"day_00_{job}" for job in ("die_yield", "root_cause", "raw_count", "euv_count")] + if raw or euv else [] + ) + + +def test_called_report_executes_without_checking_selection_flags(monkeypatch): + monkeypatch.setattr(jobs, "PostgresDB", Mock()) + repo = Mock() + monkeypatch.setattr(jobs, "ERDoseRepository", Mock(return_value=repo)) + ti = Mock() + # Execution functions need only the saved period, not the branch's flags. + window = saved_plan()[0] + del window["raw"], window["euv"] + ti.xcom_pull.return_value = [window] + jobs.report_day("die_yield", 0, ti) + repo.insert_die_yield_daily_summary.assert_called_once_with(datetime(2026, 5, 12), datetime(2026, 5, 13)) + + +def test_missing_plan_fails_instead_of_recomputing(monkeypatch): + monkeypatch.setattr(jobs, "PostgresDB", Mock(side_effect=AssertionError("do not replan"))) + ti = Mock() + ti.xcom_pull.return_value = None + with pytest.raises(ValueError, match="plan is missing"): + jobs.report_day("raw_count", 0, ti) + + +def test_statistics_retry_does_not_duplicate_committed_insert(monkeypatch): + records = {} + db = Mock() + def execute(query, params): + key = tuple(params[k] for k in ("batch_name", "target_date", "event_type", "message")) + if query.strip().lower().startswith("delete"): + records.pop(key, None) + else: + assert key not in records + records[key] = params["data"] + if execute.fail: + execute.fail = False + raise RuntimeError("worker lost after commit") + execute.fail = True + db.execute.side_effect = execute + db.select.return_value = pd.DataFrame() + monkeypatch.setattr(jobs, "PostgresDB", Mock(return_value=db)) + ti = Mock() + ti.xcom_pull.return_value = saved_plan() + with pytest.raises(RuntimeError): + jobs.report_day("raw_count", 0, ti) + jobs.report_day("raw_count", 0, ti) + jobs.report_day("euv_count", 0, ti) + assert len(records) == 2 + assert {key[0] for key in records} == {"ER_DOSE_RAW", "ER_DOSE_EUV"} + + +def test_statistics_query_failure_keeps_existing_snapshot(monkeypatch): + repo = Mock() + repo.fetch_equipment_counts.side_effect = RuntimeError("select failed") + with pytest.raises(RuntimeError): + jobs.write_equipment_count_log(repo, "ER_DOSE_RAW", datetime(2026, 5, 12), datetime(2026, 5, 13)) + repo.db.execute.assert_not_called() + repo.insert_batch_log.assert_not_called() diff --git a/tests/test_er_dose_dag.py b/tests/test_er_dose_dag.py new file mode 100644 index 0000000..f4218f6 --- /dev/null +++ b/tests/test_er_dose_dag.py @@ -0,0 +1,85 @@ +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import MagicMock, Mock + +import pytest + + +pytest.importorskip("airflow") +pytest.importorskip("airflow.providers.ssh") + +from airflow.exceptions import AirflowException +from airflow.models import DagBag +from airflow.operators.python import PythonOperator +from airflow.providers.ssh.operators.ssh import SSHOperator +from airflow.utils.dag_cycle_tester import check_cycle + + +def test_dag_separates_report_retries_from_chronological_parsing(): + path = Path(__file__).resolve().parents[1] / "dags" / "er_dose_daily_dag.py" + bag = DagBag(dag_folder=str(path), include_examples=False) + assert not bag.import_errors + dag = bag.dags["er_dose_daily"] + check_cycle(dag) + assert len(dag.tasks) == 111 + assert dag.max_active_runs == 1 + previous = "plan_dates" + for index in range(10): + for parser in ("euv", "raw"): + branch = dag.get_task(f"day_{index:02d}_choose_{parser}") + parse = dag.get_task(f"day_{index:02d}_parse_{parser}") + done = dag.get_task(f"day_{index:02d}_{parser}_done") + assert branch.upstream_task_ids == {previous} + assert parse.upstream_task_ids == {branch.task_id} + assert isinstance(parse, SSHOperator) + assert parse.ssh_conn_id == "er_dose_parser" + assert parse.cmd_timeout is None + assert parse.get_pty is True + assert parse.do_xcom_push is False + assert done.upstream_task_ids == {branch.task_id, parse.task_id} + assert done.trigger_rule == "none_failed_min_one_success" + previous = done.task_id + report_branch = dag.get_task(f"day_{index:02d}_choose_reports") + assert report_branch.upstream_task_ids == {previous} + for job in ("die_yield", "root_cause", "raw_count", "euv_count"): + report = dag.get_task(f"day_{index:02d}_{job}") + assert report.upstream_task_ids == {report_branch.task_id} + assert report.downstream_task_ids == set() + assert report.retries == 3 + assert isinstance(report, PythonOperator) + + +@pytest.mark.parametrize("parser", ["raw", "euv"]) +def test_ssh_retry_keeps_saved_date_and_propagates_failure(parser): + path = Path(__file__).resolve().parents[1] / "dags" / "er_dose_daily_dag.py" + command_prefix = 'cd "/opt/parser project" && exec /opt/venv/bin/python -m er_dose.run_er_dose_batch' + ti = SimpleNamespace(xcom_pull=lambda **kwargs: [{"target_date": "2026-05-12"}], xcom_push=Mock()) + context = {"ti": ti, "var": {"value": SimpleNamespace(er_dose_parser_command=command_prefix)}} + commands = [] + for exit_status in (1, 0): + dag = DagBag(dag_folder=str(path), include_examples=False).dags["er_dose_daily"] + task = dag.get_task(f"day_00_parse_{parser}") + task.render_template_fields(context) + task.get_ssh_client = MagicMock() + task.ssh_hook = Mock() + task.ssh_hook.exec_ssh_client_command.return_value = (exit_status, b"done", b"") + if exit_status: + with pytest.raises(AirflowException, match="exit status = 1"): + task.execute(context) + else: + task.execute(context) + commands.append(task.ssh_hook.exec_ssh_client_command.call_args.args[1]) + expected = command_prefix + f" --parser ER_DOSE_{parser.upper()} --date 2026-05-12 --chunk-size 30000" + assert commands == [expected, expected] + ti.xcom_push.assert_not_called() + + +def test_ssh_missing_plan_fails_before_connecting(): + path = Path(__file__).resolve().parents[1] / "dags" / "er_dose_daily_dag.py" + task = DagBag(dag_folder=str(path), include_examples=False).dags["er_dose_daily"].get_task("day_00_parse_raw") + ti = SimpleNamespace(xcom_pull=lambda **kwargs: None) + task.get_ssh_client = Mock() + from jinja2 import UndefinedError + with pytest.raises(UndefinedError): + task.render_template_fields({"ti": ti, "var": {"value": SimpleNamespace(er_dose_parser_command="python -m er_dose.run_er_dose_batch")}}) + task.get_ssh_client.assert_not_called() diff --git a/tests/test_er_dose_euv_numeric_precision.py b/tests/test_er_dose_euv_numeric_precision.py new file mode 100644 index 0000000..aeacab9 --- /dev/null +++ b/tests/test_er_dose_euv_numeric_precision.py @@ -0,0 +1,68 @@ +from pathlib import Path + + +SQL_DIR = Path(__file__).resolve().parents[1] / "er_dose" / "sql" +DDL_PATH = SQL_DIR / "create_er_dose_euv_parsed.sql" +MIGRATION_PATH = SQL_DIR / "migrate_er_dose_euv_parsed_numeric_precision.sql" + +NUMERIC_COLUMNS = { + "exposure_length", + "duty_cycle", + "min_dose_error", + "max_dose_error", + "on_drop_euv_energy", + "on_drop_pp_energy", + "on_drop_mp_energy", + "on_drop_pp_dlgc_1", + "on_drop_mp_dlgc_1", + "bi_cell_y_3sigma", + "fdsc_y_error", + "fdsc_y_3sigma", + "max_cross_interval", + "xint_3sigma", + "euv_3sigma", + "l2dx_maxce", + "l2dy_maxce", + "sensitivity_at_l2dx_maxce", + "sensitivity_at_l2dy_maxce", + "dose_margin", + "l2dx_qc_etdc_3sigma", + "l2dx_qc_etdc_median", + "l2dy_qc_etdc_3sigma", + "l2dy_qc_etdc_median", + "rbdy_peak_frequency_hf", + "rbdy_peak_frequency_lf", + "rbdy_peak_frequency_mf", + "rbdy_peak_power_hf", + "rbdy_qc_etdc_3sigma", + "rbdy_total_power_lf", + "rbdy_total_power_mf", +} + + +def test_euv_metric_columns_use_numeric_20_10(): + ddl = DDL_PATH.read_text(encoding="utf-8").lower() + columns = { + line.strip().split()[0] + for line in ddl.splitlines() + if "numeric(20,10)" in line + } + + assert columns == NUMERIC_COLUMNS + assert "numeric(12,7)" not in ddl + + +def test_numeric_precision_migration_alters_parent_once(): + sql = MIGRATION_PATH.read_text(encoding="utf-8").lower() + alter_lines = [ + line.strip() + for line in sql.splitlines() + if line.strip().startswith("alter column") + ] + columns = {line.split()[2] for line in alter_lines} + + assert "alter table if exists prism_common.er_dose_euv_parsed" in sql + assert "alter table only" not in sql + assert columns == NUMERIC_COLUMNS + assert len(alter_lines) == len(NUMERIC_COLUMNS) + assert all("type numeric(20,10)" in line for line in alter_lines) diff --git a/tests/test_er_dose_euv_processor.py b/tests/test_er_dose_euv_processor.py index acf1805..9d41444 100644 --- a/tests/test_er_dose_euv_processor.py +++ b/tests/test_er_dose_euv_processor.py @@ -1,8 +1,9 @@ from __future__ import annotations +import json import unittest from contextlib import redirect_stdout -from datetime import date, datetime +from datetime import date, datetime, timedelta from io import StringIO import pandas as pd @@ -24,11 +25,19 @@ def __exit__(self, exc_type, exc, traceback): class FakeDB: - def __init__(self, raw_df, source_counts=None, target_counts=None, distinct_source_counts=None): + def __init__( + self, + raw_df, + source_counts=None, + target_counts=None, + distinct_source_counts=None, + equipment_counts=None, + ): self.raw_df = raw_df self.source_counts = source_counts or {} self.target_counts = target_counts or {} self.distinct_source_counts = distinct_source_counts or {} + self.equipment_counts = equipment_counts or {} self.fetch_query = None self.fetch_params = None self.executed = [] @@ -37,10 +46,25 @@ def __init__(self, raw_df, source_counts=None, target_counts=None, distinct_sour self.connection = object() self.partition_inserts = [] - def select(self, query, params=None): + def select(self, query, params=None, connection=None): self.fetch_query = query self.fetch_params = params lowered = query.lower() + if "with source_counts as" in lowered and "full outer join target_counts" in lowered: + target_date = params["start_time"].date() + rows = self.equipment_counts.get(target_date) + if rows is None: + source_count = self.source_counts.get(target_date, 0) + target_count = self.target_counts.get(target_date, 0) + rows = ([{"eq_name": "EQ1", "source_count": source_count, "target_count": target_count}] + if source_count or target_count else []) + return pd.DataFrame([ + {**row, + "total_source_count": sum(item["source_count"] for item in rows), + "total_target_count": sum(item["target_count"] for item in rows), + "matched": int(all(item["source_count"] == item["target_count"] for item in rows))} + for row in rows + ]) if "count(distinct" in lowered: target_date = params["start_time"].date() row_count = self.distinct_source_counts.get(target_date, self.source_counts.get(target_date, 0)) @@ -182,6 +206,10 @@ def test_run_inserts_parsed_root_cause_rows(self): self.assertEqual([option["analyze"] for option in db.copy_options], [False]) analyze_queries = [query for query, _, _ in db.executed if query == "ANALYZE prism_common.er_dose_euv_parsed_1_prt_p20260504"] self.assertEqual(len(analyze_queries), 1) + summary_queries = [ + item for item in db.executed if "de_trend_root_cause_daily" in item[0].lower() + ] + self.assertEqual(summary_queries, []) self.assertNotIn("er_line", inserted_df.columns) self.assertNotIn("belong", inserted_df.columns) self.assertNotIn("type", inserted_df.columns) @@ -215,44 +243,37 @@ def test_fetch_counts_filter_active_nxe_eq_names_and_root_cause_source(self): repo.fetch_source_count(target_date, distinct=True) self.assertIn("count(distinct (r.eq_name, r.code, r.code_occur_time))", db.fetch_query) - def test_run_recent_days_skips_when_counts_match(self): + equipment_counts = repo.fetch_equipment_counts(datetime.combine(target_date, datetime.min.time()), datetime.combine(target_date + timedelta(days=1), datetime.min.time())) + self.assertEqual(equipment_counts, {"source_count": 1, "target_count": 1, "matched": True, "equipment_counts": [{"eq_name": "EQ1", "source_count": 1, "target_count": 1}]}) + self.assertIn("group by r.eq_name", db.fetch_query) + self.assertIn("group by p.eq_name", db.fetch_query) + self.assertIn("full outer join target_counts", db.fetch_query) + + + def test_euv_run_does_not_write_summary_logs(self): target_date = date(2026, 5, 4) - raw_df = pd.DataFrame( - [ - { - "eq_name": "EQ1", - "er_type": "EUV", - "code": "CODE1", - "code_occur_time": datetime(2026, 5, 4, 18, 5, 29), - "title": "Dose Error Root Cause", - "contents": SAMPLE_EUV_CONTENTS, - "reason_code": "R1", - "task": "TASK1", - "compile_script": "SCRIPT1", - }, - ] - ) db = FakeDB( - raw_df, - source_counts={target_date: 1}, - target_counts={target_date: 1}, + pd.DataFrame(), + source_counts={target_date: 7}, + target_counts={target_date: 7}, + equipment_counts={ + target_date: [ + {"eq_name": "EQ1", "source_count": 7, "target_count": 7}, + ] + }, ) - repo = ERDoseEUVRepository(db) - processor = ERDoseEUVProcessor(repo) + processor = ERDoseEUVProcessor(ERDoseEUVRepository(db)) - with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=100, + with redirect_stdout(StringIO()): + processor.run( + start_time=datetime.combine(target_date, datetime.min.time()), + end_time=datetime.combine(target_date + timedelta(days=1), datetime.min.time()), ) - self.assertEqual(db.executed, []) - self.assertEqual(len(db.partition_inserts), 0) - self.assertIn("lookback_done start_date=2026-05-04 end_date=2026-05-04", stdout.getvalue()) - self.assertIn("checked_dates=1 reloaded_dates=0 source_rows=0 inserted=0", stdout.getvalue()) + log_queries = [item for item in db.executed if "batch_event_log" in item[0].lower()] + self.assertEqual(log_queries, []) - def test_run_recent_days_truncates_and_reloads_when_counts_differ(self): + def test_explicit_date_truncates_and_reloads(self): target_date = date(2026, 5, 4) raw_df = pd.DataFrame( [ @@ -279,9 +300,8 @@ def test_run_recent_days_truncates_and_reloads_when_counts_differ(self): processor = ERDoseEUVProcessor(repo) with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=1, - reference_date=target_date, + processor.run( + target_date=target_date, chunk_size=100, ) @@ -292,30 +312,6 @@ def test_run_recent_days_truncates_and_reloads_when_counts_differ(self): analyze_queries = [item for item in db.executed if item[0] == "ANALYZE prism_common.er_dose_euv_parsed_1_prt_p20260504"] self.assertEqual(len(analyze_queries), 1) self.assertIs(analyze_queries[0][2], db.connection) - self.assertIn("lookback_done start_date=2026-05-04 end_date=2026-05-04", stdout.getvalue()) - self.assertIn("checked_dates=1 reloaded_dates=1 source_rows=1 inserted=1", stdout.getvalue()) - - def test_run_recent_days_skips_when_only_source_duplicates_differ(self): - target_date = date(2026, 5, 4) - db = FakeDB( - pd.DataFrame(), - source_counts={target_date: 2}, - target_counts={target_date: 1}, - distinct_source_counts={target_date: 1}, - ) - processor = ERDoseEUVProcessor(ERDoseEUVRepository(db)) - - with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=100, - ) - - self.assertEqual(db.executed, []) - self.assertEqual(len(db.partition_inserts), 0) - self.assertIn("checked_dates=1 reloaded_dates=0 source_rows=0 inserted=0", stdout.getvalue()) - if __name__ == "__main__": unittest.main() diff --git a/tests/test_er_dose_parser.py b/tests/test_er_dose_parser.py index 212da05..a2543de 100644 --- a/tests/test_er_dose_parser.py +++ b/tests/test_er_dose_parser.py @@ -144,6 +144,14 @@ def test_parse_lo_0052_aborted_lot(self): self.assertEqual(parsed.lot_name, "jmra22") self.assertEqual(parsed.lot_seq, 19789) + def test_parse_lo_0052_root_error_lot(self): + contents = "error getting root error for lot (name='pbs3i5.1_1905_0_mc600447', id=9616)" + raw = self._raw(contents, code="LO-0052") + parsed = parse_dose_error(raw) + self.assertEqual(parsed.lot_id, "pbs3i5.1_1905_0_mc600447") + self.assertEqual(parsed.lot_name, "pbs3i5") + self.assertEqual(parsed.lot_seq, 9616) + def test_parse_lo_0061(self): contents = "loading reticle 'gvhbrtb0v8' for lot id 2111." raw = self._raw(contents, code="LO-0061") diff --git a/tests/test_er_dose_processor.py b/tests/test_er_dose_processor.py index 085bda4..596f405 100644 --- a/tests/test_er_dose_processor.py +++ b/tests/test_er_dose_processor.py @@ -1,8 +1,9 @@ from __future__ import annotations +import json import unittest from contextlib import redirect_stdout -from datetime import datetime +from datetime import date, datetime, timedelta from io import StringIO from threading import Event @@ -44,12 +45,14 @@ def __init__( source_counts=None, target_counts=None, distinct_source_counts=None, + equipment_counts=None, ): self.raw_df = raw_df self.fetch_df_result = raw_df if fetch_df_result is None else fetch_df_result self.source_counts = source_counts or {} self.target_counts = target_counts or {} self.distinct_source_counts = distinct_source_counts or {} + self.equipment_counts = equipment_counts or {} self.executed = [] self.inserted = [] self.copy_options = [] @@ -57,10 +60,25 @@ def __init__( self.partition_inserts = [] self.queries_called = [] - def select(self, query, params=None): + def select(self, query, params=None, connection=None): self.fetch_query = query self.fetch_params = params lowered = query.lower() + if "with source_counts as" in lowered and "full outer join target_counts" in lowered: + target_date = params["start_time"].date() + rows = self.equipment_counts.get(target_date) + if rows is None: + source_count = self.source_counts.get(target_date, 0) + target_count = self.target_counts.get(target_date, 0) + rows = ([{"eq_name": "EQ1", "source_count": source_count, "target_count": target_count}] + if source_count or target_count else []) + return pd.DataFrame([ + {**row, + "total_source_count": sum(item["source_count"] for item in rows), + "total_target_count": sum(item["target_count"] for item in rows), + "matched": int(all(item["source_count"] == item["target_count"] for item in rows))} + for row in rows + ]) if "count(distinct" in lowered: target_date = params["start_time"].date() row_count = self.distinct_source_counts.get(target_date, self.source_counts.get(target_date, 0)) @@ -189,6 +207,12 @@ def test_fetch_counts_filter_active_nxe_eq_names(self): repo.fetch_source_count(target_date, distinct=True) self.assertIn("count(distinct (r.eq_name, r.code, r.code_occur_time))", db.fetch_query) + equipment_counts = repo.fetch_equipment_counts(datetime.combine(target_date, datetime.min.time()), datetime.combine(target_date + timedelta(days=1), datetime.min.time())) + self.assertEqual(equipment_counts, {"source_count": 1, "target_count": 1, "matched": True, "equipment_counts": [{"eq_name": "EQ1", "source_count": 1, "target_count": 1}]}) + self.assertIn("group by r.eq_name", db.fetch_query) + self.assertIn("group by p.eq_name", db.fetch_query) + self.assertIn("full outer join target_counts", db.fetch_query) + def test_fetch_latest_lot_states_filters_active_nxe_eq_names(self): db = FakeDB( pd.DataFrame(), @@ -282,8 +306,12 @@ def test_run_inserts_rows_without_deleting_existing_history(self): with redirect_stdout(StringIO()): processor.run(start_time=datetime(2026, 5, 1), end_time=datetime(2026, 5, 2)) - delete_queries = [query for query, _, _ in db.executed if query.strip().lower().startswith("delete")] - self.assertEqual(delete_queries, []) + parsed_delete_queries = [ + query + for query, _, _ in db.executed + if query.strip().lower().startswith("delete") and "er_dose_raw_parsed" in query.lower() + ] + self.assertEqual(parsed_delete_queries, []) parsed_insert = self._inserted_df(db, "prism_common.er_dose_raw_parsed") self.assertNotIn("parser_version", parsed_insert.columns) self.assertNotIn("parsing_status", parsed_insert.columns) @@ -450,6 +478,7 @@ def test_run_processes_multiple_chunks(self): self.assertEqual(len(db.inserted[0][1]), 2) self.assertEqual(len(db.inserted[1][1]), 1) self.assertEqual([option["target_date"] for option in db.copy_options], ["2026-05-01", "2026-05-01"]) + self.assertTrue(all(option["analyze"] for option in db.copy_options)) analyze_queries = [query for query, _, _ in db.executed if query == "ANALYZE prism_common.er_dose_raw_parsed_1_prt_p20260501"] self.assertEqual(len(analyze_queries), 1) @@ -499,6 +528,50 @@ def copy_insert_to_partition_table(self, *args, **kwargs): self.assertTrue(db.second_fetch_during_insert) self.assertEqual(len(db.partition_inserts), 2) + def test_run_prefetches_next_chunk_while_current_chunk_is_parsing(self): + class PrefetchFakeDB(FakeDB): + def __init__(self, raw_df): + super().__init__(raw_df) + self.second_fetch_started = Event() + + def select_in_chunks(self, query, params=None, chunk_size=10000): + self.fetch_query = query + self.fetch_params = params + for start in range(0, len(self.raw_df), chunk_size): + if start > 0: + self.second_fetch_started.set() + yield self.raw_df.iloc[start : start + chunk_size].copy() + + class BlockingParseProcessor(ERDoseProcessor): + def __init__(self, repository, second_fetch_started): + super().__init__(repository) + self.second_fetch_started = second_fetch_started + self.first_parse_saw_prefetch = False + + def _parse_chunk(self, raw_df): + if not self.first_parse_saw_prefetch: + self.first_parse_saw_prefetch = self.second_fetch_started.wait(timeout=1) + return super()._parse_chunk(raw_df) + + raw_df = pd.DataFrame( + [ + self._row(1, "dw-3411", SAMPLE_CONTENTS), + self._row(2, "dw-3411", SAMPLE_CONTENTS), + ] + ) + db = PrefetchFakeDB(raw_df) + processor = BlockingParseProcessor(ERDoseRepository(db), db.second_fetch_started) + + with redirect_stdout(StringIO()): + processor.run( + start_time=datetime(2026, 5, 1), + end_time=datetime(2026, 5, 2), + chunk_size=1, + ) + + self.assertTrue(processor.first_parse_saw_prefetch) + self.assertEqual(len(db.partition_inserts), 2) + def test_run_propagates_background_insert_error(self): class FailingFakeDB(FakeDB): def copy_insert_to_partition_table(self, *args, **kwargs): @@ -528,143 +601,6 @@ def test_run_accepts_target_date_and_builds_daily_window(self): self.assertEqual(db.fetch_params["start_time"], datetime(2026, 5, 1, 0, 0, 0)) self.assertEqual(db.fetch_params["end_time"], datetime(2026, 5, 2, 0, 0, 0)) - def test_run_without_window_uses_recent_days_lookback(self): - raw_df = pd.DataFrame([ - self._row(1, "dw-3411", SAMPLE_CONTENTS, code_occur_time=datetime(2026, 5, 2, 10, 0, 0)) - ]) - db = FakeDB( - raw_df, - source_counts={ - datetime(2026, 5, 1).date(): 0, - datetime(2026, 5, 2).date(): 1, - }, - target_counts={ - datetime(2026, 5, 1).date(): 0, - datetime(2026, 5, 2).date(): 1, - }, - ) - repo = ERDoseRepository(db) - processor = ERDoseProcessor(repo) - - with redirect_stdout(StringIO()) as stdout: - processor.run( - lookback_days=2, - reference_date=datetime(2026, 5, 2).date(), - chunk_size=100, - ) - - self.assertEqual(db.executed, []) - self.assertEqual(len(db.partition_inserts), 0) - self.assertIn("lookback_done start_date=2026-05-01 end_date=2026-05-02", stdout.getvalue()) - self.assertIn("checked_dates=2 reloaded_dates=0 source_rows=0 inserted=0", stdout.getvalue()) - - def test_run_recent_days_skips_when_counts_match(self): - raw_df = pd.DataFrame([ - self._row(1, "dw-3411", SAMPLE_CONTENTS, code_occur_time=datetime(2026, 5, 2, 10, 0, 0)) - ]) - db = FakeDB( - raw_df, - source_counts={ - datetime(2026, 5, 1).date(): 0, - datetime(2026, 5, 2).date(): 1, - }, - target_counts={ - datetime(2026, 5, 1).date(): 0, - datetime(2026, 5, 2).date(): 1, - }, - ) - repo = ERDoseRepository(db) - processor = ERDoseProcessor(repo) - - with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=2, - reference_date=datetime(2026, 5, 2).date(), - chunk_size=100, - ) - - self.assertEqual(db.executed, []) - self.assertEqual(len(db.partition_inserts), 0) - self.assertIn("lookback_done start_date=2026-05-01 end_date=2026-05-02", stdout.getvalue()) - self.assertIn("checked_dates=2 reloaded_dates=0 source_rows=0 inserted=0", stdout.getvalue()) - - def test_run_recent_days_truncates_and_reloads_when_counts_differ(self): - target_date = datetime(2026, 5, 1).date() - raw_df = pd.DataFrame([ - self._row(1, "dw-3411", SAMPLE_CONTENTS) - ]) - db = FakeDB( - raw_df, - source_counts={target_date: 1}, - target_counts={target_date: 0}, - distinct_source_counts={target_date: 0}, - ) - repo = ERDoseRepository(db) - processor = ERDoseProcessor(repo) - - with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=100, - ) - - truncate_queries = [query for query, _, _ in db.executed if query.strip().upper().startswith("TRUNCATE")] - self.assertEqual(truncate_queries, ["truncate table prism_common.er_dose_raw_parsed_1_prt_p20260501"]) - self.assertEqual(len(db.partition_inserts), 1) - self.assertIsNone(db.executed[0][2]) - analyze_queries = [item for item in db.executed if item[0] == "ANALYZE prism_common.er_dose_raw_parsed_1_prt_p20260501"] - self.assertEqual(len(analyze_queries), 1) - self.assertIsNone(analyze_queries[0][2]) - self.assertIn("lookback_done start_date=2026-05-01 end_date=2026-05-01", stdout.getvalue()) - self.assertIn("checked_dates=1 reloaded_dates=1 source_rows=1 inserted=1", stdout.getvalue()) - - def test_run_recent_days_skips_when_only_source_duplicates_differ(self): - target_date = datetime(2026, 5, 1).date() - db = FakeDB( - pd.DataFrame(), - source_counts={target_date: 2}, - target_counts={target_date: 1}, - distinct_source_counts={target_date: 1}, - ) - processor = ERDoseProcessor(ERDoseRepository(db)) - - with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=100, - ) - - self.assertEqual(db.executed, []) - self.assertEqual(len(db.partition_inserts), 0) - self.assertIn("checked_dates=1 reloaded_dates=0 source_rows=0 inserted=0", stdout.getvalue()) - - def test_run_recent_days_skips_when_counts_match_even_if_specific_row_is_missing(self): - target_date = datetime(2026, 5, 1).date() - raw_df = pd.DataFrame([ - self._row(1, "dw-3411", SAMPLE_CONTENTS) - ]) - db = FakeDB( - raw_df, - source_counts={target_date: 1}, - target_counts={target_date: 1}, - ) - repo = ERDoseRepository(db) - processor = ERDoseProcessor(repo) - - with redirect_stdout(StringIO()) as stdout: - processor.run_recent_days( - lookback_days=1, - reference_date=target_date, - chunk_size=100, - ) - - truncate_queries = [query for query, _, _ in db.executed if query.strip().upper().startswith("TRUNCATE")] - self.assertEqual(truncate_queries, []) - self.assertEqual(len(db.partition_inserts), 0) - self.assertIn("lookback_done start_date=2026-05-01 end_date=2026-05-01", stdout.getvalue()) - self.assertIn("checked_dates=1 reloaded_dates=0 source_rows=0 inserted=0", stdout.getvalue()) def test_insert_parsed_df_keeps_integer_columns_as_nullable_int(self): db = FakeDB(pd.DataFrame()) @@ -777,25 +713,33 @@ def _inserted_df(self, db, table_name): return inserted_df raise AssertionError(f"{table_name} was not inserted") - - def test_summary_tables_upsert_called(self): - raw_df = pd.DataFrame( - [ - self._row(1, "dw-3411", SAMPLE_CONTENTS, code_occur_time=datetime(2026, 6, 15, 10, 0, 0)), - ] + def test_explicit_date_truncates_and_reloads(self): + target_date = datetime(2026, 5, 1).date() + raw_df = pd.DataFrame([ + self._row(1, "dw-3411", SAMPLE_CONTENTS) + ]) + db = FakeDB( + raw_df, + source_counts={target_date: 1}, + target_counts={target_date: 0}, + distinct_source_counts={target_date: 0}, ) - db = FakeDB(raw_df) repo = ERDoseRepository(db) processor = ERDoseProcessor(repo) - processor.run(start_time=datetime(2026, 6, 15), end_time=datetime(2026, 6, 16), chunk_size=1000) - executed_queries = [query.lower() for query, _, _ in db.executed] - has_die_yield = any("de_trend_die_yield_daily" in q for q in executed_queries) - has_root_cause = any("de_trend_root_cause_daily" in q for q in executed_queries) - - self.assertTrue(has_die_yield, "de_trend_die_yield_daily UPSERT was not executed") - self.assertTrue(has_root_cause, "de_trend_root_cause_daily UPSERT was not executed") + with redirect_stdout(StringIO()) as stdout: + processor.run( + target_date=target_date, + chunk_size=100, + ) + truncate_queries = [query for query, _, _ in db.executed if query.strip().upper().startswith("TRUNCATE")] + self.assertEqual(truncate_queries, ["truncate table prism_common.er_dose_raw_parsed_1_prt_p20260501"]) + self.assertEqual(len(db.partition_inserts), 1) + self.assertIsNone(db.executed[0][2]) + analyze_queries = [item for item in db.executed if item[0] == "ANALYZE prism_common.er_dose_raw_parsed_1_prt_p20260501"] + self.assertEqual(len(analyze_queries), 1) + self.assertIsNone(analyze_queries[0][2]) if __name__ == "__main__": unittest.main() diff --git a/tests/test_er_dose_root_cause.py b/tests/test_er_dose_root_cause.py index 6eb1aaa..97813f8 100644 --- a/tests/test_er_dose_root_cause.py +++ b/tests/test_er_dose_root_cause.py @@ -42,10 +42,10 @@ def test_non_root_cause_contents_are_skipped(): assert parse_root_cause("system info: normal message") is None -def test_parse_euv_root_cause_tolerates_production_label_typos(): +def test_parse_euv_root_cause_with_decimal_pulse_counts(): contents = r"""Dose error detected in file: ADECetdcData_FDD_LC_EEI_SCANNER_DOSE_ERROR_EVENT_20260614_235011_0399+0900.zip. -\nRoot clause : Low dose margin -\nExposesue I D : 47737 +\nRoot cause : Low dose margin +\nExposure ID : 47737 \nTime : 2026-06-14T23:50:10.953142+09:00 \nExposure length : 0.2338 [s] \nDuty cycle : 99.62 [perc] @@ -72,7 +72,8 @@ def test_parse_euv_root_cause_tolerates_production_label_typos(): tzinfo=timezone(timedelta(hours=9)), ) assert parsed.exposure_length == Decimal("0.2338") - assert parsed.dose_error == Decimal("-1.71") - assert parsed.pulses_euv_lt_0_6dt_tot == 3 + assert parsed.min_dose_error == Decimal("-1.71") + assert parsed.max_dose_error == Decimal("-1.71") + assert parsed.pulses_euv_0_6dt_tot == 3 assert parsed.fed_pulses == 3 assert parsed.software_version == "3.0" diff --git a/tests/test_er_dose_schema.py b/tests/test_er_dose_schema.py index 7d99a81..3aab591 100644 --- a/tests/test_er_dose_schema.py +++ b/tests/test_er_dose_schema.py @@ -9,6 +9,7 @@ ROOT_CAUSE_DDL_PATH = Path(__file__).resolve().parents[1] / "er_dose" / "sql" / "create_er_dose_euv_parsed.sql" ROOT_CAUSE_RENAME_PATH = Path(__file__).resolve().parents[1] / "er_dose" / "sql" / "rename_er_dose_euv_parsed_columns.sql" ROOT_CAUSE_MIGRATION_PATH = Path(__file__).resolve().parents[1] / "er_dose" / "sql" / "migrate_er_dose_euv_parsed_schema.sql" +BATCH_EVENT_LOG_DDL_PATH = Path(__file__).resolve().parents[1] / "er_dose" / "sql" / "create_batch_event_log.sql" def _ddl() -> str: @@ -35,6 +36,10 @@ def _root_cause_migration_sql() -> str: return ROOT_CAUSE_MIGRATION_PATH.read_text(encoding="utf-8").lower() +def _batch_event_log_ddl() -> str: + return BATCH_EVENT_LOG_DDL_PATH.read_text(encoding="utf-8").lower() + + def test_parsed_table_primary_key_matches_documented_partition_key(): ddl = _ddl() @@ -135,8 +140,8 @@ def test_root_cause_table_is_fe_facing_matching_table(): assert "match_status" not in ddl assert "dose_error numeric(12,7)" not in ddl assert "dose_error_detected_in_file text" in ddl - assert "min_dose_error numeric(12,7)" in ddl - assert "max_dose_error numeric(12,7)" in ddl + assert "min_dose_error numeric(20,10)" in ddl + assert "max_dose_error numeric(20,10)" in ddl assert "pulses_euv_0_6dt_tot integer" in ddl assert "software_version text" in ddl assert "parser_version" not in ddl @@ -174,3 +179,15 @@ def test_root_cause_rename_sql_renames_existing_columns(): assert "call rename_column_if_exists('prism_common.er_dose_euv_parsed', 'pulses_euv_lt_0_6dt_tot', 'pulses_euv_0_6dt_tot')" in sql assert "call rename_column_if_exists('prism_common.er_dose_euv_parsed', 'software version', 'software_version')" in sql assert "from pg_inherits" in sql + + +def test_batch_event_log_is_generic_and_json_based(): + ddl = _batch_event_log_ddl() + + assert "create table if not exists mbeat.batch_event_log" in ddl + assert "batch_name varchar(100) not null" in ddl + assert "target_date date" in ddl + assert "event_type varchar(50) not null" in ddl + assert "message text" in ddl + assert "data jsonb not null" in ddl + assert "idx_batch_event_log_batch_date" in ddl diff --git a/tests/test_er_dose_summary.py b/tests/test_er_dose_summary.py new file mode 100644 index 0000000..672d3a4 --- /dev/null +++ b/tests/test_er_dose_summary.py @@ -0,0 +1,63 @@ +from datetime import date, datetime +from unittest.mock import Mock + +import pandas as pd +import pytest + +from airflow_modules.er_dose_jobs import write_equipment_count_log +from er_dose.euv.euv_repository import ERDoseEUVRepository +from er_dose.raw.raw_repository import ERDoseRepository + + +@pytest.mark.parametrize('repository_type', [ERDoseRepository, ERDoseEUVRepository]) +def test_equipment_counts_without_database_json_functions(repository_type): + db = Mock() + db.select.return_value = pd.DataFrame([ + {'eq_name': 'A', 'source_count': 3, 'target_count': 2, 'total_source_count': 5, 'total_target_count': 5, 'matched': 0}, + {'eq_name': 'B', 'source_count': 2, 'target_count': 3, 'total_source_count': 5, 'total_target_count': 5, 'matched': 0}, + ]) + result = repository_type(db).fetch_equipment_counts(datetime(2026, 5, 1), datetime(2026, 5, 2)) + query = db.select.call_args.args[0].lower() + assert 'jsonb_build_object' not in query and 'jsonb_agg' not in query + assert 'over ()' in query + assert result['source_count'] == result['target_count'] == 5 + assert result['matched'] is False + assert len(result['equipment_counts']) == 2 + + +def test_raw_and_euv_are_separate_rows_with_equipment_arrays(): + + repo = Mock() + raw = {"equipment_counts": [ + {"eq_name": "EQ1", "source_count": 100, "target_count": 99}, + {"eq_name": "RAW_ONLY", "source_count": 5, "target_count": 5}, + ]} + euv = {"equipment_counts": [ + {"eq_name": "EQ1", "source_count": 20, "target_count": 20}, + {"eq_name": "EUV_ONLY", "source_count": 7, "target_count": 6}, + ]} + repo.fetch_equipment_counts.side_effect = [raw, euv] + for batch in ['ER_DOSE_RAW', 'ER_DOSE_EUV']: + assert write_equipment_count_log(repo, batch, datetime(2026,5,1), datetime(2026,5,2)) is None + calls = repo.insert_batch_log.call_args_list + assert len(calls) == 2 + assert calls[0].kwargs['batch_name'] == 'ER_DOSE_RAW' + assert calls[0].kwargs['data'] == raw + assert calls[1].kwargs['batch_name'] == 'ER_DOSE_EUV' + assert calls[1].kwargs['data'] == euv + + +def test_common_log_insert_accepts_unrelated_batch_payload(): + import json + from er_dose.common.batch_log_repository import insert_batch_log + + db = Mock() + data = {"file_name": "결과.csv", "rows": 12, "details": ["ok"]} + insert_batch_log(db, 'OTHER_BATCH', date(2026, 5, 1), 'FILE_EXPORTED', 'export finished', data) + db.execute.assert_called_once() + params = db.execute.call_args.kwargs['params'] + assert params['batch_name'] == 'OTHER_BATCH' + assert params['event_type'] == 'FILE_EXPORTED' + assert params['message'] == 'export finished' + assert json.loads(params['data']) == data + assert data == {"file_name": "결과.csv", "rows": 12, "details": ["ok"]} diff --git a/tests/test_er_dose_summary_replacement.py b/tests/test_er_dose_summary_replacement.py new file mode 100644 index 0000000..d729fec --- /dev/null +++ b/tests/test_er_dose_summary_replacement.py @@ -0,0 +1,46 @@ +from datetime import datetime + +import pytest + +from er_dose.raw.raw_repository import ERDoseRepository + + +@pytest.mark.parametrize("repository_type,method,table", [ + (ERDoseRepository, "die_yield_daily_summary", "de_trend_die_yield_daily"), + (ERDoseRepository, "root_cause_daily_summary", "de_trend_root_cause_daily"), +]) +@pytest.mark.parametrize("fail_insert", [False, True]) +def test_summary_delete_and_insert_are_separate_calls(repository_type, method, table, fail_insert): + class DB: + def __init__(self): + self.connection = object() + self.events = [] + self.calls = [] + + def transaction(self): + raise AssertionError("summary must not create a shared transaction") + + def execute(self, query, params): + self.calls.append((query, params)) + action = query.strip().split()[0].lower() + self.events.append(action) + if action == "insert" and fail_insert: + raise RuntimeError("insert failed") + return None + + db = DB() + repo = repository_type(db) + start, end = datetime(2026, 5, 1), datetime(2026, 5, 2) + assert getattr(repo, "delete_" + method)(start, end) is None + if fail_insert: + with pytest.raises(RuntimeError, match="insert failed"): + getattr(repo, "insert_" + method)(start, end) + else: + assert getattr(repo, "insert_" + method)(start, end) is None + expected = ["delete", "insert"] + assert db.events == expected + delete_query, params = db.calls[0] + assert table in delete_query + assert "where occur_date >= cast(:start_time as date)" in delete_query + assert params == {"start_time": datetime(2026, 5, 1), "end_time": datetime(2026, 5, 2)} + assert "on conflict" not in db.calls[1][0].lower() diff --git a/tests/test_postgres_db.py b/tests/test_postgres_db.py index 86c286a..af290ce 100644 --- a/tests/test_postgres_db.py +++ b/tests/test_postgres_db.py @@ -102,3 +102,32 @@ def test_copy_insert_to_partition_table_can_skip_analyze(self): if __name__ == "__main__": unittest.main() + + +def test_execute_adapts_named_parameters_without_interpolating_values(): + from datetime import date + + db = PostgresDB(dsn='postgresql://unused') + connection = FakeConnection() + connection.cursor_obj.rowcount = 2 + value = "x'; DROP TABLE example; --" + result = db.execute( + "update example set value = :value where day::date = :day and label like '%test%'", + {'value': value, 'day': date(2026, 5, 1)}, connection=connection, + ) + query, params = connection.cursor_obj.executed[0] + assert query == "update example set value = %(value)s where day::date = %(day)s and label like '%%test%%'" + assert params == {'value': value, 'day': date(2026, 5, 1)} + assert value not in query + assert result == 2 + assert not connection.committed + assert not connection.closed + + +def test_execute_preserves_native_driver_parameters(): + db = PostgresDB(dsn='postgresql://unused') + connection = FakeConnection() + connection.cursor_obj.rowcount = 1 + query = 'delete from example where id = %(id)s' + db.execute(query, {'id': 7}, connection=connection) + assert connection.cursor_obj.executed == [(query, {'id': 7})]