1195 lines
49 KiB
Python
1195 lines
49 KiB
Python
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
|
||
|
||
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",
|
||
"prev_nearby_count",
|
||
)
|
||
|
||
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, prev_nearby_count)
|
||
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 = (
|
||
"SELECT metric_time, network_type, cgi, cell_name, interference_dbm, "
|
||
"prev_interference_dbm, longitude, latitude, azimuth, nearby_count, prev_nearby_count "
|
||
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,
|
||
prev_nearby_count 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",
|
||
"prev_nearby_count": "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"])),
|
||
sql_decimal_literal(row["prev_nearby_count"], "prev_nearby_count", 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"))
|
||
return {
|
||
"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": "",
|
||
}
|
||
|
||
|
||
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:
|
||
counts = [0] * len(rows)
|
||
spatial_index: dict[tuple[int, int, int], list[tuple[int, float, float]]] = {}
|
||
for index, row in enumerate(rows):
|
||
row["nearby_count"] = "0"
|
||
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']}")
|
||
|
||
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:
|
||
counts[index] += 1
|
||
counts[other_index] += 1
|
||
spatial_index.setdefault(bucket, []).append((index, longitude, latitude))
|
||
|
||
for index, count in enumerate(counts):
|
||
rows[index]["nearby_count"] = str(count)
|
||
|
||
|
||
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}
|
||
for row in rows:
|
||
previous = previous_by_cgi.get(row["cgi"])
|
||
if previous is None:
|
||
row["prev_interference_dbm"] = ""
|
||
row["prev_nearby_count"] = ""
|
||
continue
|
||
row["prev_interference_dbm"] = previous["interference_dbm"]
|
||
row["prev_nearby_count"] = previous["nearby_count"]
|
||
|
||
|
||
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
|