From 0c85dd93c22d13a0bfa012140627e3636db77c44 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Thu, 6 Aug 2026 17:02:20 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E8=87=AA=E9=80=82=E5=BA=94=E7=AD=89?= =?UTF-8?q?=E5=BE=85=E5=B9=B2=E6=89=B0=E6=95=B0=E6=8D=AE=E5=B9=B6=E5=BD=92?= =?UTF-8?q?=E6=A1=A3=E7=BB=93=E6=9E=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 34 ++- aidocs/project_context.md | 8 + main.py | 437 ++++++++++++++++++++++++++------------ tests/test_pipeline.py | 209 +++++++++++++++--- 4 files changed, 521 insertions(+), 167 deletions(-) diff --git a/README.md b/README.md index 554740e..7ec0908 100644 --- a/README.md +++ b/README.md @@ -1,36 +1,50 @@ # InterferenceETL -InterferenceETL processes the newest complete hour of interference KPI files from Metrix storage. It reads seven ZIP/XLSX source types, converts each `Sheet0` to CSV, and creates one merged interference summary. +InterferenceETL follows the newest interference KPI time group in Metrix storage. It reads seven ZIP/XLSX source types, converts each `Sheet0` to CSV, and creates one merged high-interference summary. + +Source file names are identified by one of the seven fixed prefixes plus a final `_YYYYMMDDHHMMHHMM.zip`; text between the prefix and time is unrestricted. Only the first `YYYYMMDDHHMM` start time is used for grouping. Start times are rounded down to natural 15-minute boundaries, so `12:00` through `12:14` belong to `12:00`. This also works unchanged when all providers switch together to hourly or daily delivery. + +Every run follows only the group containing the newest source start time. It never falls back or backfills an older group. If any source type is missing, the script prints `status=waiting` with the missing types and exits successfully. When one source has multiple files in the group, its latest start time wins. ## Output Each run writes one window directory: ```text -output/2026073110001100/ +output/20260731100000/ ├── converted/ │ ├── 5G下FDD干扰监控_2026073110001100.csv │ ├── 5G干扰监控_2026073110001100.csv │ └── ... seven source CSV files -├── interference_summary_2026073110001100.csv +├── interference_summary_20260731100000.csv └── manifest.json ``` -For every run, the script also reads the latest filename-dated XLSX from each configured CellData directory. It builds CellData CGI as `460-00-{eNB/gNB}-{CI}` and adds coordinates plus direction angle to matching interference rows. Unmatched rows keep empty coordinates and use `0` for azimuth. +For every processed group, the script also reads the latest filename-dated XLSX from each configured CellData directory. It builds CellData CGI as `460-00-{eNB/gNB}-{CI}` and adds coordinates plus direction angle to matching interference rows. Unmatched rows keep empty coordinates and use `0` for azimuth. The summary columns are: ```text -hour_start,hour_end,network_type,cgi,cell_name,interference_dbm,longitude,latitude,azimuth,nearby_count +metric_time,network_type,cgi,cell_name,interference_dbm,longitude,latitude,azimuth,nearby_count ``` +All selected workbook rows must belong to the selected natural 15-minute group. Every summary row receives the same normalized `metric_time`; a row outside the group fails the run before database or history output. + `network_type` is derived from the seven source types: both `5G...` sources are `2.6G`, both `700M...` sources are `700M`, and SDR/反开 sources are `4G`. The summary and database keep only high-interference rows: `2.6G >= -107 dBm` and `700M/4G >= -110 dBm`. Converted source CSV files remain full source conversions. -`nearby_count` is the number of other high-interference cells from the same selected hour within 1 km, across all network types. The cell itself is excluded, the 1 km boundary is included, and rows without coordinates use `0`. Candidate cells are found through a 1 km spatial grid index and confirmed with exact Haversine distance. +`nearby_count` is the number of other high-interference cells from the same selected time group within 1 km, across all network types. The cell itself is excluded, the 1 km boundary is included, and rows without coordinates use `0`. Candidate cells are found through a 1 km spatial grid index and confirmed with exact Haversine distance. -The script never modifies or deletes source storage files. By default it writes the complete selected hour through the Metrix Database API and keeps only that hour in the target table. +The script never modifies or deletes source storage files. Before downloading source ZIP files or CellData, it compares the normalized target time with the database. A target already present, or older than the database, is not processed again. -The database table contains one time column, `metric_time DATETIME`, which is the source KPI start time, `network_type VARCHAR(16)`, nullable `longitude` and `latitude` columns, `azimuth DECIMAL(6,2) NOT NULL DEFAULT 0`, and `nearby_count INT NOT NULL DEFAULT 0`. A same-hour rerun replaces that whole batch. After a successful insert, rows for every other hour are deleted in the same transaction. The script refuses to replace a newer database hour with an older source hour. +The database table contains normalized `metric_time DATETIME`, `network_type VARCHAR(16)`, nullable `longitude` and `latitude`, `azimuth DECIMAL(6,2) NOT NULL DEFAULT 0`, and `nearby_count INT NOT NULL DEFAULT 0`. A successful transaction replaces the target batch and deletes every other database time, so the table retains only the latest processed group. + +After the database succeeds, the same nine result columns are uploaded as UTF-8-BOM CSV to: + +```text +/网优日常优化数据文档/(勿删)干扰定时小时指标/干扰历史数据/YYYY-MM-DD/干扰数据处理结果_YYYYMMDDHHMMSS.csv +``` + +The filename time comes from normalized `metric_time`, never the server clock. If the database already contains the target time but its history CSV is missing or empty, the script exports that time from the database and repairs the history file without reprocessing source data. CGI is generated with fixed rules: @@ -49,7 +63,7 @@ python -m unittest discover -s tests -v python main.py --source-dir C:\path\to\mock-or-exported-tree --output-dir output --no-database ``` -Use `--window 2026073110001100` to require one exact source window. Without it, the script selects the newest window containing all seven source types and falls back if the newest observed window is incomplete. +Use `--window 2026073110001100` to select the natural 15-minute group containing that source start time. Without it, the newest observed source start time determines the group. Incomplete groups wait and never fall back. ## Metrix Script Management @@ -93,7 +107,7 @@ Both storage reads and database writes use the Metrix API at `http://188.5.127.1 ```text --source-dir PATH Read a local directory tree instead of Metrix API --output-dir PATH Output root, default output or INTERFERENCE_OUTPUT_DIR ---window WINDOW Require an exact 16-digit source window +--window WINDOW Select the natural 15-minute group for a 16-digit source window --lookback-days N Number of newest date directories scanned, default 3 --storage-id ID Metrix storage connection ID --root PATH Source directory in Metrix storage diff --git a/aidocs/project_context.md b/aidocs/project_context.md index 343f71a..ab182ea 100644 --- a/aidocs/project_context.md +++ b/aidocs/project_context.md @@ -79,3 +79,11 @@ - Candidate lookup uses a dependency-free 1 km Earth-centered three-dimensional grid index. Only the current and 26 adjacent buckets are checked, then Haversine distance confirms the exact radius; this avoids a full all-pairs scan while preserving distance accuracy. - Unit tests pass in both the project virtual environment and `interference-etl-runtime:1.1` image. A read-only real run for window `2026080613001400` produced 1,116 rows, including 1,110 with coordinates, 925 with non-zero nearby counts, and a maximum count of 39. All 1,116 indexed results matched a separate brute-force comparison, which found 3,973 qualifying pairs. - SSH deployment and Metrix runner execution `8f2fd51d7de14c93b4b0eb9bee8fafa0` succeeded for the same window. Database API verification found only `2026-08-06 13:00:00`, with 1,116 rows, 925 non-zero counts, a maximum of 39, and 3,973 nearby pairs. All six rows without coordinates have count `0`. + +## 2026-08-06: Adaptive source waiting and history export + +- Source recognition now fixes only the seven known prefixes and final `_YYYYMMDDHHMMHHMM.zip`; middle text may change with provider granularity. The first 12 digits are rounded down to natural 15-minute boundaries, while the final four end-time digits do not participate in grouping. +- Each run follows only the group containing the globally newest source start time. It does not fall back or backfill older groups. Missing sources or no matching files produce `status=waiting` and a successful exit; storage API failures remain task failures. Multiple files from one source in a group select the latest source start time. +- All workbook row start times must round into the target group. Result CSV and database rows use one normalized `metric_time` and the database result columns only; `hour_start` and `hour_end` were removed from summary CSV output. +- Database time is checked before source ZIP and CellData downloads. A newer or equal database time skips source processing. Successful database output is archived through the Storage API under `干扰历史数据/YYYY-MM-DD/干扰数据处理结果_YYYYMMDDHHMMSS.csv`; an equal database time with a missing/empty history file is exported from the database and repaired. +- Eighteen tests pass on Windows and in the offline runtime image. Live read-only selection saw incomplete newest group `2026-08-06 15:00:00`, correctly waited for missing `5G下FDD干扰监控`, and did not produce output. Explicit read-only validation of complete group `2026-08-06 14:00:00` produced 1,092 rows with one normalized `metric_time` and the expected nine-column result schema. diff --git a/main.py b/main.py index 11c0a0b..3db5d9f 100644 --- a/main.py +++ b/main.py @@ -9,7 +9,7 @@ import json import math import os from dataclasses import dataclass -from datetime import datetime, timedelta, timezone +from datetime import datetime, timezone from pathlib import Path import posixpath import re @@ -75,12 +75,13 @@ CELL_DATA_DIRECTORIES = ( "/网优日常优化数据文档/日常性能报表/2026年/TDD_LTE/LTE基础信息数据/LTE_小区信息表", "/网优日常优化数据文档/日常性能报表/2026年/FDD_LTE/FDD基础信息数据/FDD小区信息表", ) -FILE_RE = re.compile(r"^(?P.+)_LWP_每小时_过滤110_(?P\d{16})\.zip$", re.IGNORECASE) +FILE_TIME_RE = re.compile(r"_(?P\d{16})\.zip$", re.IGNORECASE) CELL_DATA_FILE_RE = re.compile(r"(?P20\d{6})\.xlsx$", re.IGNORECASE) LTE_PLMN = "460-00" DATABASE_NAME = "interference_etl" DATABASE_TABLE = "interference_hourly_summary" +DEFAULT_HISTORY_ROOT = f"{DEFAULT_ROOT}/干扰历史数据" HEADER_FDD = ( "开始时间", @@ -173,8 +174,7 @@ EXPECTED_HEADERS = { } SUMMARY_HEADER = ( - "hour_start", - "hour_end", + "metric_time", "network_type", "cgi", "cell_name", @@ -206,8 +206,18 @@ class Source(Protocol): class SummaryStore(Protocol): + def latest_metric_time(self) -> datetime | None: ... + def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: ... + def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: ... + + +class HistoryStore(Protocol): + def exists(self, metric_time: datetime) -> bool: ... + + def upload(self, metric_time: datetime, payload: bytes) -> str: ... + class MetrixApiClient: def __init__(self, base_url: str, token: str) -> None: @@ -224,7 +234,37 @@ class MetrixApiClient: def post_json(self, endpoint: str, payload: dict[str, object], timeout: int = 30) -> object: body = json.dumps(payload, ensure_ascii=False).encode("utf-8") - return self._decode_json(self._request("POST", endpoint, body=body, timeout=timeout)) + return self._decode_json( + self._request("POST", endpoint, body=body, content_type="application/json", timeout=timeout) + ) + + def post_file( + self, + endpoint: str, + query: dict[str, object], + filename: str, + payload: bytes, + timeout: int = 120, + ) -> object: + boundary = f"----MetrixBoundary{time.time_ns()}" + disposition = f'Content-Disposition: form-data; name="file"; filename="{filename}"\r\n' + body = ( + f"--{boundary}\r\n".encode() + + disposition.encode("utf-8") + + b"Content-Type: text/csv\r\n\r\n" + + payload + + f"\r\n--{boundary}--\r\n".encode() + ) + return self._decode_json( + self._request( + "POST", + endpoint, + query=query, + body=body, + content_type=f"multipart/form-data; boundary={boundary}", + timeout=timeout, + ) + ) def _request( self, @@ -232,14 +272,15 @@ class MetrixApiClient: endpoint: str, query: dict[str, object] | None = None, body: bytes | None = None, + content_type: str = "", timeout: int = 30, ) -> bytes: url = f"{self.base_url}{endpoint}" if query: url = f"{url}?{urlencode(query)}" headers = {"Authorization": f"Bearer {self.token}", "Accept": "application/json"} - if body is not None: - headers["Content-Type"] = "application/json" + if content_type: + headers["Content-Type"] = content_type request = Request(url, data=body, headers=headers, method=method) last_error: Exception | None = None for attempt in range(3): @@ -273,62 +314,24 @@ class ApiSummaryStore: self.connection_id = connection_id self.table = DATABASE_TABLE - def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: - if not rows: - raise ProcessingError("Refusing to replace database data with an empty batch") - metric_times = {datetime.strptime(row["hour_start"], "%Y-%m-%d %H:%M:%S") for row in rows} - if len(metric_times) != 1: - raise ProcessingError("Database batch must contain exactly one metric hour") - metric_time = next(iter(metric_times)) - self._query( - f"CREATE DATABASE IF NOT EXISTS `{DATABASE_NAME}` " - "CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci" - ) - self._query( - f""" - CREATE TABLE IF NOT EXISTS `{self.table}` ( - metric_time DATETIME NOT NULL COMMENT '指标开始时间', - network_type VARCHAR(16) NOT NULL DEFAULT '', - cgi VARCHAR(128) NOT NULL, - cell_name VARCHAR(255) NOT NULL, - interference_dbm DECIMAL(10,3) NOT NULL, - longitude DECIMAL(10,6) NULL, - latitude DECIMAL(10,6) NULL, - azimuth DECIMAL(6,2) NOT NULL DEFAULT 0, - nearby_count INT NOT NULL DEFAULT 0, - PRIMARY KEY (metric_time, cgi) - ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 - """, - database=DATABASE_NAME, - ) - columns = self.client.get_json( - self._endpoint("columns"), - {"database": DATABASE_NAME, "table": self.table}, - ) - if not isinstance(columns, list): - raise ProcessingError("Metrix Database API returned invalid column metadata") - existing_columns = {str(item.get("name")) for item in columns if isinstance(item, dict)} - migrations = { - "network_type": "VARCHAR(16) NOT NULL DEFAULT ''", - "longitude": "DECIMAL(10,6) NULL", - "latitude": "DECIMAL(10,6) NULL", - "azimuth": "DECIMAL(6,2) NOT NULL DEFAULT 0", - "nearby_count": "INT NOT NULL DEFAULT 0", - } - for column, definition in migrations.items(): - if column not in existing_columns: - self._query( - f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` {definition}", - database=DATABASE_NAME, - ) - - latest_payload = self._query( + def latest_metric_time(self) -> datetime | None: + self._ensure_schema() + payload = self._query( f"SELECT MAX(metric_time) AS latest_time FROM `{self.table}`", database=DATABASE_NAME, ) - latest_rows = latest_payload.get("rows", []) - latest_value = latest_rows[0].get("latest_time") if latest_rows else None - latest_time = datetime.fromisoformat(str(latest_value)) if latest_value else None + rows = payload.get("rows", []) + value = rows[0].get("latest_time") if rows else None + return datetime.fromisoformat(str(value)) if value else None + + def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: + if not rows: + raise ProcessingError("Refusing to replace database data with an empty batch") + metric_times = {datetime.strptime(row["metric_time"], "%Y-%m-%d %H:%M:%S") for row in rows} + if len(metric_times) != 1: + raise ProcessingError("Database batch must contain exactly one metric time") + metric_time = next(iter(metric_times)) + latest_time = self.latest_metric_time() if latest_time is not None and latest_time > metric_time: raise ProcessingError( f"Database already contains newer metric time {latest_time:%Y-%m-%d %H:%M:%S}; " @@ -374,13 +377,80 @@ class ApiSummaryStore: "old_rows_deleted": int(results[3].get("affected_rows") or 0), } + def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: + self._ensure_schema() + metric_literal = f"'{metric_time:%Y-%m-%d %H:%M:%S}'" + sql = ( + "SELECT metric_time, network_type, cgi, cell_name, interference_dbm, " + f"longitude, latitude, azimuth, nearby_count FROM `{self.table}` " + f"WHERE metric_time = {metric_literal} ORDER BY cgi, cell_name" + ) + result: list[dict[str, str]] = [] + page = 1 + while True: + payload = self._query(sql, database=DATABASE_NAME, page=page, page_size=1000) + rows = payload.get("rows", []) + if not isinstance(rows, list): + raise ProcessingError("Metrix Database API returned invalid rows") + for row in rows: + normalized = {column: normalize_cell(row.get(column)) for column in SUMMARY_HEADER} + normalized["metric_time"] = datetime.fromisoformat(normalized["metric_time"]).strftime("%Y-%m-%d %H:%M:%S") + result.append(normalized) + total = int(payload.get("total") or len(result)) + if len(result) >= total or not rows: + return result + page += 1 + + def _ensure_schema(self) -> None: + self._query( + f"CREATE DATABASE IF NOT EXISTS `{DATABASE_NAME}` " + "CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci" + ) + self._query( + f""" + CREATE TABLE IF NOT EXISTS `{self.table}` ( + metric_time DATETIME NOT NULL COMMENT '指标开始时间', + network_type VARCHAR(16) NOT NULL DEFAULT '', + cgi VARCHAR(128) NOT NULL, + cell_name VARCHAR(255) NOT NULL, + interference_dbm DECIMAL(10,3) NOT NULL, + longitude DECIMAL(10,6) NULL, + latitude DECIMAL(10,6) NULL, + azimuth DECIMAL(6,2) NOT NULL DEFAULT 0, + nearby_count INT NOT NULL DEFAULT 0, + PRIMARY KEY (metric_time, cgi) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 + """, + database=DATABASE_NAME, + ) + columns = self.client.get_json( + self._endpoint("columns"), + {"database": DATABASE_NAME, "table": self.table}, + ) + if not isinstance(columns, list): + raise ProcessingError("Metrix Database API returned invalid column metadata") + existing_columns = {str(item.get("name")) for item in columns if isinstance(item, dict)} + migrations = { + "network_type": "VARCHAR(16) NOT NULL DEFAULT ''", + "longitude": "DECIMAL(10,6) NULL", + "latitude": "DECIMAL(10,6) NULL", + "azimuth": "DECIMAL(6,2) NOT NULL DEFAULT 0", + "nearby_count": "INT NOT NULL DEFAULT 0", + } + for column, definition in migrations.items(): + if column not in existing_columns: + self._query( + f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` {definition}", + database=DATABASE_NAME, + ) + def _endpoint(self, action: str) -> str: return f"/api/databases/{quote(self.connection_id, safe='')}/{action}" - def _query(self, sql: str, database: str = "") -> dict[str, object]: + def _query(self, sql: str, database: str = "", page: int = 1, page_size: int = 100) -> dict[str, object]: payload = self.client.post_json( self._endpoint("query"), - {"sql": sql, "database": database, "page": 1, "page_size": 100}, + {"sql": sql, "database": database, "page": page, "page_size": page_size}, timeout=120, ) if not isinstance(payload, dict): @@ -404,6 +474,59 @@ class ApiSummaryStore: ) + ")" +class ApiHistoryStore: + def __init__(self, client: MetrixApiClient, storage_id: str, root: str = DEFAULT_HISTORY_ROOT) -> None: + self.client = client + self.storage_id = storage_id + self.root = root.rstrip("/") + + def exists(self, metric_time: datetime) -> bool: + directory, filename = self.target(metric_time) + base_entries = self._list_dir(self.root) + if not any(item.get("is_dir") and item.get("path") == directory for item in base_entries): + return False + return any( + not item.get("is_dir") and item.get("name") == filename and int(item.get("size") or 0) > 0 + for item in self._list_dir(directory) + ) + + def upload(self, metric_time: datetime, payload: bytes) -> str: + directory, filename = self.target(metric_time) + base_entries = self._list_dir(self.root) + if not any(item.get("is_dir") and item.get("path") == directory for item in base_entries): + response = self.client.post_json( + self._endpoint("mkdir"), + {"path": directory}, + timeout=120, + ) + if not isinstance(response, dict) or response.get("path") != directory: + raise ProcessingError("Metrix Storage API returned an invalid mkdir result") + response = self.client.post_file( + self._endpoint("upload"), + {"path": directory}, + filename, + payload, + ) + expected_path = posixpath.join(directory, filename) + if not isinstance(response, dict) or response.get("path") != expected_path: + raise ProcessingError("Metrix Storage API returned an invalid upload result") + return expected_path + + def target(self, metric_time: datetime) -> tuple[str, str]: + directory = posixpath.join(self.root, metric_time.strftime("%Y-%m-%d")) + filename = f"干扰数据处理结果_{metric_time:%Y%m%d%H%M%S}.csv" + return directory, filename + + def _endpoint(self, action: str) -> str: + return f"/api/storages/{quote(self.storage_id, safe='')}/{action}" + + def _list_dir(self, path: str) -> list[dict[str, object]]: + payload = self.client.get_json(self._endpoint("files"), {"path": path, "recursive": "false"}) + if not isinstance(payload, dict) or not isinstance(payload.get("entries"), list): + raise ProcessingError("Metrix Storage API returned an invalid file list") + return payload["entries"] + + class ApiSource: def __init__(self, client: MetrixApiClient, storage_id: str, root: str) -> None: self.client = client @@ -486,10 +609,16 @@ class LocalSource: def parse_candidate(path: str, size: int = 0) -> Candidate | None: name = posixpath.basename(path.replace("\\", "/")) - match = FILE_RE.fullmatch(name) - if not match or match.group("source_type") not in EXPECTED_TYPES: + source_type = next((item for item in EXPECTED_TYPES if name.startswith(f"{item}_")), "") + match = FILE_TIME_RE.search(name) + if not source_type or not match: return None - return Candidate(match.group("source_type"), match.group("window"), path, size) + window = match.group("window") + try: + window_start(window) + except ProcessingError: + return None + return Candidate(source_type, window, path, size) def select_latest_cell_data_file(entries: list[dict[str, object]], directory: str) -> dict[str, object]: @@ -568,33 +697,38 @@ def load_cell_metadata(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[st return metadata -def select_window(candidates: list[Candidate], requested: str = "") -> tuple[str, dict[str, Candidate], list[str]]: - grouped: dict[str, dict[str, Candidate]] = {} - duplicates: list[str] = [] - for candidate in candidates: - window_group = grouped.setdefault(candidate.window, {}) - if candidate.source_type in window_group: - duplicates.append(f"{candidate.window}/{candidate.source_type}") - window_group[candidate.source_type] = candidate - if duplicates: - raise ProcessingError(f"Duplicate source files: {', '.join(sorted(duplicates))}") - complete = sorted(window for window, items in grouped.items() if all(name in items for name in EXPECTED_TYPES)) - if requested: - if requested not in complete: - present = sorted(grouped.get(requested, {})) - missing = [name for name in EXPECTED_TYPES if name not in present] - raise ProcessingError(f"Requested window is incomplete: {requested}; missing={missing}") - selected = requested - elif complete: - selected = complete[-1] - else: - raise ProcessingError("No hour contains all seven interference source types") - warnings_out: list[str] = [] - latest_seen = max(grouped) if grouped else "" - if latest_seen and latest_seen != selected: - missing = [name for name in EXPECTED_TYPES if name not in grouped[latest_seen]] - warnings_out.append(f"Latest observed window {latest_seen} is incomplete; using {selected}; missing={missing}") - return selected, grouped[selected], warnings_out +def select_window(candidates: list[Candidate], requested: str = "") -> tuple[datetime | None, dict[str, Candidate], list[str]]: + if not candidates: + return None, {}, list(EXPECTED_TYPES) + target_time = metric_time_for_window(requested) if requested else metric_time_for_window(max(candidates, key=candidate_key).window) + grouped = [candidate for candidate in candidates if metric_time_for_window(candidate.window) == target_time] + selected: dict[str, Candidate] = {} + for source_type in EXPECTED_TYPES: + matching = [candidate for candidate in grouped if candidate.source_type == source_type] + if matching: + selected[source_type] = max(matching, key=candidate_key) + missing = [source_type for source_type in EXPECTED_TYPES if source_type not in selected] + return target_time, selected, missing + + +def candidate_key(candidate: Candidate) -> tuple[datetime, str, str]: + return window_start(candidate.window), candidate.window, candidate.path + + +def window_start(window: str) -> datetime: + if not re.fullmatch(r"\d{16}", window): + raise ProcessingError(f"Invalid window: {window}") + try: + start = datetime.strptime(window[:12], "%Y%m%d%H%M") + datetime.strptime(window[12:], "%H%M") + except ValueError as exc: + raise ProcessingError(f"Invalid window: {window}") from exc + return start + + +def metric_time_for_window(window: str) -> datetime: + start = window_start(window) + return start.replace(minute=(start.minute // 15) * 15, second=0, microsecond=0) def parse_workbook(raw_zip: bytes, source_type: str) -> tuple[tuple[str, ...], list[tuple[object, ...]], str]: @@ -637,13 +771,43 @@ def process( lookback_days: int, requested_window: str = "", store: SummaryStore | None = None, -) -> Path: + history: HistoryStore | None = None, +) -> Path | None: candidates = source.candidates(lookback_days) - window, selected, warnings_out = select_window(candidates, requested_window) + metric_time, selected, missing = select_window(candidates, requested_window) + if metric_time is None: + print("status=waiting") + print("reason=no_source_files") + return None + metric_time_text = metric_time.strftime("%Y-%m-%d %H:%M:%S") + if missing: + print("status=waiting") + print(f"target_time={metric_time_text}") + print(f"missing_sources={','.join(missing)}") + return None + + if store is not None: + latest_time = store.latest_metric_time() + if latest_time is not None and latest_time >= metric_time: + if latest_time == metric_time and history is not None and not history.exists(metric_time): + stored_rows = store.rows_for_time(metric_time) + if not stored_rows: + raise ProcessingError(f"Database contains no rows for {metric_time_text}") + history_path = history.upload(metric_time, dict_csv_bytes(SUMMARY_HEADER, stored_rows)) + print("status=history_repaired") + print(f"target_time={metric_time_text}") + print(f"history_output={history_path}") + else: + print("status=skipped") + print(f"target_time={metric_time_text}") + print(f"database_metric_time={latest_time:%Y-%m-%d %H:%M:%S}") + return None + cell_data_workbooks = source.cell_data_workbooks() cell_metadata = load_cell_metadata(cell_data_workbooks) - temp_dir = output_root.resolve() / f".{window}.tmp-{os.getpid()}" - final_dir = output_root.resolve() / window + output_key = metric_time.strftime("%Y%m%d%H%M%S") + temp_dir = output_root.resolve() / f".{output_key}.tmp-{os.getpid()}" + final_dir = output_root.resolve() / output_key ensure_scoped(output_root.resolve(), temp_dir) if temp_dir.exists(): shutil.rmtree(temp_dir) @@ -654,16 +818,17 @@ def process( threshold_filtered_rows = 0 manifest_files: list[dict[str, object]] = [] database_result: dict[str, object] = {"enabled": False} + history_result: dict[str, object] = {"enabled": False} try: for source_type in EXPECTED_TYPES: candidate = selected[source_type] raw_zip = source.download(candidate.path) header, rows, member_name = parse_workbook(raw_zip, source_type) - csv_name = f"{source_type}_{window}.csv" + csv_name = f"{source_type}_{candidate.window}.csv" write_csv(converted_dir / csv_name, header, rows) for row_number, row in enumerate(rows, start=2): record = dict(zip(header, row, strict=True)) - summary = summary_record(record, source_type, candidate.path, window, row_number, cell_metadata) + summary = summary_record(record, source_type, candidate.path, metric_time, row_number, cell_metadata) if passes_interference_threshold(summary): summary_rows.append(summary) else: @@ -672,6 +837,7 @@ def process( { "source_size": candidate.size, "sha256": hashlib.sha256(raw_zip).hexdigest(), + "source_window": candidate.window, "xlsx_member": member_name, "rows": len(rows), "converted_csv": f"converted/{csv_name}", @@ -680,26 +846,27 @@ def process( populate_nearby_counts(summary_rows) summary_rows.sort(key=lambda item: (item["cgi"], item["cell_name"])) - summary_name = f"interference_summary_{window}.csv" + summary_name = f"interference_summary_{output_key}.csv" write_dict_csv(temp_dir / summary_name, SUMMARY_HEADER, summary_rows) if store is not None: database_result = store.replace_latest(summary_rows) + if history is not None: + history_path = history.upload(metric_time, dict_csv_bytes(SUMMARY_HEADER, summary_rows)) + history_result = {"enabled": True, "path": history_path} matched_coordinates = sum(bool(row["longitude"] and row["latitude"]) for row in summary_rows) manifest = { "generated_at": datetime.now(timezone.utc).isoformat(), - "window": window, - "hour_start": window_bounds(window)[0], - "hour_end": window_bounds(window)[1], + "metric_time": metric_time_text, "source_file_count": len(manifest_files), "cell_data_file_count": len(cell_data_workbooks), "summary_rows": len(summary_rows), "threshold_filtered_rows": threshold_filtered_rows, "coordinate_matched_rows": matched_coordinates, "coordinate_unmatched_rows": len(summary_rows) - matched_coordinates, - "warnings": warnings_out, "files": manifest_files, "summary_csv": summary_name, "database": database_result, + "history": history_result, } (temp_dir / "manifest.json").write_text(json.dumps(manifest, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") output_root.resolve().mkdir(parents=True, exist_ok=True) @@ -711,7 +878,8 @@ def process( shutil.rmtree(temp_dir, ignore_errors=True) raise - print(f"selected_window={window}") + print("status=processed") + print(f"target_time={metric_time_text}") print(f"source_files={len(manifest_files)}") print(f"cell_data_files={len(cell_data_workbooks)}") print(f"summary_rows={len(summary_rows)}") @@ -722,8 +890,8 @@ def process( print(f"database_metric_time={database_result['metric_time']}") print(f"database_inserted_rows={database_result['inserted_rows']}") print(f"database_old_rows_deleted={database_result['old_rows_deleted']}") - for message in warnings_out: - print(f"warning={message}") + if history_result["enabled"]: + print(f"history_output={history_result['path']}") print(f"output={final_dir}") return final_dir @@ -732,14 +900,21 @@ def summary_record( record: dict[str, object], source_type: str, source_path: str, - window: str, + metric_time: datetime, row_number: int, cell_metadata: dict[str, tuple[str, str, str]] | None = None, ) -> dict[str, str]: - hour_start = normalize_cell(record["开始时间"]) - expected_start, hour_end = window_bounds(window) - if hour_start != expected_start: - raise ProcessingError(f"{source_path}: row {row_number} time {hour_start!r} does not match {expected_start!r}") + row_time_text = normalize_cell(record["开始时间"]) + try: + row_time = datetime.strptime(row_time_text, "%Y-%m-%d %H:%M:%S") + except ValueError as exc: + raise ProcessingError(f"{source_path}: row {row_number} has invalid time {row_time_text!r}") from exc + row_metric_time = row_time.replace(minute=(row_time.minute // 15) * 15, second=0, microsecond=0) + if row_metric_time != metric_time: + raise ProcessingError( + f"{source_path}: row {row_number} time {row_time_text!r} is outside metric group " + f"{metric_time:%Y-%m-%d %H:%M:%S}" + ) if source_type in ("5G干扰监控", "700M干扰监控"): cgi = build_cgi("", record, ("gNBplmn", "gNBId", "cellId"), source_path, row_number) @@ -760,8 +935,7 @@ def summary_record( longitude, latitude, azimuth = (cell_metadata or {}).get(cgi, ("", "", "0")) return { - "hour_start": hour_start, - "hour_end": hour_end, + "metric_time": metric_time.strftime("%Y-%m-%d %H:%M:%S"), "network_type": NETWORK_TYPE_BY_SOURCE[source_type], "cgi": cgi, "cell_name": cell_name, @@ -850,18 +1024,6 @@ def required(record: dict[str, object], column: str, path: str, row: int) -> str return value -def window_bounds(window: str) -> tuple[str, str]: - if not re.fullmatch(r"\d{16}", window): - raise ProcessingError(f"Invalid window: {window}") - start = datetime.strptime(window[:12], "%Y%m%d%H%M") - end_hour = int(window[12:14]) - end_minute = int(window[14:16]) - end = start.replace(hour=end_hour, minute=end_minute) - if end <= start: - end += timedelta(days=1) - return start.strftime("%Y-%m-%d %H:%M:%S"), end.strftime("%Y-%m-%d %H:%M:%S") - - def normalize_cell(value: object) -> str: if value is None: return "" @@ -899,10 +1061,15 @@ def write_csv(path: Path, header: tuple[str, ...], rows: list[tuple[object, ...] def write_dict_csv(path: Path, header: tuple[str, ...], rows: list[dict[str, str]]) -> None: - with path.open("w", encoding="utf-8-sig", newline="") as file: - writer = csv.DictWriter(file, fieldnames=header, extrasaction="raise") - writer.writeheader() - writer.writerows(rows) + path.write_bytes(dict_csv_bytes(header, rows)) + + +def dict_csv_bytes(header: tuple[str, ...], rows: list[dict[str, str]]) -> bytes: + output = io.StringIO(newline="") + writer = csv.DictWriter(output, fieldnames=header, extrasaction="raise") + writer.writeheader() + writer.writerows(rows) + return output.getvalue().encode("utf-8-sig") def ensure_scoped(root: Path, target: Path) -> None: @@ -911,10 +1078,14 @@ def ensure_scoped(root: Path, target: Path) -> None: def build_parser() -> argparse.ArgumentParser: - parser = argparse.ArgumentParser(description="Process the latest complete hour of interference KPI files") + parser = argparse.ArgumentParser(description="Process the latest complete interference KPI time group") parser.add_argument("--source-dir", type=Path, help="Use a local source tree instead of the Metrix storage API") parser.add_argument("--output-dir", type=Path, default=Path(os.getenv("INTERFERENCE_OUTPUT_DIR", "output"))) - parser.add_argument("--window", default=os.getenv("INTERFERENCE_WINDOW", ""), help="Optional exact 16-digit source window") + parser.add_argument( + "--window", + default=os.getenv("INTERFERENCE_WINDOW", ""), + help="Optional 16-digit source window used to select a natural 15-minute group", + ) parser.add_argument("--lookback-days", type=int, default=int(os.getenv("INTERFERENCE_LOOKBACK_DAYS", "3"))) parser.add_argument("--storage-id", default=os.getenv("METRIX_STORAGE_ID", DEFAULT_STORAGE_ID)) parser.add_argument("--root", default=os.getenv("INTERFERENCE_SOURCE_ROOT", DEFAULT_ROOT)) @@ -934,10 +1105,12 @@ def main(argv: list[str] | None = None) -> int: source = ApiSource(client, args.storage_id, args.root) if args.no_database: store = None + history = None else: client = client or MetrixApiClient(METRIX_API_BASE_URL, METRIX_API_TOKEN) store = ApiSummaryStore(client, METRIX_DATABASE_CONNECTION_ID) - process(source, args.output_dir, args.lookback_days, args.window, store=store) + history = ApiHistoryStore(client, args.storage_id, f"{args.root.rstrip('/')}/干扰历史数据") + process(source, args.output_dir, args.lookback_days, args.window, store=store, history=history) return 0 diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index cd3a7f4..0040eba 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -15,32 +15,36 @@ from openpyxl import Workbook COMPLETE_WINDOW = "2026073110001100" INCOMPLETE_WINDOW = "2026073111001200" +TARGET_KEY = "20260731100000" class PipelineTest(unittest.TestCase): - def test_latest_complete_window_is_converted_and_merged(self) -> None: + def test_latest_group_is_converted_merged_stored_and_archived(self) -> None: with tempfile.TemporaryDirectory() as temp: root = Path(temp) / "source" output = Path(temp) / "output" - for source_type in main.EXPECTED_TYPES: - create_archive(root, source_type, COMPLETE_WINDOW) - create_archive(root, main.EXPECTED_TYPES[0], INCOMPLETE_WINDOW) + windows = ["2026073110001100", "2026073110141100", "2026073110051100"] + for index, source_type in enumerate(main.EXPECTED_TYPES): + create_archive(root, source_type, windows[index % len(windows)], middle="任意粒度") + create_archive(root, main.EXPECTED_TYPES[0], "2026073110091100", middle="更新版本") create_cell_data_sources(root) store = RecordingStore() + history = RecordingHistory() - result = main.process(main.LocalSource(root), output, lookback_days=3, store=store) + result = main.process(main.LocalSource(root), output, lookback_days=3, store=store, history=history) - self.assertEqual(result.name, COMPLETE_WINDOW) + self.assertIsNotNone(result) + assert result is not None + self.assertEqual(result.name, TARGET_KEY) converted = sorted((result / "converted").glob("*.csv")) self.assertEqual(len(converted), 7) - with (result / f"interference_summary_{COMPLETE_WINDOW}.csv").open(encoding="utf-8-sig", newline="") as file: + with (result / f"interference_summary_{TARGET_KEY}.csv").open(encoding="utf-8-sig", newline="") as file: rows = list(csv.DictReader(file)) self.assertEqual(len(rows), 7) self.assertEqual( list(rows[0]), [ - "hour_start", - "hour_end", + "metric_time", "network_type", "cgi", "cell_name", @@ -55,6 +59,7 @@ class PipelineTest(unittest.TestCase): self.assertEqual(rows_by_type["5G干扰监控"]["network_type"], "2.6G") self.assertEqual(rows_by_type["700M干扰监控"]["network_type"], "700M") self.assertEqual(rows_by_type["SDR_FDD干扰监控"]["network_type"], "4G") + self.assertTrue(all(row["metric_time"] == "2026-07-31 10:00:00" for row in rows)) self.assertEqual(rows_by_type["5G干扰监控"]["cgi"], "460-00-200-1") self.assertEqual(rows_by_type["700M干扰监控"]["cgi"], "460-00-200-1") for source_type in set(main.EXPECTED_TYPES) - {"5G干扰监控", "700M干扰监控"}: @@ -76,12 +81,18 @@ class PipelineTest(unittest.TestCase): self.assertEqual(manifest["threshold_filtered_rows"], 0) self.assertEqual(manifest["coordinate_matched_rows"], 2) self.assertEqual(manifest["coordinate_unmatched_rows"], 5) - self.assertEqual(len(manifest["warnings"]), 1) - self.assertIn(INCOMPLETE_WINDOW, manifest["warnings"][0]) + self.assertEqual(manifest["metric_time"], "2026-07-31 10:00:00") + self.assertNotIn("hour_start", manifest) + self.assertNotIn("hour_end", manifest) self.assertNotIn("source_types", manifest) self.assertTrue(all("source_type" not in item and "source_path" not in item for item in manifest["files"])) + self.assertIn("2026073110091100", [item["source_window"] for item in manifest["files"]]) self.assertEqual(manifest["database"]["metric_time"], "2026-07-31 10:00:00") self.assertEqual(len(store.rows), 7) + self.assertEqual(history.path, "/history/2026-07-31/干扰数据处理结果_20260731100000.csv") + archived_rows = list(csv.DictReader(io.StringIO(history.payload.decode("utf-8-sig")))) + self.assertEqual(list(archived_rows[0]), list(main.SUMMARY_HEADER)) + self.assertEqual(len(archived_rows), 7) def test_interference_thresholds_remove_only_lower_values(self) -> None: for network_type, threshold in (("2.6G", "-107"), ("700M", "-110"), ("4G", "-110")): @@ -105,13 +116,30 @@ class PipelineTest(unittest.TestCase): self.assertEqual(manifest["summary_rows"], 6) self.assertEqual(manifest["threshold_filtered_rows"], 1) - def test_requested_incomplete_window_is_rejected(self) -> None: + def test_latest_incomplete_group_waits_without_falling_back(self) -> None: with tempfile.TemporaryDirectory() as temp: root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive(root, source_type, COMPLETE_WINDOW) create_archive(root, main.EXPECTED_TYPES[0], INCOMPLETE_WINDOW) + store = RecordingStore() - with self.assertRaisesRegex(main.ProcessingError, "Requested window is incomplete"): - main.process(main.LocalSource(root), Path(temp) / "output", 3, INCOMPLETE_WINDOW) + result = main.process(main.LocalSource(root), Path(temp) / "output", 3, store=store) + + self.assertIsNone(result) + self.assertEqual(store.latest_calls, 0) + self.assertEqual(store.rows, []) + + def test_no_source_files_waits_successfully(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + root.mkdir() + store = RecordingStore() + + result = main.process(main.LocalSource(root), Path(temp) / "output", 3, store=store) + + self.assertIsNone(result) + self.assertEqual(store.latest_calls, 0) def test_schema_change_is_rejected(self) -> None: with tempfile.TemporaryDirectory() as temp: @@ -122,12 +150,34 @@ class PipelineTest(unittest.TestCase): with self.assertRaisesRegex(main.ProcessingError, "Unexpected Sheet0 header"): main.process(main.LocalSource(root), Path(temp) / "output", 3) - def test_cross_midnight_window(self) -> None: - self.assertEqual( - main.window_bounds("2026073023000000"), - ("2026-07-30 23:00:00", "2026-07-31 00:00:00"), + def test_workbook_row_outside_selected_quarter_is_rejected(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive( + root, + source_type, + COMPLETE_WINDOW, + row_window="2026073110151100" if source_type == main.EXPECTED_TYPES[0] else "", + ) + create_cell_data_sources(root) + + with self.assertRaisesRegex(main.ProcessingError, "is outside metric group"): + main.process(main.LocalSource(root), Path(temp) / "output", 3) + + def test_file_name_middle_is_flexible_and_time_uses_natural_quarter(self) -> None: + candidate = main.parse_candidate( + "/source/700M干扰监控_任意描述_2026073112141300.zip", + 123, ) + self.assertIsNotNone(candidate) + assert candidate is not None + self.assertEqual(candidate.source_type, "700M干扰监控") + self.assertEqual(candidate.size, 123) + self.assertEqual(main.metric_time_for_window(candidate.window), datetime(2026, 7, 31, 12, 0)) + self.assertEqual(main.metric_time_for_window("2026073112151300"), datetime(2026, 7, 31, 12, 15)) + def test_latest_dated_cell_data_file_is_selected(self) -> None: entries = [ {"name": "江门5G小区信息表20260727.xlsx", "path": "/old.xlsx", "is_dir": False}, @@ -222,22 +272,93 @@ class PipelineTest(unittest.TestCase): with self.assertRaisesRegex(main.ProcessingError, "Invalid decimal value for interference_dbm"): store.replace_latest([row]) + def test_api_store_exports_normalized_rows_for_history(self) -> None: + client = FakeApiClient(None) + store = main.ApiSummaryStore(client, "db_share_mysql") + + rows = store.rows_for_time(datetime(2026, 7, 31, 10, 0)) + + self.assertEqual(rows, [database_row()]) + export_query = next( + payload for endpoint, payload in client.posts + if endpoint.endswith("/query") and str(payload["sql"]).startswith("SELECT metric_time") + ) + self.assertEqual(export_query["page_size"], 1000) + + def test_existing_database_time_repairs_missing_history_without_reprocessing(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive(root, source_type, COMPLETE_WINDOW) + stored_row = database_row() + store = RecordingStore(datetime(2026, 7, 31, 10, 0), [stored_row]) + history = RecordingHistory(exists=False) + + result = main.process(main.LocalSource(root), Path(temp) / "output", 3, store=store, history=history) + + self.assertIsNone(result) + self.assertEqual(store.rows_for_time_calls, 1) + self.assertEqual(history.path, "/history/2026-07-31/干扰数据处理结果_20260731100000.csv") + archived_rows = list(csv.DictReader(io.StringIO(history.payload.decode("utf-8-sig")))) + self.assertEqual(archived_rows, [stored_row]) + + def test_api_history_store_creates_date_directory_and_uploads_csv(self) -> None: + client = FakeHistoryApiClient() + history = main.ApiHistoryStore(client, "stg_test", "/history") + metric_time = datetime(2026, 8, 6, 12, 0) + + self.assertFalse(history.exists(metric_time)) + path = history.upload(metric_time, b"csv-data") + + self.assertEqual(path, "/history/2026-08-06/干扰数据处理结果_20260806120000.csv") + self.assertEqual(client.mkdir_paths, ["/history/2026-08-06"]) + self.assertEqual(client.uploads, [("/history/2026-08-06", "干扰数据处理结果_20260806120000.csv", b"csv-data")]) + class RecordingStore: - def __init__(self) -> None: + def __init__(self, latest_time: datetime | None = None, stored_rows: list[dict[str, str]] | None = None) -> None: self.rows: list[dict[str, str]] = [] + self.latest_time = latest_time + self.stored_rows = stored_rows or [] + self.latest_calls = 0 + self.rows_for_time_calls = 0 + + def latest_metric_time(self) -> datetime | None: + self.latest_calls += 1 + return self.latest_time def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: self.rows = list(rows) return { "enabled": True, "table": main.DATABASE_TABLE, - "metric_time": rows[0]["hour_start"], + "metric_time": rows[0]["metric_time"], "inserted_rows": len(rows), "refreshed_rows": 0, "old_rows_deleted": 0, } + def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: + self.rows_for_time_calls += 1 + self.asserted_metric_time = metric_time + return list(self.stored_rows) + + +class RecordingHistory: + def __init__(self, exists: bool = False) -> None: + self.exists_value = exists + self.path = "" + self.payload = b"" + + def exists(self, metric_time: datetime) -> bool: + self.checked_metric_time = metric_time + return self.exists_value + + def upload(self, metric_time: datetime, payload: bytes) -> str: + self.payload = payload + self.path = f"/history/{metric_time:%Y-%m-%d}/干扰数据处理结果_{metric_time:%Y%m%d%H%M%S}.csv" + return self.path + class FakeApiClient: def __init__(self, latest_time: str | None, fail_script: bool = False) -> None: @@ -258,8 +379,13 @@ class FakeApiClient: del timeout self.posts.append((endpoint, payload)) if endpoint.endswith("/query"): - if str(payload["sql"]).lstrip().startswith("SELECT MAX"): + sql = str(payload["sql"]).lstrip() + if sql.startswith("SELECT MAX"): return {"rows": [{"latest_time": self.latest_time}]} + if sql.startswith("SELECT metric_time"): + row = database_row() + row["metric_time"] = "2026-07-31T10:00:00" + return {"rows": [row], "total": 1} return {"affected_rows": 0} if self.fail_script: return { @@ -282,10 +408,41 @@ class FakeApiClient: } +class FakeHistoryApiClient: + def __init__(self) -> None: + self.mkdir_paths: list[str] = [] + self.uploads: list[tuple[str, str, bytes]] = [] + + def get_json(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> object: + del endpoint, timeout + path = str((query or {}).get("path") or "") + if path == "/history": + return {"entries": []} + return {"entries": []} + + def post_json(self, endpoint: str, payload: dict[str, object], timeout: int = 30) -> object: + del endpoint, timeout + path = str(payload["path"]) + self.mkdir_paths.append(path) + return {"path": path} + + def post_file( + self, + endpoint: str, + query: dict[str, object], + filename: str, + payload: bytes, + timeout: int = 120, + ) -> object: + del endpoint, timeout + directory = str(query["path"]) + self.uploads.append((directory, filename, payload)) + return {"path": f"{directory}/{filename}"} + + def database_row() -> dict[str, str]: return { - "hour_start": "2026-07-31 10:00:00", - "hour_end": "2026-07-31 11:00:00", + "metric_time": "2026-07-31 10:00:00", "network_type": "2.6G", "cgi": "460-00-200-1", "cell_name": "测试小区", @@ -311,10 +468,12 @@ def create_archive( window: str, bad_header: bool = False, interference_dbm: float = -100.5, + middle: str = "LWP_每小时_过滤110", + row_window: str = "", ) -> None: date_dir = root / f"{window[:4]}-{window[4:6]}-{window[6:8]}" date_dir.mkdir(parents=True, exist_ok=True) - filename = f"{source_type}_LWP_每小时_过滤110_{window}" + filename = f"{source_type}_{middle}_{window}" workbook = Workbook() sheet = workbook.active sheet.title = "Sheet0" @@ -322,7 +481,7 @@ def create_archive( if bad_header: header[-1] = "unexpected" sheet.append(header) - sheet.append(mock_row(source_type, window, interference_dbm)) + sheet.append(mock_row(source_type, row_window or window, interference_dbm)) metadata = workbook.create_sheet("指标(计数器)") metadata.append(["指标或计数器", "指标或计数器描述", "指标公式", "指标或计数器状态"]) content = io.BytesIO()