feat: 增加制式和高干扰筛选

This commit is contained in:
2026-08-05 16:02:32 +08:00
parent 60e5c1ce6a
commit 6ce3d1748f
4 changed files with 102 additions and 12 deletions
+4 -2
View File
@@ -21,12 +21,14 @@ For every run, the script also reads the latest filename-dated XLSX from each co
The summary columns are:
```text
hour_start,hour_end,cgi,cell_name,interference_dbm,longitude,latitude,azimuth
hour_start,hour_end,network_type,cgi,cell_name,interference_dbm,longitude,latitude,azimuth
```
`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.
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, 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.
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, 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:
+7
View File
@@ -64,3 +64,10 @@
- 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.
## 2026-08-05: Network type and high-interference filtering
- Summary CSV and `interference_hourly_summary` add `network_type`: the two `5G...` source types map to `2.6G`, the two `700M...` source types map to `700M`, and SDR FDD/TDD plus reverse-activated RD map to `4G`.
- Only high-interference rows enter the summary and database. `2.6G` keeps values greater than or equal to `-107 dBm`; `700M` and `4G` keep values greater than or equal to `-110 dBm`. Values strictly below those thresholds are discarded; equal values remain.
- The seven converted source CSV files remain unfiltered source conversions. `manifest.json` and stdout record `threshold_filtered_rows` for the summary filter.
- Real read-only validation of window `2026080514001500` reduced 1,868 source rows to 1,109 summary rows: `2.6G=137`, `700M=108`, `4G=864`, with 759 lower-interference rows removed. Thirteen containerized unit tests pass.
+41 -4
View File
@@ -46,6 +46,22 @@ EXPECTED_TYPES = (
"反开RD干扰监控",
)
NETWORK_TYPE_BY_SOURCE = {
"5G下FDD干扰监控": "2.6G",
"5G干扰监控": "2.6G",
"700M下FDD干扰监控": "700M",
"700M干扰监控": "700M",
"SDR_FDD干扰监控": "4G",
"SDR_TDD干扰监控": "4G",
"反开RD干扰监控": "4G",
}
MIN_INTERFERENCE_BY_NETWORK_TYPE = {
"2.6G": Decimal("-107"),
"700M": Decimal("-110"),
"4G": Decimal("-110"),
}
DEFAULT_STORAGE_ID = "stg_4d9a910d72"
DEFAULT_ROOT = "/网优日常优化数据文档/(勿删)干扰定时小时指标"
CELL_DATA_DIRECTORIES = (
@@ -155,6 +171,7 @@ EXPECTED_HEADERS = {
SUMMARY_HEADER = (
"hour_start",
"hour_end",
"network_type",
"cgi",
"cell_name",
"interference_dbm",
@@ -266,6 +283,7 @@ class ApiSummaryStore:
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,
@@ -285,6 +303,7 @@ class ApiSummaryStore:
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",
@@ -318,7 +337,7 @@ class ApiSummaryStore:
START TRANSACTION;
DELETE FROM `{self.table}` WHERE metric_time = {metric_literal};
INSERT INTO `{self.table}`
(metric_time, cgi, cell_name, interference_dbm, longitude, latitude, azimuth)
(metric_time, network_type, cgi, cell_name, interference_dbm, longitude, latitude, azimuth)
VALUES
{values};
DELETE FROM `{self.table}` WHERE metric_time <> {metric_literal};
@@ -366,6 +385,7 @@ class ApiSummaryStore:
return "(" + ", ".join(
(
metric_literal,
sql_text_literal(row["network_type"]),
sql_text_literal(row["cgi"]),
sql_text_literal(row["cell_name"]),
sql_decimal_literal(row["interference_dbm"], "interference_dbm"),
@@ -623,6 +643,7 @@ def process(
converted_dir.mkdir(parents=True)
summary_rows: list[dict[str, str]] = []
threshold_filtered_rows = 0
manifest_files: list[dict[str, object]] = []
database_result: dict[str, object] = {"enabled": False}
try:
@@ -634,9 +655,11 @@ 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, cell_metadata)
)
summary = summary_record(record, source_type, candidate.path, window, row_number, cell_metadata)
if passes_interference_threshold(summary):
summary_rows.append(summary)
else:
threshold_filtered_rows += 1
manifest_files.append(
{
"source_size": candidate.size,
@@ -661,6 +684,7 @@ def process(
"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,
@@ -682,6 +706,7 @@ def process(
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"threshold_filtered_rows={threshold_filtered_rows}")
print(f"coordinate_matched_rows={matched_coordinates}")
print(f"coordinate_unmatched_rows={len(summary_rows) - matched_coordinates}")
if database_result["enabled"]:
@@ -728,6 +753,7 @@ def summary_record(
return {
"hour_start": hour_start,
"hour_end": hour_end,
"network_type": NETWORK_TYPE_BY_SOURCE[source_type],
"cgi": cgi,
"cell_name": cell_name,
"interference_dbm": interference,
@@ -737,6 +763,17 @@ def summary_record(
}
def passes_interference_threshold(record: dict[str, str]) -> bool:
network_type = record["network_type"]
try:
interference = Decimal(record["interference_dbm"])
except InvalidOperation as exc:
raise ProcessingError(f"Invalid interference value: {record['interference_dbm']}") from exc
if not interference.is_finite():
raise ProcessingError(f"Invalid interference value: {record['interference_dbm']}")
return interference >= MIN_INTERFERENCE_BY_NETWORK_TYPE[network_type]
def build_cgi(prefix: str, record: dict[str, object], columns: tuple[str, ...], path: str, row: int) -> str:
parts = [required(record, column, path, row) for column in columns]
return "-".join(([prefix] if prefix else []) + parts)
+50 -6
View File
@@ -38,9 +38,22 @@ class PipelineTest(unittest.TestCase):
self.assertEqual(len(rows), 7)
self.assertEqual(
list(rows[0]),
["hour_start", "hour_end", "cgi", "cell_name", "interference_dbm", "longitude", "latitude", "azimuth"],
[
"hour_start",
"hour_end",
"network_type",
"cgi",
"cell_name",
"interference_dbm",
"longitude",
"latitude",
"azimuth",
],
)
rows_by_type = {row["cell_name"].removesuffix("-小区"): row for row in rows}
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.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干扰监控"}:
@@ -56,6 +69,7 @@ class PipelineTest(unittest.TestCase):
self.assertEqual(manifest["source_file_count"], 7)
self.assertEqual(manifest["cell_data_file_count"], 5)
self.assertEqual(manifest["summary_rows"], 7)
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)
@@ -65,6 +79,28 @@ class PipelineTest(unittest.TestCase):
self.assertEqual(manifest["database"]["metric_time"], "2026-07-31 10:00:00")
self.assertEqual(len(store.rows), 7)
def test_interference_thresholds_remove_only_lower_values(self) -> None:
for network_type, threshold in (("2.6G", "-107"), ("700M", "-110"), ("4G", "-110")):
row = database_row()
row["network_type"] = network_type
row["interference_dbm"] = threshold
self.assertTrue(main.passes_interference_threshold(row))
row["interference_dbm"] = str(float(threshold) - 0.1)
self.assertFalse(main.passes_interference_threshold(row))
with tempfile.TemporaryDirectory() as temp:
root = Path(temp) / "source"
for source_type in main.EXPECTED_TYPES:
interference = -107.1 if source_type == main.EXPECTED_TYPES[0] else -100.5
create_archive(root, source_type, COMPLETE_WINDOW, interference_dbm=interference)
create_cell_data_sources(root)
result = main.process(main.LocalSource(root), Path(temp) / "output", 3)
manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8"))
self.assertEqual(manifest["summary_rows"], 6)
self.assertEqual(manifest["threshold_filtered_rows"], 1)
def test_requested_incomplete_window_is_rejected(self) -> None:
with tempfile.TemporaryDirectory() as temp:
root = Path(temp) / "source"
@@ -129,6 +165,7 @@ class PipelineTest(unittest.TestCase):
self.assertEqual(result["old_rows_deleted"], 20)
queries = [payload["sql"] for endpoint, payload in client.posts if endpoint.endswith("/query")]
self.assertTrue(any("CREATE DATABASE IF NOT EXISTS `interference_etl`" in query for query in queries))
self.assertTrue(any("ADD COLUMN `network_type`" 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 `azimuth`" in query for query in queries))
@@ -229,6 +266,7 @@ def database_row() -> dict[str, str]:
return {
"hour_start": "2026-07-31 10:00:00",
"hour_end": "2026-07-31 11:00:00",
"network_type": "2.6G",
"cgi": "460-00-200-1",
"cell_name": "测试小区",
"interference_dbm": "-100.5",
@@ -238,7 +276,13 @@ def database_row() -> dict[str, str]:
}
def create_archive(root: Path, source_type: str, window: str, bad_header: bool = False) -> None:
def create_archive(
root: Path,
source_type: str,
window: str,
bad_header: bool = False,
interference_dbm: float = -100.5,
) -> 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}"
@@ -249,7 +293,7 @@ def create_archive(root: Path, source_type: str, window: str, bad_header: bool =
if bad_header:
header[-1] = "unexpected"
sheet.append(header)
sheet.append(mock_row(source_type, window))
sheet.append(mock_row(source_type, window, interference_dbm))
metadata = workbook.create_sheet("指标(计数器)")
metadata.append(["指标或计数器", "指标或计数器描述", "指标公式", "指标或计数器状态"])
content = io.BytesIO()
@@ -259,7 +303,7 @@ def create_archive(root: Path, source_type: str, window: str, bad_header: bool =
archive.writestr(f"{filename}.xlsx", content.getvalue())
def mock_row(source_type: str, window: str) -> list[object]:
def mock_row(source_type: str, window: str, interference_dbm: float = -100.5) -> list[object]:
header = main.EXPECTED_HEADERS[source_type]
values: dict[str, object] = {column: "mock" for column in header}
values.update(
@@ -277,8 +321,8 @@ def mock_row(source_type: str, window: str) -> list[object]:
"E-UTRAN TDD小区名称": f"{source_type}-小区",
"CU小区配置名称": f"{source_type}-小区",
"小区名称": f"{source_type}-小区",
"载波平均噪声干扰(dBm)": -100.5,
"小区上行平均干扰电平(dBm)": -100.5,
"载波平均噪声干扰(dBm)": interference_dbm,
"小区上行平均干扰电平(dBm)": interference_dbm,
}
)
return [values[column] for column in header]