feat: 增加小区经纬度匹配

This commit is contained in:
2026-08-03 10:42:45 +08:00
parent 1a6658dd03
commit 470763a818
4 changed files with 222 additions and 7 deletions
+139 -4
View File
@@ -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<source_type>.+)_LWP_每小时_过滤110_(?P<window>\d{16})\.zip$", re.IGNORECASE)
CELL_DATA_FILE_RE = re.compile(r"(?P<date>20\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,
}