feat: 自适应等待干扰数据并归档结果
This commit is contained in:
@@ -9,7 +9,7 @@ import json
|
||||
import math
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
import posixpath
|
||||
import re
|
||||
@@ -75,12 +75,13 @@ CELL_DATA_DIRECTORIES = (
|
||||
"/网优日常优化数据文档/日常性能报表/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)
|
||||
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}/干扰历史数据"
|
||||
|
||||
HEADER_FDD = (
|
||||
"开始时间",
|
||||
@@ -173,8 +174,7 @@ EXPECTED_HEADERS = {
|
||||
}
|
||||
|
||||
SUMMARY_HEADER = (
|
||||
"hour_start",
|
||||
"hour_end",
|
||||
"metric_time",
|
||||
"network_type",
|
||||
"cgi",
|
||||
"cell_name",
|
||||
@@ -206,8 +206,18 @@ class Source(Protocol):
|
||||
|
||||
|
||||
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:
|
||||
@@ -224,7 +234,37 @@ class MetrixApiClient:
|
||||
|
||||
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, timeout=timeout))
|
||||
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,
|
||||
@@ -232,14 +272,15 @@ class MetrixApiClient:
|
||||
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 body is not None:
|
||||
headers["Content-Type"] = "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):
|
||||
@@ -273,62 +314,24 @@ class ApiSummaryStore:
|
||||
self.connection_id = connection_id
|
||||
self.table = DATABASE_TABLE
|
||||
|
||||
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["hour_start"], "%Y-%m-%d %H:%M:%S") for row in rows}
|
||||
if len(metric_times) != 1:
|
||||
raise ProcessingError("Database batch must contain exactly one metric hour")
|
||||
metric_time = next(iter(metric_times))
|
||||
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,
|
||||
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,
|
||||
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 ''",
|
||||
"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",
|
||||
}
|
||||
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,
|
||||
)
|
||||
|
||||
latest_payload = self._query(
|
||||
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,
|
||||
)
|
||||
latest_rows = latest_payload.get("rows", [])
|
||||
latest_value = latest_rows[0].get("latest_time") if latest_rows else None
|
||||
latest_time = datetime.fromisoformat(str(latest_value)) if latest_value else None
|
||||
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}; "
|
||||
@@ -374,13 +377,80 @@ class ApiSummaryStore:
|
||||
"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, "
|
||||
f"longitude, latitude, azimuth, nearby_count 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,
|
||||
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,
|
||||
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 ''",
|
||||
"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",
|
||||
}
|
||||
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 = "") -> dict[str, object]:
|
||||
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": 1, "page_size": 100},
|
||||
{"sql": sql, "database": database, "page": page, "page_size": page_size},
|
||||
timeout=120,
|
||||
)
|
||||
if not isinstance(payload, dict):
|
||||
@@ -404,6 +474,59 @@ class ApiSummaryStore:
|
||||
) + ")"
|
||||
|
||||
|
||||
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
|
||||
@@ -486,10 +609,16 @@ class LocalSource:
|
||||
|
||||
def parse_candidate(path: str, size: int = 0) -> Candidate | None:
|
||||
name = posixpath.basename(path.replace("\\", "/"))
|
||||
match = FILE_RE.fullmatch(name)
|
||||
if not match or match.group("source_type") not in EXPECTED_TYPES:
|
||||
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
|
||||
return Candidate(match.group("source_type"), match.group("window"), path, size)
|
||||
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]:
|
||||
@@ -568,33 +697,38 @@ def load_cell_metadata(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[st
|
||||
return metadata
|
||||
|
||||
|
||||
def select_window(candidates: list[Candidate], requested: str = "") -> tuple[str, dict[str, Candidate], list[str]]:
|
||||
grouped: dict[str, dict[str, Candidate]] = {}
|
||||
duplicates: list[str] = []
|
||||
for candidate in candidates:
|
||||
window_group = grouped.setdefault(candidate.window, {})
|
||||
if candidate.source_type in window_group:
|
||||
duplicates.append(f"{candidate.window}/{candidate.source_type}")
|
||||
window_group[candidate.source_type] = candidate
|
||||
if duplicates:
|
||||
raise ProcessingError(f"Duplicate source files: {', '.join(sorted(duplicates))}")
|
||||
complete = sorted(window for window, items in grouped.items() if all(name in items for name in EXPECTED_TYPES))
|
||||
if requested:
|
||||
if requested not in complete:
|
||||
present = sorted(grouped.get(requested, {}))
|
||||
missing = [name for name in EXPECTED_TYPES if name not in present]
|
||||
raise ProcessingError(f"Requested window is incomplete: {requested}; missing={missing}")
|
||||
selected = requested
|
||||
elif complete:
|
||||
selected = complete[-1]
|
||||
else:
|
||||
raise ProcessingError("No hour contains all seven interference source types")
|
||||
warnings_out: list[str] = []
|
||||
latest_seen = max(grouped) if grouped else ""
|
||||
if latest_seen and latest_seen != selected:
|
||||
missing = [name for name in EXPECTED_TYPES if name not in grouped[latest_seen]]
|
||||
warnings_out.append(f"Latest observed window {latest_seen} is incomplete; using {selected}; missing={missing}")
|
||||
return selected, grouped[selected], warnings_out
|
||||
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]:
|
||||
@@ -637,13 +771,43 @@ def process(
|
||||
lookback_days: int,
|
||||
requested_window: str = "",
|
||||
store: SummaryStore | None = None,
|
||||
) -> Path:
|
||||
history: HistoryStore | None = None,
|
||||
) -> Path | None:
|
||||
candidates = source.candidates(lookback_days)
|
||||
window, selected, warnings_out = select_window(candidates, requested_window)
|
||||
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
|
||||
|
||||
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))
|
||||
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
|
||||
|
||||
cell_data_workbooks = source.cell_data_workbooks()
|
||||
cell_metadata = load_cell_metadata(cell_data_workbooks)
|
||||
temp_dir = output_root.resolve() / f".{window}.tmp-{os.getpid()}"
|
||||
final_dir = output_root.resolve() / window
|
||||
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)
|
||||
@@ -654,16 +818,17 @@ def process(
|
||||
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}_{window}.csv"
|
||||
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, window, row_number, cell_metadata)
|
||||
summary = summary_record(record, source_type, candidate.path, metric_time, row_number, cell_metadata)
|
||||
if passes_interference_threshold(summary):
|
||||
summary_rows.append(summary)
|
||||
else:
|
||||
@@ -672,6 +837,7 @@ def process(
|
||||
{
|
||||
"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}",
|
||||
@@ -680,26 +846,27 @@ def process(
|
||||
|
||||
populate_nearby_counts(summary_rows)
|
||||
summary_rows.sort(key=lambda item: (item["cgi"], item["cell_name"]))
|
||||
summary_name = f"interference_summary_{window}.csv"
|
||||
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_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(),
|
||||
"window": window,
|
||||
"hour_start": window_bounds(window)[0],
|
||||
"hour_end": window_bounds(window)[1],
|
||||
"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,
|
||||
"coordinate_matched_rows": matched_coordinates,
|
||||
"coordinate_unmatched_rows": len(summary_rows) - matched_coordinates,
|
||||
"warnings": warnings_out,
|
||||
"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)
|
||||
@@ -711,7 +878,8 @@ def process(
|
||||
shutil.rmtree(temp_dir, ignore_errors=True)
|
||||
raise
|
||||
|
||||
print(f"selected_window={window}")
|
||||
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)}")
|
||||
@@ -722,8 +890,8 @@ def process(
|
||||
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']}")
|
||||
for message in warnings_out:
|
||||
print(f"warning={message}")
|
||||
if history_result["enabled"]:
|
||||
print(f"history_output={history_result['path']}")
|
||||
print(f"output={final_dir}")
|
||||
return final_dir
|
||||
|
||||
@@ -732,14 +900,21 @@ def summary_record(
|
||||
record: dict[str, object],
|
||||
source_type: str,
|
||||
source_path: str,
|
||||
window: str,
|
||||
metric_time: datetime,
|
||||
row_number: int,
|
||||
cell_metadata: dict[str, tuple[str, str, str]] | None = None,
|
||||
) -> dict[str, str]:
|
||||
hour_start = normalize_cell(record["开始时间"])
|
||||
expected_start, hour_end = window_bounds(window)
|
||||
if hour_start != expected_start:
|
||||
raise ProcessingError(f"{source_path}: row {row_number} time {hour_start!r} does not match {expected_start!r}")
|
||||
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)
|
||||
@@ -760,8 +935,7 @@ def summary_record(
|
||||
|
||||
longitude, latitude, azimuth = (cell_metadata or {}).get(cgi, ("", "", "0"))
|
||||
return {
|
||||
"hour_start": hour_start,
|
||||
"hour_end": hour_end,
|
||||
"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,
|
||||
@@ -850,18 +1024,6 @@ def required(record: dict[str, object], column: str, path: str, row: int) -> str
|
||||
return value
|
||||
|
||||
|
||||
def window_bounds(window: str) -> tuple[str, str]:
|
||||
if not re.fullmatch(r"\d{16}", window):
|
||||
raise ProcessingError(f"Invalid window: {window}")
|
||||
start = datetime.strptime(window[:12], "%Y%m%d%H%M")
|
||||
end_hour = int(window[12:14])
|
||||
end_minute = int(window[14:16])
|
||||
end = start.replace(hour=end_hour, minute=end_minute)
|
||||
if end <= start:
|
||||
end += timedelta(days=1)
|
||||
return start.strftime("%Y-%m-%d %H:%M:%S"), end.strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
|
||||
def normalize_cell(value: object) -> str:
|
||||
if value is None:
|
||||
return ""
|
||||
@@ -899,10 +1061,15 @@ def write_csv(path: Path, header: tuple[str, ...], rows: list[tuple[object, ...]
|
||||
|
||||
|
||||
def write_dict_csv(path: Path, header: tuple[str, ...], rows: list[dict[str, str]]) -> None:
|
||||
with path.open("w", encoding="utf-8-sig", newline="") as file:
|
||||
writer = csv.DictWriter(file, fieldnames=header, extrasaction="raise")
|
||||
writer.writeheader()
|
||||
writer.writerows(rows)
|
||||
path.write_bytes(dict_csv_bytes(header, rows))
|
||||
|
||||
|
||||
def dict_csv_bytes(header: tuple[str, ...], rows: list[dict[str, str]]) -> bytes:
|
||||
output = io.StringIO(newline="")
|
||||
writer = csv.DictWriter(output, fieldnames=header, extrasaction="raise")
|
||||
writer.writeheader()
|
||||
writer.writerows(rows)
|
||||
return output.getvalue().encode("utf-8-sig")
|
||||
|
||||
|
||||
def ensure_scoped(root: Path, target: Path) -> None:
|
||||
@@ -911,10 +1078,14 @@ def ensure_scoped(root: Path, target: Path) -> None:
|
||||
|
||||
|
||||
def build_parser() -> argparse.ArgumentParser:
|
||||
parser = argparse.ArgumentParser(description="Process the latest complete hour of interference KPI files")
|
||||
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 exact 16-digit source window")
|
||||
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))
|
||||
@@ -934,10 +1105,12 @@ def main(argv: list[str] | None = None) -> int:
|
||||
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)
|
||||
process(source, args.output_dir, args.lookback_days, args.window, store=store)
|
||||
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
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user