diff --git a/README.md b/README.md index 323b9e0..a6c6b96 100644 --- a/README.md +++ b/README.md @@ -16,15 +16,17 @@ output/2026073110001100/ └── 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 to matching interference rows. Unmatched rows keep empty coordinates. + The summary columns are: ```text -hour_start,hour_end,cgi,cell_name,interference_dbm +hour_start,hour_end,cgi,cell_name,interference_dbm,longitude,latitude ``` The script never modifies or deletes source storage files. By default it writes the complete selected hour to MySQL and keeps only that hour in the target table. -The database table contains one time column, `metric_time DATETIME`, which is the source KPI start time. 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 one time column, `metric_time DATETIME`, which is the source KPI start time, plus nullable `longitude` and `latitude` columns. 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. CGI is generated with fixed rules: @@ -33,6 +35,8 @@ CGI is generated with fixed rules: MySQL connection values, database `interference_etl`, and table `interference_hourly_summary` are constants at the top of `main.py`; they are not project environment variables. The script creates its database and table automatically when they do not exist. +The five CellData directories are fixed in `CELL_DATA_DIRECTORIES`. In each directory, only XLSX files ending with a valid `YYYYMMDD.xlsx` date are considered, and the latest date is selected. + ## Local development ```powershell diff --git a/aidocs/project_context.md b/aidocs/project_context.md index fe92d79..de3b891 100644 --- a/aidocs/project_context.md +++ b/aidocs/project_context.md @@ -39,3 +39,10 @@ - The default MySQL path first connects without selecting a database, runs `CREATE DATABASE IF NOT EXISTS interference_etl` with `utf8mb4`, then reconnects to that database and creates the table if needed. The configured account therefore needs database creation permission on first run. - Offline dependencies are stored below the workspace `vendor/` directory instead of separate package directories at the project root. `main.py` prepends this directory to `sys.path`, so the run command remains `python main.py`. - Summary CSV rows and database rows do not expose `source_type` or `source_path`. The database primary key is `(metric_time, cgi)`; the current online hour was checked for CGI uniqueness before migrating from the former source-aware key. + +## 2026-08-03: CellData coordinate enrichment + +- Each run selects the latest filename-dated XLSX from five fixed CellData directories for 5G, 700M, reverse-activated 5G, TDD LTE, and FDD LTE. The filename date must end in `YYYYMMDD.xlsx`; source files remain read-only. +- CellData workbooks use sheet `小区信息表` and required columns `eNB/gNB`, `CI`, `经度`, and `纬度`. CellData CGI is always `460-00-{eNB/gNB}-{CI}` and conflicting coordinates for the same CGI stop the run instead of silently overwriting data. +- The summary CSV now ends with `longitude,latitude`. The script-owned `interference_hourly_summary` table has nullable `DECIMAL(10,6)` columns with the same names; existing tables are migrated automatically and unmatched CGI values are stored as `NULL`. +- Read-only validation of the five `20260728` workbooks produced 66,250 unique CGI coordinates with no conflicts. Against the sampled latest-hour summary, 1,529 of 1,539 rows matched (99.35%); the remaining 10 rows correctly stay empty. diff --git a/main.py b/main.py index 51abc20..4dbd862 100644 --- a/main.py +++ b/main.py @@ -40,7 +40,15 @@ EXPECTED_TYPES = ( DEFAULT_STORAGE_ID = "stg_4d9a910d72" DEFAULT_ROOT = "/网优日常优化数据文档/(勿删)干扰定时小时指标" +CELL_DATA_DIRECTORIES = ( + "/网优日常优化数据文档/日常性能报表/2026年/5G/5G小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/700M/700M小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/5G_反开/反开小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/TDD_LTE/LTE基础信息数据/LTE_小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/FDD_LTE/FDD基础信息数据/FDD小区信息表", +) FILE_RE = re.compile(r"^(?P.+)_LWP_每小时_过滤110_(?P\d{16})\.zip$", re.IGNORECASE) +CELL_DATA_FILE_RE = re.compile(r"(?P20\d{6})\.xlsx$", re.IGNORECASE) LTE_PLMN = "460-00" DATABASE_HOST = "172.17.0.1" @@ -146,6 +154,8 @@ SUMMARY_HEADER = ( "cgi", "cell_name", "interference_dbm", + "longitude", + "latitude", ) class ProcessingError(RuntimeError): @@ -165,6 +175,8 @@ class Source(Protocol): def download(self, path: str) -> bytes: ... + def cell_data_workbooks(self) -> list[tuple[str, bytes]]: ... + class SummaryStore(Protocol): def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: ... @@ -194,10 +206,17 @@ class MySQLSummaryStore: 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, PRIMARY KEY (metric_time, cgi) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 """ ) + cursor.execute(f"SHOW COLUMNS FROM {table}") + existing_columns = {column[0] for column in cursor.fetchall()} + for column in ("longitude", "latitude"): + if column not in existing_columns: + cursor.execute(f"ALTER TABLE {table} ADD COLUMN `{column}` DECIMAL(10,6) NULL") cursor.execute(f"SELECT MAX(metric_time) FROM {table}") latest_row = cursor.fetchone() latest_time = latest_row[0] if latest_row else None @@ -214,8 +233,8 @@ class MySQLSummaryStore: cursor.executemany( f""" INSERT INTO {table} - (metric_time, cgi, cell_name, interference_dbm) - VALUES (%s, %s, %s, %s) + (metric_time, cgi, cell_name, interference_dbm, longitude, latitude) + VALUES (%s, %s, %s, %s, %s, %s) """, [ ( @@ -223,6 +242,8 @@ class MySQLSummaryStore: row["cgi"], row["cell_name"], row["interference_dbm"], + row["longitude"] or None, + row["latitude"] or None, ) for row in rows ], @@ -326,6 +347,14 @@ class ApiSource: endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/download" return self._request(endpoint, {"path": path}, timeout=120) + def cell_data_workbooks(self) -> list[tuple[str, bytes]]: + workbooks: list[tuple[str, bytes]] = [] + for directory in CELL_DATA_DIRECTORIES: + item = select_latest_cell_data_file(self._list_dir(directory), directory) + path = str(item["path"]) + workbooks.append((path, self.download(path))) + return workbooks + def _list_dir(self, path: str) -> list[dict[str, object]]: endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/files" payload = self._request(endpoint, {"path": path, "recursive": "false"}) @@ -369,6 +398,24 @@ class LocalSource: def download(self, path: str) -> bytes: return Path(path).read_bytes() + def cell_data_workbooks(self) -> list[tuple[str, bytes]]: + directories = [self.root / Path(directory.lstrip("/")) for directory in CELL_DATA_DIRECTORIES] + existing = [directory for directory in directories if directory.is_dir()] + if not existing: + return [] + if len(existing) != len(directories): + raise ProcessingError("Local CellData source must contain all five configured directories") + workbooks: list[tuple[str, bytes]] = [] + for directory in directories: + entries = [ + {"name": path.name, "path": str(path.resolve()), "is_dir": path.is_dir()} + for path in directory.iterdir() + ] + item = select_latest_cell_data_file(entries, str(directory)) + path = str(item["path"]) + workbooks.append((path, self.download(path))) + return workbooks + def parse_candidate(path: str, size: int = 0) -> Candidate | None: name = posixpath.basename(path.replace("\\", "/")) @@ -378,6 +425,75 @@ def parse_candidate(path: str, size: int = 0) -> Candidate | None: return Candidate(match.group("source_type"), match.group("window"), path, size) +def select_latest_cell_data_file(entries: list[dict[str, object]], directory: str) -> dict[str, object]: + dated_files: list[tuple[datetime, str, dict[str, object]]] = [] + for item in entries: + if item.get("is_dir"): + continue + name = str(item.get("name") or "") + match = CELL_DATA_FILE_RE.search(name) + if not match: + continue + try: + file_date = datetime.strptime(match.group("date"), "%Y%m%d") + except ValueError: + continue + dated_files.append((file_date, name, item)) + if not dated_files: + raise ProcessingError(f"No dated CellData XLSX found in {directory}") + return max(dated_files, key=lambda entry: (entry[0], entry[1]))[2] + + +def parse_cell_data_workbook(raw_xlsx: bytes, path: str) -> dict[str, tuple[str, str]]: + try: + with warnings.catch_warnings(): + warnings.filterwarnings("ignore", message="Workbook contains no default style") + workbook = load_workbook(io.BytesIO(raw_xlsx), read_only=True, data_only=True) + except Exception as exc: + raise ProcessingError(f"Invalid CellData workbook {path}: {exc}") from exc + try: + if "小区信息表" not in workbook.sheetnames: + raise ProcessingError(f"CellData workbook does not contain 小区信息表: {path}") + rows = workbook["小区信息表"].iter_rows(values_only=True) + try: + header = tuple(normalize_cell(value) for value in next(rows)) + except StopIteration as exc: + raise ProcessingError(f"CellData workbook is empty: {path}") from exc + required_columns = ("eNB/gNB", "CI", "经度", "纬度") + missing = [column for column in required_columns if column not in header] + if missing: + raise ProcessingError(f"CellData workbook missing columns {missing}: {path}") + indexes = {column: header.index(column) for column in required_columns} + coordinates: dict[str, tuple[str, str]] = {} + for row_number, row in enumerate(rows, start=2): + node = normalize_cell(row[indexes["eNB/gNB"]]) + cell = normalize_cell(row[indexes["CI"]]) + longitude = normalize_cell(row[indexes["经度"]]) + latitude = normalize_cell(row[indexes["纬度"]]) + if not node or not cell or not longitude or not latitude: + continue + cgi = f"{LTE_PLMN}-{node}-{cell}" + value = (longitude, latitude) + existing = coordinates.get(cgi) + if existing is not None and existing != value: + raise ProcessingError(f"Conflicting CellData coordinates for {cgi} in {path}, row {row_number}") + coordinates[cgi] = value + return coordinates + finally: + workbook.close() + + +def load_cell_coordinates(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[str, str]]: + coordinates: dict[str, tuple[str, str]] = {} + for path, raw_xlsx in workbooks: + for cgi, value in parse_cell_data_workbook(raw_xlsx, path).items(): + existing = coordinates.get(cgi) + if existing is not None and existing != value: + raise ProcessingError(f"Conflicting CellData coordinates for {cgi} across workbooks") + coordinates[cgi] = value + return coordinates + + def select_window(candidates: list[Candidate], requested: str = "") -> tuple[str, dict[str, Candidate], list[str]]: grouped: dict[str, dict[str, Candidate]] = {} duplicates: list[str] = [] @@ -450,6 +566,8 @@ def process( ) -> Path: candidates = source.candidates(lookback_days) window, selected, warnings_out = select_window(candidates, requested_window) + cell_data_workbooks = source.cell_data_workbooks() + cell_coordinates = load_cell_coordinates(cell_data_workbooks) temp_dir = output_root.resolve() / f".{window}.tmp-{os.getpid()}" final_dir = output_root.resolve() / window ensure_scoped(output_root.resolve(), temp_dir) @@ -470,7 +588,9 @@ def process( 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_rows.append(summary_record(record, source_type, candidate.path, window, row_number)) + summary_rows.append( + summary_record(record, source_type, candidate.path, window, row_number, cell_coordinates) + ) manifest_files.append( { "source_size": candidate.size, @@ -486,13 +606,17 @@ def process( write_dict_csv(temp_dir / summary_name, SUMMARY_HEADER, summary_rows) if store is not None: database_result = store.replace_latest(summary_rows) + 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], "source_file_count": len(manifest_files), + "cell_data_file_count": len(cell_data_workbooks), "summary_rows": len(summary_rows), + "coordinate_matched_rows": matched_coordinates, + "coordinate_unmatched_rows": len(summary_rows) - matched_coordinates, "warnings": warnings_out, "files": manifest_files, "summary_csv": summary_name, @@ -510,7 +634,10 @@ def process( print(f"selected_window={window}") print(f"source_files={len(manifest_files)}") + print(f"cell_data_files={len(cell_data_workbooks)}") print(f"summary_rows={len(summary_rows)}") + print(f"coordinate_matched_rows={matched_coordinates}") + print(f"coordinate_unmatched_rows={len(summary_rows) - matched_coordinates}") if database_result["enabled"]: print(f"database_metric_time={database_result['metric_time']}") print(f"database_inserted_rows={database_result['inserted_rows']}") @@ -522,7 +649,12 @@ def process( def summary_record( - record: dict[str, object], source_type: str, source_path: str, window: str, row_number: int + record: dict[str, object], + source_type: str, + source_path: str, + window: str, + row_number: int, + cell_coordinates: dict[str, tuple[str, str]] | None = None, ) -> dict[str, str]: hour_start = normalize_cell(record["开始时间"]) expected_start, hour_end = window_bounds(window) @@ -546,12 +678,15 @@ def summary_record( cell_name = required(record, "E-UTRAN FDD小区名称", source_path, row_number) interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number) + longitude, latitude = (cell_coordinates or {}).get(cgi, ("", "")) return { "hour_start": hour_start, "hour_end": hour_end, "cgi": cgi, "cell_name": cell_name, "interference_dbm": interference, + "longitude": longitude, + "latitude": latitude, } diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index 8c83f3c..8fb810c 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -27,6 +27,7 @@ class PipelineTest(unittest.TestCase): for source_type in main.EXPECTED_TYPES: create_archive(root, source_type, COMPLETE_WINDOW) create_archive(root, main.EXPECTED_TYPES[0], INCOMPLETE_WINDOW) + create_cell_data_sources(root) store = RecordingStore() result = main.process(main.LocalSource(root), output, lookback_days=3, store=store) @@ -37,17 +38,26 @@ class PipelineTest(unittest.TestCase): with (result / f"interference_summary_{COMPLETE_WINDOW}.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", "cgi", "cell_name", "interference_dbm"]) + self.assertEqual( + list(rows[0]), + ["hour_start", "hour_end", "cgi", "cell_name", "interference_dbm", "longitude", "latitude"], + ) rows_by_type = {row["cell_name"].removesuffix("-小区"): row 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干扰监控"}: self.assertEqual(rows_by_type[source_type]["cgi"], "460-00-100-1") self.assertTrue(all(row["interference_dbm"] == "-100.5" for row in rows)) + self.assertEqual(rows_by_type["5G干扰监控"]["longitude"], "113.123456") + self.assertEqual(rows_by_type["5G干扰监控"]["latitude"], "22.654321") + self.assertTrue(all(not row["longitude"] and not row["latitude"] for row in rows if row["cgi"] == "460-00-100-1")) manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8")) self.assertEqual(manifest["source_file_count"], 7) + self.assertEqual(manifest["cell_data_file_count"], 5) self.assertEqual(manifest["summary_rows"], 7) + 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.assertNotIn("source_types", manifest) @@ -78,6 +88,30 @@ class PipelineTest(unittest.TestCase): ("2026-07-30 23:00:00", "2026-07-31 00:00:00"), ) + def test_latest_dated_cell_data_file_is_selected(self) -> None: + entries = [ + {"name": "江门5G小区信息表20260727.xlsx", "path": "/old.xlsx", "is_dir": False}, + {"name": "江门5G小区信息表20260728.xlsx", "path": "/latest.xlsx", "is_dir": False}, + {"name": "说明.txt", "path": "/说明.txt", "is_dir": False}, + ] + + selected = main.select_latest_cell_data_file(entries, "/cell-data") + + self.assertEqual(selected["path"], "/latest.xlsx") + + def test_cell_data_schema_change_is_rejected(self) -> None: + raw = create_cell_data_workbook(include_latitude=False) + + with self.assertRaisesRegex(main.ProcessingError, "missing columns.*纬度"): + main.parse_cell_data_workbook(raw, "bad.xlsx") + + def test_conflicting_cell_data_coordinates_are_rejected(self) -> None: + first = create_cell_data_workbook(longitude=113.1) + second = create_cell_data_workbook(longitude=113.2) + + with self.assertRaisesRegex(main.ProcessingError, "Conflicting CellData coordinates"): + main.load_cell_coordinates([("first.xlsx", first), ("second.xlsx", second)]) + def test_mysql_store_replaces_same_hour_and_deletes_other_hours(self) -> None: connection = FakeConnection(datetime(2026, 7, 31, 10), same_hour_rows=7, old_rows=20) store = main.MySQLSummaryStore(connect_factory=lambda: connection) @@ -89,11 +123,15 @@ class PipelineTest(unittest.TestCase): self.assertEqual(result["refreshed_rows"], 7) self.assertEqual(result["old_rows_deleted"], 20) self.assertEqual(connection.cursor_instance.inserted[0][0], datetime(2026, 7, 31, 10)) + self.assertIsNone(connection.cursor_instance.inserted[0][4]) + self.assertIsNone(connection.cursor_instance.inserted[0][5]) self.assertTrue(connection.committed) self.assertFalse(connection.rolled_back) self.assertTrue(connection.closed) self.assertNotIn("source_type", " ".join(connection.cursor_instance.queries)) self.assertNotIn("source_path", " ".join(connection.cursor_instance.queries)) + self.assertTrue(any("ADD COLUMN `longitude`" in query for query in connection.cursor_instance.queries)) + self.assertTrue(any("ADD COLUMN `latitude`" in query for query in connection.cursor_instance.queries)) def test_mysql_store_refuses_to_replace_a_newer_hour(self) -> None: connection = FakeConnection(datetime(2026, 7, 31, 11)) @@ -169,6 +207,9 @@ class FakeCursor: def fetchone(self) -> tuple[datetime | None]: return (self.latest_time,) + def fetchall(self) -> list[tuple[str]]: + return [(name,) for name in ("metric_time", "cgi", "cell_name", "interference_dbm")] + def close(self) -> None: self.closed = True @@ -224,6 +265,8 @@ def database_row() -> dict[str, str]: "cgi": "460-00-200-1", "cell_name": "测试小区", "interference_dbm": "-100.5", + "longitude": "", + "latitude": "", } @@ -273,5 +316,31 @@ def mock_row(source_type: str, window: str) -> list[object]: return [values[column] for column in header] +def create_cell_data_sources(root: Path) -> None: + raw = create_cell_data_workbook() + for index, configured_directory in enumerate(main.CELL_DATA_DIRECTORIES, start=1): + directory = root / Path(configured_directory.lstrip("/")) + directory.mkdir(parents=True, exist_ok=True) + (directory / f"江门小区信息表{index}-20260728.xlsx").write_bytes(raw) + + +def create_cell_data_workbook(longitude: float = 113.123456, include_latitude: bool = True) -> bytes: + workbook = Workbook() + sheet = workbook.active + sheet.title = "小区信息表" + header = ["小区名称", "eNB/gNB", "CI", "经度"] + if include_latitude: + header.append("纬度") + sheet.append(header) + row: list[object] = ["测试小区", 200, 1, longitude] + if include_latitude: + row.append(22.654321) + sheet.append(row) + content = io.BytesIO() + workbook.save(content) + workbook.close() + return content.getvalue() + + if __name__ == "__main__": unittest.main()