feat: 增加干扰小区方位角

This commit is contained in:
2026-08-03 15:25:51 +08:00
parent 9776391ad1
commit 60e5c1ce6a
4 changed files with 71 additions and 32 deletions
+4 -4
View File
@@ -16,17 +16,17 @@ output/2026073110001100/
└── manifest.json └── 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. 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.
The summary columns are: The summary columns are:
```text ```text
hour_start,hour_end,cgi,cell_name,interference_dbm,longitude,latitude hour_start,hour_end,cgi,cell_name,interference_dbm,longitude,latitude,azimuth
``` ```
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. By default it writes the complete selected hour through the Metrix Database API 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, 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. The database table contains one time column, `metric_time DATETIME`, which is the source KPI start time, nullable `longitude` and `latitude` columns, and `azimuth DECIMAL(6,2) 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.
CGI is generated with fixed rules: CGI is generated with fixed rules:
@@ -35,7 +35,7 @@ CGI is generated with fixed rules:
The Metrix API address, API Token, and database connection ID are constants in `runtime_config.py`; they are not project environment variables. Copy `runtime_config.example.py` to `runtime_config.py` and fill in the actual Token and the `conn_id` of `ShareMySQL`. The fixed database is `interference_etl`, the fixed table is `interference_hourly_summary`, and the script creates both automatically when they do not exist. The Metrix API address, API Token, and database connection ID are constants in `runtime_config.py`; they are not project environment variables. Copy `runtime_config.example.py` to `runtime_config.py` and fill in the actual Token and the `conn_id` of `ShareMySQL`. The fixed database is `interference_etl`, the fixed table is `interference_hourly_summary`, and the script creates both 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. 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. CellData column `方向角` maps to summary/database column `azimuth`; an empty source value becomes `0`.
## Local development ## Local development
+8
View File
@@ -56,3 +56,11 @@
- Eleven containerized unit tests cover the pipeline, CellData enrichment, API database transaction, newer-hour protection, decimal validation, and API failure propagation. - Eleven containerized unit tests cover the pipeline, CellData enrichment, API database transaction, newer-hour protection, decimal validation, and API failure propagation.
- SSH deployment and a real Metrix runner execution succeeded with run `eead5c315fb04836a6327be201f561e6`. Window `2026080310001100` produced 1,991 rows, matched 1,977 coordinates, left 14 unmatched, and deleted 1,789 rows from the prior hour. - SSH deployment and a real Metrix runner execution succeeded with run `eead5c315fb04836a6327be201f561e6`. Window `2026080310001100` produced 1,991 rows, matched 1,977 coordinates, left 14 unmatched, and deleted 1,789 rows from the prior hour.
- Final API verification found exactly one stored hour (`2026-08-03 10:00:00`) and exactly six columns: `metric_time`, `cgi`, `cell_name`, `interference_dbm`, `longitude`, and `latitude`. The obsolete online source backup and environment-level API URL/Token entries were removed after the successful run. - Final API verification found exactly one stored hour (`2026-08-03 10:00:00`) and exactly six columns: `metric_time`, `cgi`, `cell_name`, `interference_dbm`, `longitude`, and `latitude`. The obsolete online source backup and environment-level API URL/Token entries were removed after the successful run.
## 2026-08-03: CellData azimuth enrichment
- All five latest CellData workbooks expose direction angle through column `方向角`. CellData metadata now maps each `460-00-{eNB/gNB}-{CI}` to longitude, latitude, and azimuth.
- Summary CSV and `interference_hourly_summary` add `azimuth`. The database type is `DECIMAL(6,2) NOT NULL DEFAULT 0`; existing tables add the column automatically. Empty CellData direction angles and unmatched CGI values both become `0`, while unmatched coordinates remain empty/NULL.
- Real CellData validation loaded 66,524 CGI metadata rows: 54,129 non-zero azimuth values and 12,395 zero/default values. Full window `2026080313001400` produced 1,863 summary rows, including 1,677 non-zero azimuth values and 186 zero values.
- SSH deployment run `e5462179fa8c4f9f9c2dfb2e4541e048` succeeded. Database API verification found exactly one hour (`2026-08-03 13:00:00`), 1,863 rows, and columns `metric_time,cgi,cell_name,interference_dbm,longitude,latitude,azimuth`.
- Twelve containerized unit tests pass, including empty azimuth defaulting and CellData metadata conflict detection.
+38 -22
View File
@@ -160,6 +160,7 @@ SUMMARY_HEADER = (
"interference_dbm", "interference_dbm",
"longitude", "longitude",
"latitude", "latitude",
"azimuth",
) )
class ProcessingError(RuntimeError): class ProcessingError(RuntimeError):
@@ -270,6 +271,7 @@ class ApiSummaryStore:
interference_dbm DECIMAL(10,3) NOT NULL, interference_dbm DECIMAL(10,3) NOT NULL,
longitude DECIMAL(10,6) NULL, longitude DECIMAL(10,6) NULL,
latitude DECIMAL(10,6) NULL, latitude DECIMAL(10,6) NULL,
azimuth DECIMAL(6,2) NOT NULL DEFAULT 0,
PRIMARY KEY (metric_time, cgi) PRIMARY KEY (metric_time, cgi)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
""", """,
@@ -282,10 +284,15 @@ class ApiSummaryStore:
if not isinstance(columns, list): if not isinstance(columns, list):
raise ProcessingError("Metrix Database API returned invalid column metadata") raise ProcessingError("Metrix Database API returned invalid column metadata")
existing_columns = {str(item.get("name")) for item in columns if isinstance(item, dict)} existing_columns = {str(item.get("name")) for item in columns if isinstance(item, dict)}
for column in ("longitude", "latitude"): migrations = {
"longitude": "DECIMAL(10,6) NULL",
"latitude": "DECIMAL(10,6) NULL",
"azimuth": "DECIMAL(6,2) NOT NULL DEFAULT 0",
}
for column, definition in migrations.items():
if column not in existing_columns: if column not in existing_columns:
self._query( self._query(
f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` DECIMAL(10,6) NULL", f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` {definition}",
database=DATABASE_NAME, database=DATABASE_NAME,
) )
@@ -311,7 +318,7 @@ class ApiSummaryStore:
START TRANSACTION; START TRANSACTION;
DELETE FROM `{self.table}` WHERE metric_time = {metric_literal}; DELETE FROM `{self.table}` WHERE metric_time = {metric_literal};
INSERT INTO `{self.table}` INSERT INTO `{self.table}`
(metric_time, cgi, cell_name, interference_dbm, longitude, latitude) (metric_time, cgi, cell_name, interference_dbm, longitude, latitude, azimuth)
VALUES VALUES
{values}; {values};
DELETE FROM `{self.table}` WHERE metric_time <> {metric_literal}; DELETE FROM `{self.table}` WHERE metric_time <> {metric_literal};
@@ -364,6 +371,7 @@ class ApiSummaryStore:
sql_decimal_literal(row["interference_dbm"], "interference_dbm"), sql_decimal_literal(row["interference_dbm"], "interference_dbm"),
sql_decimal_literal(row["longitude"], "longitude", nullable=True), sql_decimal_literal(row["longitude"], "longitude", nullable=True),
sql_decimal_literal(row["latitude"], "latitude", nullable=True), sql_decimal_literal(row["latitude"], "latitude", nullable=True),
sql_decimal_literal(row["azimuth"], "azimuth"),
) )
) + ")" ) + ")"
@@ -475,7 +483,7 @@ def select_latest_cell_data_file(entries: list[dict[str, object]], directory: st
return max(dated_files, key=lambda entry: (entry[0], entry[1]))[2] 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]]: def parse_cell_data_workbook(raw_xlsx: bytes, path: str) -> dict[str, tuple[str, str, str]]:
try: try:
with warnings.catch_warnings(): with warnings.catch_warnings():
warnings.filterwarnings("ignore", message="Workbook contains no default style") warnings.filterwarnings("ignore", message="Workbook contains no default style")
@@ -490,39 +498,46 @@ def parse_cell_data_workbook(raw_xlsx: bytes, path: str) -> dict[str, tuple[str,
header = tuple(normalize_cell(value) for value in next(rows)) header = tuple(normalize_cell(value) for value in next(rows))
except StopIteration as exc: except StopIteration as exc:
raise ProcessingError(f"CellData workbook is empty: {path}") from exc raise ProcessingError(f"CellData workbook is empty: {path}") from exc
required_columns = ("eNB/gNB", "CI", "经度", "纬度") required_columns = ("eNB/gNB", "CI", "经度", "纬度", "方向角")
missing = [column for column in required_columns if column not in header] missing = [column for column in required_columns if column not in header]
if missing: if missing:
raise ProcessingError(f"CellData workbook missing columns {missing}: {path}") raise ProcessingError(f"CellData workbook missing columns {missing}: {path}")
indexes = {column: header.index(column) for column in required_columns} indexes = {column: header.index(column) for column in required_columns}
coordinates: dict[str, tuple[str, str]] = {} metadata: dict[str, tuple[str, str, str]] = {}
for row_number, row in enumerate(rows, start=2): for row_number, row in enumerate(rows, start=2):
node = normalize_cell(row[indexes["eNB/gNB"]]) node = normalize_cell(row[indexes["eNB/gNB"]])
cell = normalize_cell(row[indexes["CI"]]) cell = normalize_cell(row[indexes["CI"]])
longitude = normalize_cell(row[indexes["经度"]]) longitude = normalize_cell(row[indexes["经度"]])
latitude = normalize_cell(row[indexes["纬度"]]) latitude = normalize_cell(row[indexes["纬度"]])
if not node or not cell or not longitude or not latitude: azimuth = normalize_cell(row[indexes["方向角"]]) or "0"
if not node or not cell:
continue continue
try:
azimuth_number = Decimal(azimuth)
except InvalidOperation as exc:
raise ProcessingError(f"Invalid CellData azimuth {azimuth!r} in {path}, row {row_number}") from exc
if not azimuth_number.is_finite():
raise ProcessingError(f"Invalid CellData azimuth {azimuth!r} in {path}, row {row_number}")
cgi = f"{LTE_PLMN}-{node}-{cell}" cgi = f"{LTE_PLMN}-{node}-{cell}"
value = (longitude, latitude) value = (longitude, latitude, format(azimuth_number, "f"))
existing = coordinates.get(cgi) existing = metadata.get(cgi)
if existing is not None and existing != value: if existing is not None and existing != value:
raise ProcessingError(f"Conflicting CellData coordinates for {cgi} in {path}, row {row_number}") raise ProcessingError(f"Conflicting CellData metadata for {cgi} in {path}, row {row_number}")
coordinates[cgi] = value metadata[cgi] = value
return coordinates return metadata
finally: finally:
workbook.close() workbook.close()
def load_cell_coordinates(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[str, str]]: def load_cell_metadata(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[str, str, str]]:
coordinates: dict[str, tuple[str, str]] = {} metadata: dict[str, tuple[str, str, str]] = {}
for path, raw_xlsx in workbooks: for path, raw_xlsx in workbooks:
for cgi, value in parse_cell_data_workbook(raw_xlsx, path).items(): for cgi, value in parse_cell_data_workbook(raw_xlsx, path).items():
existing = coordinates.get(cgi) existing = metadata.get(cgi)
if existing is not None and existing != value: if existing is not None and existing != value:
raise ProcessingError(f"Conflicting CellData coordinates for {cgi} across workbooks") raise ProcessingError(f"Conflicting CellData metadata for {cgi} across workbooks")
coordinates[cgi] = value metadata[cgi] = value
return coordinates return metadata
def select_window(candidates: list[Candidate], requested: str = "") -> tuple[str, dict[str, Candidate], list[str]]: def select_window(candidates: list[Candidate], requested: str = "") -> tuple[str, dict[str, Candidate], list[str]]:
@@ -598,7 +613,7 @@ def process(
candidates = source.candidates(lookback_days) candidates = source.candidates(lookback_days)
window, selected, warnings_out = select_window(candidates, requested_window) window, selected, warnings_out = select_window(candidates, requested_window)
cell_data_workbooks = source.cell_data_workbooks() cell_data_workbooks = source.cell_data_workbooks()
cell_coordinates = load_cell_coordinates(cell_data_workbooks) cell_metadata = load_cell_metadata(cell_data_workbooks)
temp_dir = output_root.resolve() / f".{window}.tmp-{os.getpid()}" temp_dir = output_root.resolve() / f".{window}.tmp-{os.getpid()}"
final_dir = output_root.resolve() / window final_dir = output_root.resolve() / window
ensure_scoped(output_root.resolve(), temp_dir) ensure_scoped(output_root.resolve(), temp_dir)
@@ -620,7 +635,7 @@ def process(
for row_number, row in enumerate(rows, start=2): for row_number, row in enumerate(rows, start=2):
record = dict(zip(header, row, strict=True)) record = dict(zip(header, row, strict=True))
summary_rows.append( summary_rows.append(
summary_record(record, source_type, candidate.path, window, row_number, cell_coordinates) summary_record(record, source_type, candidate.path, window, row_number, cell_metadata)
) )
manifest_files.append( manifest_files.append(
{ {
@@ -685,7 +700,7 @@ def summary_record(
source_path: str, source_path: str,
window: str, window: str,
row_number: int, row_number: int,
cell_coordinates: dict[str, tuple[str, str]] | None = None, cell_metadata: dict[str, tuple[str, str, str]] | None = None,
) -> dict[str, str]: ) -> dict[str, str]:
hour_start = normalize_cell(record["开始时间"]) hour_start = normalize_cell(record["开始时间"])
expected_start, hour_end = window_bounds(window) expected_start, hour_end = window_bounds(window)
@@ -709,7 +724,7 @@ def summary_record(
cell_name = required(record, "E-UTRAN FDD小区名称", source_path, row_number) cell_name = required(record, "E-UTRAN FDD小区名称", source_path, row_number)
interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number) interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number)
longitude, latitude = (cell_coordinates or {}).get(cgi, ("", "")) longitude, latitude, azimuth = (cell_metadata or {}).get(cgi, ("", "", "0"))
return { return {
"hour_start": hour_start, "hour_start": hour_start,
"hour_end": hour_end, "hour_end": hour_end,
@@ -718,6 +733,7 @@ def summary_record(
"interference_dbm": interference, "interference_dbm": interference,
"longitude": longitude, "longitude": longitude,
"latitude": latitude, "latitude": latitude,
"azimuth": azimuth,
} }
+21 -6
View File
@@ -38,7 +38,7 @@ class PipelineTest(unittest.TestCase):
self.assertEqual(len(rows), 7) self.assertEqual(len(rows), 7)
self.assertEqual( self.assertEqual(
list(rows[0]), list(rows[0]),
["hour_start", "hour_end", "cgi", "cell_name", "interference_dbm", "longitude", "latitude"], ["hour_start", "hour_end", "cgi", "cell_name", "interference_dbm", "longitude", "latitude", "azimuth"],
) )
rows_by_type = {row["cell_name"].removesuffix("-小区"): row for row in rows} 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["5G干扰监控"]["cgi"], "460-00-200-1")
@@ -48,7 +48,9 @@ class PipelineTest(unittest.TestCase):
self.assertTrue(all(row["interference_dbm"] == "-100.5" for row in rows)) 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干扰监控"]["longitude"], "113.123456")
self.assertEqual(rows_by_type["5G干扰监控"]["latitude"], "22.654321") self.assertEqual(rows_by_type["5G干扰监控"]["latitude"], "22.654321")
self.assertEqual(rows_by_type["5G干扰监控"]["azimuth"], "30")
self.assertTrue(all(not row["longitude"] and not row["latitude"] for row in rows if row["cgi"] == "460-00-100-1")) self.assertTrue(all(not row["longitude"] and not row["latitude"] for row in rows if row["cgi"] == "460-00-100-1"))
self.assertTrue(all(row["azimuth"] == "0" for row in rows if row["cgi"] == "460-00-100-1"))
manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8")) manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8"))
self.assertEqual(manifest["source_file_count"], 7) self.assertEqual(manifest["source_file_count"], 7)
@@ -103,12 +105,17 @@ class PipelineTest(unittest.TestCase):
with self.assertRaisesRegex(main.ProcessingError, "missing columns.*纬度"): with self.assertRaisesRegex(main.ProcessingError, "missing columns.*纬度"):
main.parse_cell_data_workbook(raw, "bad.xlsx") main.parse_cell_data_workbook(raw, "bad.xlsx")
def test_conflicting_cell_data_coordinates_are_rejected(self) -> None: def test_conflicting_cell_data_metadata_are_rejected(self) -> None:
first = create_cell_data_workbook(longitude=113.1) first = create_cell_data_workbook(longitude=113.1)
second = create_cell_data_workbook(longitude=113.2) second = create_cell_data_workbook(longitude=113.2)
with self.assertRaisesRegex(main.ProcessingError, "Conflicting CellData coordinates"): with self.assertRaisesRegex(main.ProcessingError, "Conflicting CellData metadata"):
main.load_cell_coordinates([("first.xlsx", first), ("second.xlsx", second)]) main.load_cell_metadata([("first.xlsx", first), ("second.xlsx", second)])
def test_empty_cell_data_azimuth_defaults_to_zero(self) -> None:
metadata = main.parse_cell_data_workbook(create_cell_data_workbook(azimuth=None), "cell-data.xlsx")
self.assertEqual(metadata["460-00-200-1"], ("113.123456", "22.654321", "0"))
def test_api_store_replaces_same_hour_and_deletes_other_hours(self) -> None: def test_api_store_replaces_same_hour_and_deletes_other_hours(self) -> None:
client = FakeApiClient("2026-07-31T10:00:00") client = FakeApiClient("2026-07-31T10:00:00")
@@ -124,10 +131,11 @@ class PipelineTest(unittest.TestCase):
self.assertTrue(any("CREATE DATABASE IF NOT EXISTS `interference_etl`" in query for query in queries)) self.assertTrue(any("CREATE DATABASE IF NOT EXISTS `interference_etl`" in query for query in queries))
self.assertTrue(any("ADD COLUMN `longitude`" in query for query in queries)) self.assertTrue(any("ADD COLUMN `longitude`" in query for query in queries))
self.assertTrue(any("ADD COLUMN `latitude`" in query for query in queries)) self.assertTrue(any("ADD COLUMN `latitude`" in query for query in queries))
self.assertTrue(any("ADD COLUMN `azimuth`" in query for query in queries))
script = next(payload for endpoint, payload in client.posts if endpoint.endswith("/run-script")) script = next(payload for endpoint, payload in client.posts if endpoint.endswith("/run-script"))
self.assertTrue(script["single_session"]) self.assertTrue(script["single_session"])
self.assertIn("START TRANSACTION", script["content"]) self.assertIn("START TRANSACTION", script["content"])
self.assertIn("NULL, NULL)", script["content"]) self.assertIn("NULL, NULL, 0", script["content"])
self.assertIn("CONVERT(0x", script["content"]) self.assertIn("CONVERT(0x", script["content"])
self.assertNotIn("source_type", script["content"]) self.assertNotIn("source_type", script["content"])
self.assertNotIn("source_path", script["content"]) self.assertNotIn("source_path", script["content"])
@@ -226,6 +234,7 @@ def database_row() -> dict[str, str]:
"interference_dbm": "-100.5", "interference_dbm": "-100.5",
"longitude": "", "longitude": "",
"latitude": "", "latitude": "",
"azimuth": "0",
} }
@@ -283,17 +292,23 @@ def create_cell_data_sources(root: Path) -> None:
(directory / f"江门小区信息表{index}-20260728.xlsx").write_bytes(raw) (directory / f"江门小区信息表{index}-20260728.xlsx").write_bytes(raw)
def create_cell_data_workbook(longitude: float = 113.123456, include_latitude: bool = True) -> bytes: def create_cell_data_workbook(
longitude: float = 113.123456,
include_latitude: bool = True,
azimuth: float | None = 30,
) -> bytes:
workbook = Workbook() workbook = Workbook()
sheet = workbook.active sheet = workbook.active
sheet.title = "小区信息表" sheet.title = "小区信息表"
header = ["小区名称", "eNB/gNB", "CI", "经度"] header = ["小区名称", "eNB/gNB", "CI", "经度"]
if include_latitude: if include_latitude:
header.append("纬度") header.append("纬度")
header.append("方向角")
sheet.append(header) sheet.append(header)
row: list[object] = ["测试小区", 200, 1, longitude] row: list[object] = ["测试小区", 200, 1, longitude]
if include_latitude: if include_latitude:
row.append(22.654321) row.append(22.654321)
row.append(azimuth)
sheet.append(row) sheet.append(row)
content = io.BytesIO() content = io.BytesIO()
workbook.save(content) workbook.save(content)