Files

1236 lines
51 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
import argparse
import csv
from decimal import Decimal, InvalidOperation
import hashlib
import io
import json
import math
import os
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
import posixpath
import re
import shutil
import sys
import time
from typing import Protocol
from urllib.error import HTTPError, URLError
from urllib.parse import quote, urlencode
from urllib.request import ProxyHandler, Request, build_opener
import warnings
import zipfile
VENDOR_DIR = Path(__file__).resolve().parent / "vendor"
sys.path.insert(0, str(VENDOR_DIR))
from openpyxl import load_workbook
try:
from runtime_config import METRIX_API_BASE_URL, METRIX_API_TOKEN, METRIX_DATABASE_CONNECTION_ID
except ImportError:
METRIX_API_BASE_URL = "http://188.5.127.115:18271"
METRIX_API_TOKEN = ""
METRIX_DATABASE_CONNECTION_ID = ""
EXPECTED_TYPES = (
"5G下FDD干扰监控",
"5G干扰监控",
"700M下FDD干扰监控",
"700M干扰监控",
"SDR_FDD干扰监控",
"SDR_TDD干扰监控",
"反开RD干扰监控",
)
NETWORK_TYPE_BY_SOURCE = {
"5G下FDD干扰监控": "FDD",
"5G干扰监控": "2.6G",
"700M下FDD干扰监控": "FDD",
"700M干扰监控": "700M",
"SDR_FDD干扰监控": "FDD",
"SDR_TDD干扰监控": "TDD",
"反开RD干扰监控": "TDD",
}
# Thresholds keep the original band grouping (the two 5G-directory sources use
# -107, all others -110); they are independent of the network_type label.
MIN_INTERFERENCE_BY_SOURCE = {
"5G下FDD干扰监控": Decimal("-107"),
"5G干扰监控": Decimal("-107"),
"700M下FDD干扰监控": Decimal("-110"),
"700M干扰监控": Decimal("-110"),
"SDR_FDD干扰监控": Decimal("-110"),
"SDR_TDD干扰监控": Decimal("-110"),
"反开RD干扰监控": Decimal("-110"),
}
EARTH_RADIUS_KM = 6371.0088
NEARBY_RADIUS_KM = 1.0
NEARBY_COLUMN_BY_NETWORK_TYPE = {
"2.6G": "nearby_26g",
"700M": "nearby_700m",
"TDD": "nearby_tdd",
"FDD": "nearby_fdd",
}
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_TIME_RE = re.compile(r"_(?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_NAME = "interference_etl"
DATABASE_TABLE = "interference_hourly_summary"
DEFAULT_HISTORY_ROOT = f"{DEFAULT_ROOT}/干扰历史数据"
LOCAL_CSV_ENCODING = "utf-8-sig"
HISTORY_CSV_ENCODING = "gbk"
HEADER_FDD = (
"开始时间",
"粒度",
"子网ID",
"子网名称",
"网元ID",
"管理网元",
"eNodeB CUID",
"eNodeB CU名称",
"LTEID",
"LTE名称",
"E-UTRAN FDD小区ID",
"E-UTRAN FDD小区名称",
"cellId",
"eNodeBId",
"载波平均噪声干扰(dBm)",
"集团-上下行总业务量(GB)",
"RRC连接建立最大用户数",
)
HEADER_NR = (
"开始时间",
"粒度",
"子网ID",
"子网名称",
"网元ID",
"管理网元",
"gNB CU-CP功能配置ID",
"gNB CU-CP功能配置名称",
"CU小区配置ID",
"CU小区配置名称",
"cellId",
"duMeMoId",
"gNBId",
"gNBIdLength",
"gNBplmn",
"masterOperatorId",
"nrCarrierGroupId",
"nrPhysicalCellDUId",
"小区上行平均干扰电平(dBm)",
"5G上下行总流量(上行PDCP PDU数据量+下行PDCP成功发送数据量)(GB)",
"RRC连接最大连接用户数",
)
HEADER_SDR = (
"开始时间",
"粒度",
"子网ID",
"子网名称",
"网元ID",
"管理网元",
"eNodeBID",
"eNodeB名称",
"小区ID",
"小区名称",
"载波平均噪声干扰(dBm)",
"集团-上下行总业务量(GB)",
"RRC连接建立最大用户数",
)
HEADER_RD = (
"开始时间",
"粒度",
"子网ID",
"子网名称",
"网元ID",
"管理网元",
"eNodeB CUID",
"eNodeB CU名称",
"LTEID",
"LTE名称",
"E-UTRAN TDD小区ID",
"E-UTRAN TDD小区名称",
"cellId",
"eNodeBId",
"载波平均噪声干扰(dBm)",
"上下行总业务量(GB)",
"RRC连接建立最大用户数",
)
EXPECTED_HEADERS = {
"5G下FDD干扰监控": HEADER_FDD,
"5G干扰监控": HEADER_NR,
"700M下FDD干扰监控": HEADER_FDD,
"700M干扰监控": HEADER_NR,
"SDR_FDD干扰监控": HEADER_SDR,
"SDR_TDD干扰监控": HEADER_SDR,
"反开RD干扰监控": HEADER_RD,
}
SUMMARY_HEADER = (
"metric_time",
"network_type",
"cgi",
"cell_name",
"interference_dbm",
"prev_interference_dbm",
"longitude",
"latitude",
"azimuth",
"nearby_count",
"nearby_26g",
"nearby_700m",
"nearby_tdd",
"nearby_fdd",
"prev_nearby_count",
"prev_nearby_26g",
"prev_nearby_700m",
"prev_nearby_tdd",
"prev_nearby_fdd",
)
class ProcessingError(RuntimeError):
pass
@dataclass(frozen=True)
class Candidate:
source_type: str
window: str
path: str
size: int = 0
class Source(Protocol):
def candidates(self, lookback_days: int) -> list[Candidate]: ...
def download(self, path: str) -> bytes: ...
def cell_data_workbooks(self) -> list[tuple[str, bytes]]: ...
class SummaryStore(Protocol):
def latest_metric_time(self) -> datetime | None: ...
def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: ...
def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: ...
class HistoryStore(Protocol):
def exists(self, metric_time: datetime) -> bool: ...
def upload(self, metric_time: datetime, payload: bytes) -> str: ...
class MetrixApiClient:
def __init__(self, base_url: str, token: str) -> None:
if not token:
raise ProcessingError("METRIX_API_TOKEN is required for Metrix API access")
self.base_url = base_url.rstrip("/")
self.token = token
self.opener = build_opener(ProxyHandler({}))
def get_bytes(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> bytes:
return self._request("GET", endpoint, query=query, timeout=timeout)
def get_json(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> object:
return self._decode_json(self.get_bytes(endpoint, query, timeout))
def post_json(self, endpoint: str, payload: dict[str, object], timeout: int = 30) -> object:
body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
return self._decode_json(
self._request("POST", endpoint, body=body, content_type="application/json", timeout=timeout)
)
def post_file(
self,
endpoint: str,
query: dict[str, object],
filename: str,
payload: bytes,
timeout: int = 120,
) -> object:
boundary = f"----MetrixBoundary{time.time_ns()}"
disposition = f'Content-Disposition: form-data; name="file"; filename="{filename}"\r\n'
body = (
f"--{boundary}\r\n".encode()
+ disposition.encode("utf-8")
+ b"Content-Type: text/csv\r\n\r\n"
+ payload
+ f"\r\n--{boundary}--\r\n".encode()
)
return self._decode_json(
self._request(
"POST",
endpoint,
query=query,
body=body,
content_type=f"multipart/form-data; boundary={boundary}",
timeout=timeout,
)
)
def _request(
self,
method: str,
endpoint: str,
query: dict[str, object] | None = None,
body: bytes | None = None,
content_type: str = "",
timeout: int = 30,
) -> bytes:
url = f"{self.base_url}{endpoint}"
if query:
url = f"{url}?{urlencode(query)}"
headers = {"Authorization": f"Bearer {self.token}", "Accept": "application/json"}
if content_type:
headers["Content-Type"] = content_type
request = Request(url, data=body, headers=headers, method=method)
last_error: Exception | None = None
for attempt in range(3):
try:
with self.opener.open(request, timeout=timeout) as response:
return response.read()
except HTTPError as exc:
detail = exc.read().decode("utf-8", "replace")[:1000]
if exc.code < 500:
raise ProcessingError(f"Metrix API {exc.code}: {detail}") from exc
last_error = exc
except (URLError, TimeoutError) as exc:
last_error = exc
if attempt < 2:
time.sleep(2**attempt)
raise ProcessingError(f"Metrix API request failed: {last_error}")
@staticmethod
def _decode_json(payload: bytes) -> object:
try:
return json.loads(payload.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
raise ProcessingError("Metrix API returned invalid JSON") from exc
class ApiSummaryStore:
def __init__(self, client: MetrixApiClient, connection_id: str) -> None:
if not connection_id:
raise ProcessingError("METRIX_DATABASE_CONNECTION_ID is required for database output")
self.client = client
self.connection_id = connection_id
self.table = DATABASE_TABLE
def latest_metric_time(self) -> datetime | None:
self._ensure_schema()
payload = self._query(
f"SELECT MAX(metric_time) AS latest_time FROM `{self.table}`",
database=DATABASE_NAME,
)
rows = payload.get("rows", [])
value = rows[0].get("latest_time") if rows else None
return datetime.fromisoformat(str(value)) if value else None
def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]:
if not rows:
raise ProcessingError("Refusing to replace database data with an empty batch")
metric_times = {datetime.strptime(row["metric_time"], "%Y-%m-%d %H:%M:%S") for row in rows}
if len(metric_times) != 1:
raise ProcessingError("Database batch must contain exactly one metric time")
metric_time = next(iter(metric_times))
latest_time = self.latest_metric_time()
if latest_time is not None and latest_time > metric_time:
raise ProcessingError(
f"Database already contains newer metric time {latest_time:%Y-%m-%d %H:%M:%S}; "
f"refusing to replace it with {metric_time:%Y-%m-%d %H:%M:%S}"
)
metric_literal = f"'{metric_time:%Y-%m-%d %H:%M:%S}'"
values = ",\n".join(self._row_values(metric_literal, row) for row in rows)
script_payload = self.client.post_json(
self._endpoint("run-script"),
{
"content": f"""
START TRANSACTION;
DELETE FROM `{self.table}` WHERE metric_time = {metric_literal};
INSERT INTO `{self.table}`
(metric_time, network_type, cgi, cell_name, interference_dbm, prev_interference_dbm, longitude, latitude, azimuth, nearby_count, nearby_26g, nearby_700m, nearby_tdd, nearby_fdd, prev_nearby_count, prev_nearby_26g, prev_nearby_700m, prev_nearby_tdd, prev_nearby_fdd)
VALUES
{values};
DELETE FROM `{self.table}` WHERE metric_time <> {metric_literal};
COMMIT;
""",
"database": DATABASE_NAME,
"stop_on_error": True,
"single_session": True,
},
timeout=120,
)
if not isinstance(script_payload, dict) or not isinstance(script_payload.get("results"), list):
raise ProcessingError("Metrix Database API returned an invalid script result")
results = script_payload["results"]
failed = next((item for item in results if not item.get("ok")), None)
if failed is not None:
raise ProcessingError(f"Database script failed at statement {failed.get('index')}: {failed.get('message', '')}")
if len(results) != 5:
raise ProcessingError(f"Database script returned {len(results)} results; expected 5")
return {
"enabled": True,
"table": self.table,
"metric_time": metric_time.strftime("%Y-%m-%d %H:%M:%S"),
"inserted_rows": int(results[2].get("affected_rows") or 0),
"refreshed_rows": int(results[1].get("affected_rows") or 0),
"old_rows_deleted": int(results[3].get("affected_rows") or 0),
}
def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]:
self._ensure_schema()
metric_literal = f"'{metric_time:%Y-%m-%d %H:%M:%S}'"
sql = (
f"SELECT {', '.join(SUMMARY_HEADER)} "
f"FROM `{self.table}` "
f"WHERE metric_time = {metric_literal} ORDER BY cgi, cell_name"
)
result: list[dict[str, str]] = []
page = 1
while True:
payload = self._query(sql, database=DATABASE_NAME, page=page, page_size=1000)
rows = payload.get("rows", [])
if not isinstance(rows, list):
raise ProcessingError("Metrix Database API returned invalid rows")
for row in rows:
normalized = {column: normalize_cell(row.get(column)) for column in SUMMARY_HEADER}
normalized["metric_time"] = datetime.fromisoformat(normalized["metric_time"]).strftime("%Y-%m-%d %H:%M:%S")
result.append(normalized)
total = int(payload.get("total") or len(result))
if len(result) >= total or not rows:
return result
page += 1
def _ensure_schema(self) -> None:
self._query(
f"CREATE DATABASE IF NOT EXISTS `{DATABASE_NAME}` "
"CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci"
)
self._query(
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,
prev_interference_dbm DECIMAL(10,3) NULL,
longitude DECIMAL(10,6) NULL,
latitude DECIMAL(10,6) NULL,
azimuth DECIMAL(6,2) NOT NULL DEFAULT 0,
nearby_count INT NOT NULL DEFAULT 0,
nearby_26g INT NOT NULL DEFAULT 0,
nearby_700m INT NOT NULL DEFAULT 0,
nearby_tdd INT NOT NULL DEFAULT 0,
nearby_fdd INT NOT NULL DEFAULT 0,
prev_nearby_count INT NULL,
prev_nearby_26g INT NULL,
prev_nearby_700m INT NULL,
prev_nearby_tdd INT NULL,
prev_nearby_fdd INT NULL,
PRIMARY KEY (metric_time, cgi)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
""",
database=DATABASE_NAME,
)
columns = self.client.get_json(
self._endpoint("columns"),
{"database": DATABASE_NAME, "table": self.table},
)
if not isinstance(columns, list):
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 ''",
"prev_interference_dbm": "DECIMAL(10,3) NULL",
"longitude": "DECIMAL(10,6) NULL",
"latitude": "DECIMAL(10,6) NULL",
"azimuth": "DECIMAL(6,2) NOT NULL DEFAULT 0",
"nearby_count": "INT NOT NULL DEFAULT 0",
"nearby_26g": "INT NOT NULL DEFAULT 0",
"nearby_700m": "INT NOT NULL DEFAULT 0",
"nearby_tdd": "INT NOT NULL DEFAULT 0",
"nearby_fdd": "INT NOT NULL DEFAULT 0",
"prev_nearby_count": "INT NULL",
"prev_nearby_26g": "INT NULL",
"prev_nearby_700m": "INT NULL",
"prev_nearby_tdd": "INT NULL",
"prev_nearby_fdd": "INT NULL",
}
for column, definition in migrations.items():
if column not in existing_columns:
self._query(
f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` {definition}",
database=DATABASE_NAME,
)
def _endpoint(self, action: str) -> str:
return f"/api/databases/{quote(self.connection_id, safe='')}/{action}"
def _query(self, sql: str, database: str = "", page: int = 1, page_size: int = 100) -> dict[str, object]:
payload = self.client.post_json(
self._endpoint("query"),
{"sql": sql, "database": database, "page": page, "page_size": page_size},
timeout=120,
)
if not isinstance(payload, dict):
raise ProcessingError("Metrix Database API returned an invalid query result")
return payload
@staticmethod
def _row_values(metric_literal: str, row: dict[str, str]) -> str:
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"),
sql_decimal_literal(row["prev_interference_dbm"], "prev_interference_dbm", nullable=True),
sql_decimal_literal(row["longitude"], "longitude", nullable=True),
sql_decimal_literal(row["latitude"], "latitude", nullable=True),
sql_decimal_literal(row["azimuth"], "azimuth"),
str(int(row["nearby_count"])),
str(int(row["nearby_26g"])),
str(int(row["nearby_700m"])),
str(int(row["nearby_tdd"])),
str(int(row["nearby_fdd"])),
sql_decimal_literal(row["prev_nearby_count"], "prev_nearby_count", nullable=True),
sql_decimal_literal(row["prev_nearby_26g"], "prev_nearby_26g", nullable=True),
sql_decimal_literal(row["prev_nearby_700m"], "prev_nearby_700m", nullable=True),
sql_decimal_literal(row["prev_nearby_tdd"], "prev_nearby_tdd", nullable=True),
sql_decimal_literal(row["prev_nearby_fdd"], "prev_nearby_fdd", nullable=True),
)
) + ")"
class ApiHistoryStore:
def __init__(self, client: MetrixApiClient, storage_id: str, root: str = DEFAULT_HISTORY_ROOT) -> None:
self.client = client
self.storage_id = storage_id
self.root = root.rstrip("/")
def exists(self, metric_time: datetime) -> bool:
directory, filename = self.target(metric_time)
base_entries = self._list_dir(self.root)
if not any(item.get("is_dir") and item.get("path") == directory for item in base_entries):
return False
return any(
not item.get("is_dir") and item.get("name") == filename and int(item.get("size") or 0) > 0
for item in self._list_dir(directory)
)
def upload(self, metric_time: datetime, payload: bytes) -> str:
directory, filename = self.target(metric_time)
base_entries = self._list_dir(self.root)
if not any(item.get("is_dir") and item.get("path") == directory for item in base_entries):
response = self.client.post_json(
self._endpoint("mkdir"),
{"path": directory},
timeout=120,
)
if not isinstance(response, dict) or response.get("path") != directory:
raise ProcessingError("Metrix Storage API returned an invalid mkdir result")
response = self.client.post_file(
self._endpoint("upload"),
{"path": directory},
filename,
payload,
)
expected_path = posixpath.join(directory, filename)
if not isinstance(response, dict) or response.get("path") != expected_path:
raise ProcessingError("Metrix Storage API returned an invalid upload result")
return expected_path
def target(self, metric_time: datetime) -> tuple[str, str]:
directory = posixpath.join(self.root, metric_time.strftime("%Y-%m-%d"))
filename = f"干扰数据处理结果_{metric_time:%Y%m%d%H%M%S}.csv"
return directory, filename
def _endpoint(self, action: str) -> str:
return f"/api/storages/{quote(self.storage_id, safe='')}/{action}"
def _list_dir(self, path: str) -> list[dict[str, object]]:
payload = self.client.get_json(self._endpoint("files"), {"path": path, "recursive": "false"})
if not isinstance(payload, dict) or not isinstance(payload.get("entries"), list):
raise ProcessingError("Metrix Storage API returned an invalid file list")
return payload["entries"]
class ApiSource:
def __init__(self, client: MetrixApiClient, storage_id: str, root: str) -> None:
self.client = client
self.storage_id = storage_id
self.root = root.rstrip("/") or "/"
def candidates(self, lookback_days: int) -> list[Candidate]:
root_entries = self._list_dir(self.root)
date_dirs = sorted(
(item for item in root_entries if item.get("is_dir") and re.fullmatch(r"\d{4}-\d{2}-\d{2}", item.get("name", ""))),
key=lambda item: item["name"],
)
selected_dirs = date_dirs[-max(1, lookback_days) :]
result: list[Candidate] = []
for directory in selected_dirs:
for item in self._list_dir(directory["path"]):
if item.get("is_dir"):
continue
candidate = parse_candidate(item["path"], int(item.get("size") or 0))
if candidate is not None:
result.append(candidate)
return result
def download(self, path: str) -> bytes:
endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/download"
return self.client.get_bytes(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.client.get_json(endpoint, {"path": path, "recursive": "false"})
if not isinstance(payload, dict) or not isinstance(payload.get("entries"), list):
raise ProcessingError("Metrix Storage API returned an invalid file list")
return payload["entries"]
class LocalSource:
def __init__(self, root: Path) -> None:
self.root = root.resolve()
if not self.root.is_dir():
raise ProcessingError(f"Local source directory does not exist: {self.root}")
def candidates(self, lookback_days: int) -> list[Candidate]:
del lookback_days
result: list[Candidate] = []
for path in self.root.rglob("*.zip"):
candidate = parse_candidate(str(path.resolve()), path.stat().st_size)
if candidate is not None:
result.append(candidate)
return result
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("\\", "/"))
source_type = next((item for item in EXPECTED_TYPES if name.startswith(f"{item}_")), "")
match = FILE_TIME_RE.search(name)
if not source_type or not match:
return None
window = match.group("window")
try:
window_start(window)
except ProcessingError:
return None
return Candidate(source_type, 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, 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}
metadata: dict[str, tuple[str, 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["纬度"]])
azimuth = normalize_cell(row[indexes["方向角"]]) or "0"
if not node or not cell:
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}"
value = (longitude, latitude, format(azimuth_number, "f"))
existing = metadata.get(cgi)
if existing is not None and existing != value:
raise ProcessingError(f"Conflicting CellData metadata for {cgi} in {path}, row {row_number}")
metadata[cgi] = value
return metadata
finally:
workbook.close()
def load_cell_metadata(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[str, str, str]]:
metadata: dict[str, tuple[str, str, str]] = {}
for path, raw_xlsx in workbooks:
for cgi, value in parse_cell_data_workbook(raw_xlsx, path).items():
existing = metadata.get(cgi)
if existing is not None and existing != value:
raise ProcessingError(f"Conflicting CellData metadata for {cgi} across workbooks")
metadata[cgi] = value
return metadata
def select_window(candidates: list[Candidate], requested: str = "") -> tuple[datetime | None, dict[str, Candidate], list[str]]:
if not candidates:
return None, {}, list(EXPECTED_TYPES)
target_time = metric_time_for_window(requested) if requested else metric_time_for_window(max(candidates, key=candidate_key).window)
grouped = [candidate for candidate in candidates if metric_time_for_window(candidate.window) == target_time]
selected: dict[str, Candidate] = {}
for source_type in EXPECTED_TYPES:
matching = [candidate for candidate in grouped if candidate.source_type == source_type]
if matching:
selected[source_type] = max(matching, key=candidate_key)
missing = [source_type for source_type in EXPECTED_TYPES if source_type not in selected]
return target_time, selected, missing
def candidate_key(candidate: Candidate) -> tuple[datetime, str, str]:
return window_start(candidate.window), candidate.window, candidate.path
def window_start(window: str) -> datetime:
if not re.fullmatch(r"\d{16}", window):
raise ProcessingError(f"Invalid window: {window}")
try:
start = datetime.strptime(window[:12], "%Y%m%d%H%M")
datetime.strptime(window[12:], "%H%M")
except ValueError as exc:
raise ProcessingError(f"Invalid window: {window}") from exc
return start
def metric_time_for_window(window: str) -> datetime:
start = window_start(window)
return start.replace(minute=(start.minute // 15) * 15, second=0, microsecond=0)
def parse_workbook(raw_zip: bytes, source_type: str) -> tuple[tuple[str, ...], list[tuple[object, ...]], str]:
try:
with zipfile.ZipFile(io.BytesIO(raw_zip)) as archive:
bad_member = archive.testzip()
if bad_member:
raise ProcessingError(f"ZIP CRC check failed: {bad_member}")
xlsx_members = [item for item in archive.infolist() if not item.is_dir() and item.filename.lower().endswith(".xlsx")]
if len(xlsx_members) != 1:
raise ProcessingError(f"Expected one XLSX member, found {len(xlsx_members)}")
member = xlsx_members[0]
workbook_bytes = archive.read(member)
except zipfile.BadZipFile as exc:
raise ProcessingError("Invalid ZIP archive") from exc
with warnings.catch_warnings():
warnings.filterwarnings("ignore", message="Workbook contains no default style")
workbook = load_workbook(io.BytesIO(workbook_bytes), read_only=True, data_only=True)
try:
if "Sheet0" not in workbook.sheetnames:
raise ProcessingError("Workbook does not contain Sheet0")
rows = workbook["Sheet0"].iter_rows(values_only=True)
try:
header = tuple(normalize_cell(value) for value in next(rows))
except StopIteration as exc:
raise ProcessingError("Sheet0 is empty") from exc
expected = EXPECTED_HEADERS[source_type]
if header != expected:
raise ProcessingError(f"Unexpected Sheet0 header for {source_type}: {header}")
data = [tuple(row) for row in rows if any(value is not None and normalize_cell(value) != "" for value in row)]
return header, data, member.filename
finally:
workbook.close()
def process(
source: Source,
output_root: Path,
lookback_days: int,
requested_window: str = "",
store: SummaryStore | None = None,
history: HistoryStore | None = None,
) -> Path | None:
candidates = source.candidates(lookback_days)
metric_time, selected, missing = select_window(candidates, requested_window)
if metric_time is None:
print("status=waiting")
print("reason=no_source_files")
return None
metric_time_text = metric_time.strftime("%Y-%m-%d %H:%M:%S")
if missing:
print("status=waiting")
print(f"target_time={metric_time_text}")
print(f"missing_sources={','.join(missing)}")
return None
latest_time: datetime | None = None
previous_rows: list[dict[str, str]] = []
if store is not None:
latest_time = store.latest_metric_time()
if latest_time is not None and latest_time >= metric_time:
if latest_time == metric_time and history is not None and not history.exists(metric_time):
stored_rows = store.rows_for_time(metric_time)
if not stored_rows:
raise ProcessingError(f"Database contains no rows for {metric_time_text}")
history_path = history.upload(
metric_time,
dict_csv_bytes(SUMMARY_HEADER, stored_rows, HISTORY_CSV_ENCODING),
)
print("status=history_repaired")
print(f"target_time={metric_time_text}")
print(f"history_output={history_path}")
else:
print("status=skipped")
print(f"target_time={metric_time_text}")
print(f"database_metric_time={latest_time:%Y-%m-%d %H:%M:%S}")
return None
if latest_time is not None:
previous_rows = store.rows_for_time(latest_time)
cell_data_workbooks = source.cell_data_workbooks()
cell_metadata = load_cell_metadata(cell_data_workbooks)
output_key = metric_time.strftime("%Y%m%d%H%M%S")
temp_dir = output_root.resolve() / f".{output_key}.tmp-{os.getpid()}"
final_dir = output_root.resolve() / output_key
ensure_scoped(output_root.resolve(), temp_dir)
if temp_dir.exists():
shutil.rmtree(temp_dir)
converted_dir = temp_dir / "converted"
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}
history_result: dict[str, object] = {"enabled": False}
try:
for source_type in EXPECTED_TYPES:
candidate = selected[source_type]
raw_zip = source.download(candidate.path)
header, rows, member_name = parse_workbook(raw_zip, source_type)
csv_name = f"{source_type}_{candidate.window}.csv"
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 = summary_record(record, source_type, candidate.path, metric_time, row_number, cell_metadata)
if passes_interference_threshold(summary, source_type):
summary_rows.append(summary)
else:
threshold_filtered_rows += 1
manifest_files.append(
{
"source_size": candidate.size,
"sha256": hashlib.sha256(raw_zip).hexdigest(),
"source_window": candidate.window,
"xlsx_member": member_name,
"rows": len(rows),
"converted_csv": f"converted/{csv_name}",
}
)
duplicate_cgi_rows = deduplicate_by_cgi(summary_rows)
populate_nearby_counts(summary_rows)
populate_previous_values(summary_rows, previous_rows)
summary_rows.sort(key=lambda item: (item["cgi"], item["cell_name"]))
summary_name = f"interference_summary_{output_key}.csv"
write_dict_csv(temp_dir / summary_name, SUMMARY_HEADER, summary_rows)
if store is not None:
database_result = store.replace_latest(summary_rows)
if history is not None:
history_path = history.upload(
metric_time,
dict_csv_bytes(SUMMARY_HEADER, summary_rows, HISTORY_CSV_ENCODING),
)
history_result = {"enabled": True, "path": history_path}
matched_coordinates = sum(bool(row["longitude"] and row["latitude"]) for row in summary_rows)
manifest = {
"generated_at": datetime.now(timezone.utc).isoformat(),
"metric_time": metric_time_text,
"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,
"duplicate_cgi_rows": duplicate_cgi_rows,
"coordinate_matched_rows": matched_coordinates,
"coordinate_unmatched_rows": len(summary_rows) - matched_coordinates,
"files": manifest_files,
"summary_csv": summary_name,
"database": database_result,
"history": history_result,
}
(temp_dir / "manifest.json").write_text(json.dumps(manifest, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
output_root.resolve().mkdir(parents=True, exist_ok=True)
if final_dir.exists():
ensure_scoped(output_root.resolve(), final_dir)
shutil.rmtree(final_dir)
temp_dir.replace(final_dir)
except Exception:
shutil.rmtree(temp_dir, ignore_errors=True)
raise
print("status=processed")
print(f"target_time={metric_time_text}")
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"duplicate_cgi_rows={duplicate_cgi_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']}")
print(f"database_old_rows_deleted={database_result['old_rows_deleted']}")
if history_result["enabled"]:
print(f"history_output={history_result['path']}")
print(f"output={final_dir}")
return final_dir
def summary_record(
record: dict[str, object],
source_type: str,
source_path: str,
metric_time: datetime,
row_number: int,
cell_metadata: dict[str, tuple[str, str, str]] | None = None,
) -> dict[str, str]:
row_time_text = normalize_cell(record["开始时间"])
try:
row_time = datetime.strptime(row_time_text, "%Y-%m-%d %H:%M:%S")
except ValueError as exc:
raise ProcessingError(f"{source_path}: row {row_number} has invalid time {row_time_text!r}") from exc
row_metric_time = row_time.replace(minute=(row_time.minute // 15) * 15, second=0, microsecond=0)
if row_metric_time != metric_time:
raise ProcessingError(
f"{source_path}: row {row_number} time {row_time_text!r} is outside metric group "
f"{metric_time:%Y-%m-%d %H:%M:%S}"
)
if source_type in ("5G干扰监控", "700M干扰监控"):
cgi = build_cgi("", record, ("gNBplmn", "gNBId", "cellId"), source_path, row_number)
cell_name = required(record, "CU小区配置名称", source_path, row_number)
interference = required(record, "小区上行平均干扰电平(dBm)", source_path, row_number)
elif source_type in ("SDR_FDD干扰监控", "SDR_TDD干扰监控"):
cgi = build_cgi(LTE_PLMN, record, ("eNodeBID", "小区ID"), source_path, row_number)
cell_name = required(record, "小区名称", source_path, row_number)
interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number)
elif source_type == "反开RD干扰监控":
cgi = build_cgi(LTE_PLMN, record, ("eNodeBId", "cellId"), source_path, row_number)
cell_name = required(record, "E-UTRAN TDD小区名称", source_path, row_number)
interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number)
else:
cgi = build_cgi(LTE_PLMN, record, ("eNodeBId", "cellId"), source_path, row_number)
cell_name = required(record, "E-UTRAN FDD小区名称", source_path, row_number)
interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number)
longitude, latitude, azimuth = (cell_metadata or {}).get(cgi, ("", "", "0"))
record_out = {
"metric_time": metric_time.strftime("%Y-%m-%d %H:%M:%S"),
"network_type": NETWORK_TYPE_BY_SOURCE[source_type],
"cgi": cgi,
"cell_name": cell_name,
"interference_dbm": interference,
"prev_interference_dbm": "",
"longitude": longitude,
"latitude": latitude,
"azimuth": azimuth,
"nearby_count": "0",
"prev_nearby_count": "",
}
for column in NEARBY_COLUMN_BY_NETWORK_TYPE.values():
record_out[column] = "0"
record_out[f"prev_{column}"] = ""
return record_out
def passes_interference_threshold(record: dict[str, str], source_type: str) -> bool:
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_SOURCE[source_type]
def populate_nearby_counts(rows: list[dict[str, str]]) -> None:
type_counts = [dict.fromkeys(NEARBY_COLUMN_BY_NETWORK_TYPE.values(), 0) for _ in rows]
spatial_index: dict[tuple[int, int, int], list[tuple[int, float, float]]] = {}
for index, row in enumerate(rows):
if not row["longitude"] or not row["latitude"]:
continue
try:
longitude = float(row["longitude"])
latitude = float(row["latitude"])
except ValueError as exc:
raise ProcessingError(f"Invalid coordinates for CGI {row['cgi']}") from exc
if not math.isfinite(longitude) or not math.isfinite(latitude) or not -180 <= longitude <= 180 or not -90 <= latitude <= 90:
raise ProcessingError(f"Invalid coordinates for CGI {row['cgi']}")
column = NEARBY_COLUMN_BY_NETWORK_TYPE[row["network_type"]]
bucket = spatial_bucket(longitude, latitude)
for x_offset in (-1, 0, 1):
for y_offset in (-1, 0, 1):
for z_offset in (-1, 0, 1):
nearby_bucket = (bucket[0] + x_offset, bucket[1] + y_offset, bucket[2] + z_offset)
for other_index, other_longitude, other_latitude in spatial_index.get(nearby_bucket, ()):
if haversine_distance_km(longitude, latitude, other_longitude, other_latitude) <= NEARBY_RADIUS_KM:
type_counts[index][NEARBY_COLUMN_BY_NETWORK_TYPE[rows[other_index]["network_type"]]] += 1
type_counts[other_index][column] += 1
spatial_index.setdefault(bucket, []).append((index, longitude, latitude))
for row, counts in zip(rows, type_counts, strict=True):
for column, count in counts.items():
row[column] = str(count)
row["nearby_count"] = str(sum(counts.values()))
def deduplicate_by_cgi(rows: list[dict[str, str]]) -> int:
selected: dict[str, dict[str, str]] = {}
removed = 0
for row in rows:
previous = selected.get(row["cgi"])
if previous is None:
selected[row["cgi"]] = row
continue
if Decimal(row["interference_dbm"]) > Decimal(previous["interference_dbm"]):
selected[row["cgi"]] = row
removed += 1
rows[:] = selected.values()
return removed
def populate_previous_values(rows: list[dict[str, str]], previous_rows: list[dict[str, str]]) -> None:
previous_by_cgi = {row["cgi"]: row for row in previous_rows}
copied_columns = ("interference_dbm", "nearby_count", *NEARBY_COLUMN_BY_NETWORK_TYPE.values())
for row in rows:
previous = previous_by_cgi.get(row["cgi"])
for column in copied_columns:
row[f"prev_{column}"] = previous[column] if previous is not None else ""
def spatial_bucket(longitude: float, latitude: float) -> tuple[int, int, int]:
longitude = math.radians(longitude)
latitude = math.radians(latitude)
radius_at_latitude = EARTH_RADIUS_KM * math.cos(latitude)
return (
math.floor(radius_at_latitude * math.cos(longitude) / NEARBY_RADIUS_KM),
math.floor(radius_at_latitude * math.sin(longitude) / NEARBY_RADIUS_KM),
math.floor(EARTH_RADIUS_KM * math.sin(latitude) / NEARBY_RADIUS_KM),
)
def haversine_distance_km(longitude1: float, latitude1: float, longitude2: float, latitude2: float) -> float:
longitude1, latitude1, longitude2, latitude2 = map(
math.radians,
(longitude1, latitude1, longitude2, latitude2),
)
longitude_delta = longitude2 - longitude1
latitude_delta = latitude2 - latitude1
haversine = math.sin(latitude_delta / 2) ** 2 + (
math.cos(latitude1) * math.cos(latitude2) * math.sin(longitude_delta / 2) ** 2
)
return 2 * EARTH_RADIUS_KM * math.asin(math.sqrt(min(1.0, haversine)))
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)
def required(record: dict[str, object], column: str, path: str, row: int) -> str:
value = normalize_cell(record[column])
if not value:
raise ProcessingError(f"{path}: row {row} has empty {column}")
return value
def normalize_cell(value: object) -> str:
if value is None:
return ""
if isinstance(value, datetime):
return value.strftime("%Y-%m-%d %H:%M:%S")
if isinstance(value, float) and value.is_integer():
return str(int(value))
return str(value).strip()
def normalize_csv_cell(value: object) -> str:
return normalize_cell(value).replace("\u00a0", " ")
def sql_text_literal(value: str) -> str:
encoded = value.encode("utf-8").hex()
return "''" if not encoded else f"CONVERT(0x{encoded} USING utf8mb4)"
def sql_decimal_literal(value: str, field: str, nullable: bool = False) -> str:
normalized = value.strip()
if not normalized and nullable:
return "NULL"
try:
number = Decimal(normalized)
except InvalidOperation as exc:
raise ProcessingError(f"Invalid decimal value for {field}: {value}") from exc
if not number.is_finite():
raise ProcessingError(f"Invalid decimal value for {field}: {value}")
return format(number, "f")
def write_csv(path: Path, header: tuple[str, ...], rows: list[tuple[object, ...]]) -> None:
with path.open("w", encoding=LOCAL_CSV_ENCODING, newline="") as file:
writer = csv.writer(file)
writer.writerow(header)
for row in rows:
writer.writerow(normalize_cell(value) for value in row)
def write_dict_csv(path: Path, header: tuple[str, ...], rows: list[dict[str, str]]) -> None:
path.write_bytes(dict_csv_bytes(header, rows))
def dict_csv_bytes(
header: tuple[str, ...],
rows: list[dict[str, str]],
encoding: str = LOCAL_CSV_ENCODING,
) -> bytes:
output = io.StringIO(newline="")
writer = csv.DictWriter(output, fieldnames=header, extrasaction="raise")
writer.writeheader()
for row in rows:
if encoding == HISTORY_CSV_ENCODING:
writer.writerow({column: normalize_csv_cell(row[column]) for column in header})
else:
writer.writerow(row)
return output.getvalue().encode(encoding)
def ensure_scoped(root: Path, target: Path) -> None:
if target == root or root not in target.parents:
raise ProcessingError(f"Refusing to modify path outside output root: {target}")
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="Process the latest complete interference KPI time group")
parser.add_argument("--source-dir", type=Path, help="Use a local source tree instead of the Metrix storage API")
parser.add_argument("--output-dir", type=Path, default=Path(os.getenv("INTERFERENCE_OUTPUT_DIR", "output")))
parser.add_argument(
"--window",
default=os.getenv("INTERFERENCE_WINDOW", ""),
help="Optional 16-digit source window used to select a natural 15-minute group",
)
parser.add_argument("--lookback-days", type=int, default=int(os.getenv("INTERFERENCE_LOOKBACK_DAYS", "3")))
parser.add_argument("--storage-id", default=os.getenv("METRIX_STORAGE_ID", DEFAULT_STORAGE_ID))
parser.add_argument("--root", default=os.getenv("INTERFERENCE_SOURCE_ROOT", DEFAULT_ROOT))
parser.add_argument("--no-database", action="store_true", help="Generate files without writing the database")
return parser
def main(argv: list[str] | None = None) -> int:
args = build_parser().parse_args(argv)
if args.lookback_days < 1:
raise ProcessingError("lookback-days must be at least 1")
client: MetrixApiClient | None = None
if args.source_dir:
source: Source = LocalSource(args.source_dir)
else:
client = MetrixApiClient(METRIX_API_BASE_URL, METRIX_API_TOKEN)
source = ApiSource(client, args.storage_id, args.root)
if args.no_database:
store = None
history = None
else:
client = client or MetrixApiClient(METRIX_API_BASE_URL, METRIX_API_TOKEN)
store = ApiSummaryStore(client, METRIX_DATABASE_CONNECTION_ID)
history = ApiHistoryStore(client, args.storage_id, f"{args.root.rstrip('/')}/干扰历史数据")
process(source, args.output_dir, args.lookback_days, args.window, store=store, history=history)
return 0
if __name__ == "__main__":
try:
raise SystemExit(main())
except ProcessingError as exc:
print(f"error={exc}", file=sys.stderr)
raise SystemExit(1) from exc