From 1a6658dd038cbc67fd3f4aeafd85ee4df68ccca8 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Mon, 3 Aug 2026 10:10:07 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E7=A7=BB=E9=99=A4=E5=B9=B2=E6=89=B0?= =?UTF-8?q?=E6=9D=A5=E6=BA=90=E5=AD=97=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 2 +- aidocs/project_context.md | 3 ++- main.py | 19 ++++--------------- tests/test_pipeline.py | 15 +++++++++------ 4 files changed, 16 insertions(+), 23 deletions(-) diff --git a/README.md b/README.md index 7e35ab5..323b9e0 100644 --- a/README.md +++ b/README.md @@ -19,7 +19,7 @@ output/2026073110001100/ The summary columns are: ```text -hour_start,hour_end,source_type,cgi,cell_name,interference_dbm,source_path +hour_start,hour_end,cgi,cell_name,interference_dbm ``` 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. diff --git a/aidocs/project_context.md b/aidocs/project_context.md index 8444544..fe92d79 100644 --- a/aidocs/project_context.md +++ b/aidocs/project_context.md @@ -4,7 +4,7 @@ - This repository is developed as the `InterferenceETL` submodule under Metrix and is intended to run later in Metrix Script Management on an hourly schedule. - `main.py` reads the configured Metrix SFTP storage through Metrix API Token authentication or a local directory for tests. It selects the newest hour containing all seven known interference source types and falls back from a newer incomplete hour. -- Each selected ZIP must contain exactly one XLSX. The script strictly validates the known `Sheet0` header, writes one full UTF-8-BOM CSV per source type, and writes one merged CSV containing `hour_start`, `hour_end`, `source_type`, `cgi`, `cell_name`, `interference_dbm`, and `source_path`. +- Each selected ZIP must contain exactly one XLSX. The script strictly validates the known `Sheet0` header, writes one full UTF-8-BOM CSV per source type, and writes one merged CSV containing only `hour_start`, `hour_end`, `cgi`, `cell_name`, and `interference_dbm`. - NR CGI uses `{gNBplmn}-{gNBId}-{cellId}`. All 4G sources use `460-00-{node}-{cell}`, mapped to the actual node and cell column names in each workbook schema. - Runs are idempotent at the output-window directory level. Generation happens in a scoped temporary directory, then replaces only the same window below the configured output root. Source storage is never modified. - `manifest.json` records input paths, sizes, SHA-256 hashes, row counts, warnings, generated files, and the database write result for ingestion auditing. @@ -38,3 +38,4 @@ - InterferenceETL no longer writes into the Metrix platform database. Its fixed database is `interference_etl`, while the table remains `interference_hourly_summary`. - 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. diff --git a/main.py b/main.py index e7906eb..51abc20 100644 --- a/main.py +++ b/main.py @@ -143,11 +143,9 @@ EXPECTED_HEADERS = { SUMMARY_HEADER = ( "hour_start", "hour_end", - "source_type", "cgi", "cell_name", "interference_dbm", - "source_path", ) class ProcessingError(RuntimeError): @@ -193,12 +191,10 @@ class MySQLSummaryStore: f""" CREATE TABLE IF NOT EXISTS {table} ( metric_time DATETIME NOT NULL COMMENT '指标开始时间', - source_type VARCHAR(64) NOT NULL, cgi VARCHAR(128) NOT NULL, cell_name VARCHAR(255) NOT NULL, interference_dbm DECIMAL(10,3) NOT NULL, - source_path VARCHAR(1024) NOT NULL, - PRIMARY KEY (metric_time, source_type, cgi) + PRIMARY KEY (metric_time, cgi) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 """ ) @@ -218,17 +214,15 @@ class MySQLSummaryStore: cursor.executemany( f""" INSERT INTO {table} - (metric_time, source_type, cgi, cell_name, interference_dbm, source_path) - VALUES (%s, %s, %s, %s, %s, %s) + (metric_time, cgi, cell_name, interference_dbm) + VALUES (%s, %s, %s, %s) """, [ ( metric_time, - row["source_type"], row["cgi"], row["cell_name"], row["interference_dbm"], - row["source_path"], ) for row in rows ], @@ -479,8 +473,6 @@ def process( summary_rows.append(summary_record(record, source_type, candidate.path, window, row_number)) manifest_files.append( { - "source_type": source_type, - "source_path": candidate.path, "source_size": candidate.size, "sha256": hashlib.sha256(raw_zip).hexdigest(), "xlsx_member": member_name, @@ -489,7 +481,7 @@ def process( } ) - summary_rows.sort(key=lambda item: (EXPECTED_TYPES.index(item["source_type"]), item["cgi"], item["cell_name"])) + summary_rows.sort(key=lambda item: (item["cgi"], item["cell_name"])) summary_name = f"interference_summary_{window}.csv" write_dict_csv(temp_dir / summary_name, SUMMARY_HEADER, summary_rows) if store is not None: @@ -499,7 +491,6 @@ def process( "window": window, "hour_start": window_bounds(window)[0], "hour_end": window_bounds(window)[1], - "source_types": list(EXPECTED_TYPES), "source_file_count": len(manifest_files), "summary_rows": len(summary_rows), "warnings": warnings_out, @@ -558,11 +549,9 @@ def summary_record( return { "hour_start": hour_start, "hour_end": hour_end, - "source_type": source_type, "cgi": cgi, "cell_name": cell_name, "interference_dbm": interference, - "source_path": source_path, } diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index 191bfd6..8c83f3c 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -11,9 +11,8 @@ import unittest from unittest.mock import patch import zipfile -from openpyxl import Workbook - import main +from openpyxl import Workbook COMPLETE_WINDOW = "2026073110001100" @@ -38,8 +37,8 @@ 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({row["source_type"] for row in rows}, set(main.EXPECTED_TYPES)) - rows_by_type = {row["source_type"]: row for row in rows} + self.assertEqual(list(rows[0]), ["hour_start", "hour_end", "cgi", "cell_name", "interference_dbm"]) + 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干扰监控"}: @@ -51,6 +50,8 @@ class PipelineTest(unittest.TestCase): self.assertEqual(manifest["summary_rows"], 7) self.assertEqual(len(manifest["warnings"]), 1) self.assertIn(INCOMPLETE_WINDOW, manifest["warnings"][0]) + 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.assertEqual(manifest["database"]["metric_time"], "2026-07-31 10:00:00") self.assertEqual(len(store.rows), 7) @@ -91,6 +92,8 @@ class PipelineTest(unittest.TestCase): 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)) def test_mysql_store_refuses_to_replace_a_newer_hour(self) -> None: connection = FakeConnection(datetime(2026, 7, 31, 11)) @@ -147,10 +150,12 @@ class FakeCursor: self.old_rows = old_rows self.rowcount = 0 self.inserted: list[tuple[object, ...]] = [] + self.queries: list[str] = [] self.closed = False def execute(self, query: str, params: tuple[object, ...] | None = None) -> None: normalized = " ".join(query.split()) + self.queries.append(normalized) if normalized.startswith("DELETE") and "metric_time = %s" in normalized: self.rowcount = self.same_hour_rows elif normalized.startswith("DELETE") and "metric_time <> %s" in normalized: @@ -216,11 +221,9 @@ def database_row() -> dict[str, str]: return { "hour_start": "2026-07-31 10:00:00", "hour_end": "2026-07-31 11:00:00", - "source_type": "5G干扰监控", "cgi": "460-00-200-1", "cell_name": "测试小区", "interference_dbm": "-100.5", - "source_path": "/source.zip", }