fix: 移除干扰来源字段
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user