feat: 改用Metrix API写入干扰数据
This commit is contained in:
@@ -2,6 +2,7 @@ from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import csv
|
||||
from decimal import Decimal, InvalidOperation
|
||||
import hashlib
|
||||
import io
|
||||
import json
|
||||
@@ -14,7 +15,7 @@ import re
|
||||
import shutil
|
||||
import sys
|
||||
import time
|
||||
from typing import Callable, Protocol
|
||||
from typing import Protocol
|
||||
from urllib.error import HTTPError, URLError
|
||||
from urllib.parse import quote, urlencode
|
||||
from urllib.request import Request, urlopen
|
||||
@@ -27,6 +28,13 @@ 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干扰监控",
|
||||
@@ -51,10 +59,6 @@ FILE_RE = re.compile(r"^(?P<source_type>.+)_LWP_每小时_过滤110_(?P<window>\
|
||||
CELL_DATA_FILE_RE = re.compile(r"(?P<date>20\d{6})\.xlsx$", re.IGNORECASE)
|
||||
|
||||
LTE_PLMN = "460-00"
|
||||
DATABASE_HOST = "172.17.0.1"
|
||||
DATABASE_PORT = 3306
|
||||
DATABASE_USER = "root"
|
||||
DATABASE_PASSWORD = "OSp!jmgm@26"
|
||||
DATABASE_NAME = "interference_etl"
|
||||
DATABASE_TABLE = "interference_hourly_summary"
|
||||
|
||||
@@ -182,10 +186,69 @@ class SummaryStore(Protocol):
|
||||
def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: ...
|
||||
|
||||
|
||||
class MySQLSummaryStore:
|
||||
def __init__(self, connect_factory: Callable[[], object] | None = None) -> None:
|
||||
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
|
||||
|
||||
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, timeout=timeout))
|
||||
|
||||
def _request(
|
||||
self,
|
||||
method: str,
|
||||
endpoint: str,
|
||||
query: dict[str, object] | None = None,
|
||||
body: bytes | None = None,
|
||||
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"
|
||||
request = Request(url, data=body, headers=headers, method=method)
|
||||
last_error: Exception | None = None
|
||||
for attempt in range(3):
|
||||
try:
|
||||
with urlopen(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
|
||||
self.connect_factory = connect_factory
|
||||
|
||||
def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]:
|
||||
if not rows:
|
||||
@@ -194,14 +257,13 @@ class MySQLSummaryStore:
|
||||
if len(metric_times) != 1:
|
||||
raise ProcessingError("Database batch must contain exactly one metric hour")
|
||||
metric_time = next(iter(metric_times))
|
||||
table = f"`{self.table}`"
|
||||
connection = self._connect()
|
||||
cursor = None
|
||||
try:
|
||||
cursor = connection.cursor()
|
||||
cursor.execute(
|
||||
f"""
|
||||
CREATE TABLE IF NOT EXISTS {table} (
|
||||
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 '指标开始时间',
|
||||
cgi VARCHAR(128) NOT NULL,
|
||||
cell_name VARCHAR(255) NOT NULL,
|
||||
@@ -210,119 +272,105 @@ class MySQLSummaryStore:
|
||||
latitude DECIMAL(10,6) NULL,
|
||||
PRIMARY KEY (metric_time, cgi)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
|
||||
"""
|
||||
)
|
||||
cursor.execute(f"SHOW COLUMNS FROM {table}")
|
||||
existing_columns = {column[0] for column in cursor.fetchall()}
|
||||
for column in ("longitude", "latitude"):
|
||||
if column not in existing_columns:
|
||||
cursor.execute(f"ALTER TABLE {table} ADD COLUMN `{column}` DECIMAL(10,6) NULL")
|
||||
cursor.execute(f"SELECT MAX(metric_time) FROM {table}")
|
||||
latest_row = cursor.fetchone()
|
||||
latest_time = latest_row[0] if latest_row else None
|
||||
if isinstance(latest_time, str):
|
||||
latest_time = datetime.fromisoformat(latest_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}"
|
||||
""",
|
||||
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)}
|
||||
for column in ("longitude", "latitude"):
|
||||
if column not in existing_columns:
|
||||
self._query(
|
||||
f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` DECIMAL(10,6) NULL",
|
||||
database=DATABASE_NAME,
|
||||
)
|
||||
|
||||
cursor.execute(f"DELETE FROM {table} WHERE metric_time = %s", (metric_time,))
|
||||
refreshed_rows = cursor.rowcount
|
||||
cursor.executemany(
|
||||
f"""
|
||||
INSERT INTO {table}
|
||||
(metric_time, cgi, cell_name, interference_dbm, longitude, latitude)
|
||||
VALUES (%s, %s, %s, %s, %s, %s)
|
||||
""",
|
||||
[
|
||||
(
|
||||
metric_time,
|
||||
row["cgi"],
|
||||
row["cell_name"],
|
||||
row["interference_dbm"],
|
||||
row["longitude"] or None,
|
||||
row["latitude"] or None,
|
||||
)
|
||||
for row in rows
|
||||
],
|
||||
latest_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
|
||||
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}"
|
||||
)
|
||||
cursor.execute(f"DELETE FROM {table} WHERE metric_time <> %s", (metric_time,))
|
||||
old_rows_deleted = cursor.rowcount
|
||||
connection.commit()
|
||||
except ProcessingError:
|
||||
connection.rollback()
|
||||
raise
|
||||
except Exception as exc:
|
||||
connection.rollback()
|
||||
raise ProcessingError(f"MySQL write failed: {exc}") from exc
|
||||
finally:
|
||||
if cursor is not None:
|
||||
cursor.close()
|
||||
connection.close()
|
||||
|
||||
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, cgi, cell_name, interference_dbm, longitude, latitude)
|
||||
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": len(rows),
|
||||
"refreshed_rows": refreshed_rows,
|
||||
"old_rows_deleted": old_rows_deleted,
|
||||
"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 _connect(self) -> object:
|
||||
if self.connect_factory is not None:
|
||||
return self.connect_factory()
|
||||
try:
|
||||
import pymysql
|
||||
except ImportError as exc:
|
||||
raise ProcessingError("PyMySQL is required for database output") from exc
|
||||
try:
|
||||
bootstrap = pymysql.connect(
|
||||
host=DATABASE_HOST,
|
||||
port=DATABASE_PORT,
|
||||
user=DATABASE_USER,
|
||||
password=DATABASE_PASSWORD,
|
||||
charset="utf8mb4",
|
||||
autocommit=True,
|
||||
connect_timeout=15,
|
||||
read_timeout=60,
|
||||
write_timeout=60,
|
||||
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]:
|
||||
payload = self.client.post_json(
|
||||
self._endpoint("query"),
|
||||
{"sql": sql, "database": database, "page": 1, "page_size": 100},
|
||||
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["cgi"]),
|
||||
sql_text_literal(row["cell_name"]),
|
||||
sql_decimal_literal(row["interference_dbm"], "interference_dbm"),
|
||||
sql_decimal_literal(row["longitude"], "longitude", nullable=True),
|
||||
sql_decimal_literal(row["latitude"], "latitude", nullable=True),
|
||||
)
|
||||
try:
|
||||
cursor = bootstrap.cursor()
|
||||
try:
|
||||
cursor.execute(
|
||||
f"CREATE DATABASE IF NOT EXISTS `{DATABASE_NAME}` "
|
||||
"CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci"
|
||||
)
|
||||
finally:
|
||||
cursor.close()
|
||||
finally:
|
||||
bootstrap.close()
|
||||
return pymysql.connect(
|
||||
host=DATABASE_HOST,
|
||||
port=DATABASE_PORT,
|
||||
user=DATABASE_USER,
|
||||
password=DATABASE_PASSWORD,
|
||||
database=DATABASE_NAME,
|
||||
charset="utf8mb4",
|
||||
autocommit=False,
|
||||
connect_timeout=15,
|
||||
read_timeout=60,
|
||||
write_timeout=60,
|
||||
)
|
||||
except Exception as exc:
|
||||
raise ProcessingError(f"MySQL database setup failed: {exc}") from exc
|
||||
) + ")"
|
||||
|
||||
|
||||
class ApiSource:
|
||||
def __init__(self, base_url: str, token: str, storage_id: str, root: str) -> None:
|
||||
if not token:
|
||||
raise ProcessingError("METRIX_API_TOKEN is required in API mode")
|
||||
self.base_url = base_url.rstrip("/")
|
||||
self.token = token
|
||||
def __init__(self, client: MetrixApiClient, storage_id: str, root: str) -> None:
|
||||
self.client = client
|
||||
self.storage_id = storage_id
|
||||
self.root = root.rstrip("/") or "/"
|
||||
|
||||
@@ -345,7 +393,7 @@ class ApiSource:
|
||||
|
||||
def download(self, path: str) -> bytes:
|
||||
endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/download"
|
||||
return self._request(endpoint, {"path": path}, timeout=120)
|
||||
return self.client.get_bytes(endpoint, {"path": path}, timeout=120)
|
||||
|
||||
def cell_data_workbooks(self) -> list[tuple[str, bytes]]:
|
||||
workbooks: list[tuple[str, bytes]] = []
|
||||
@@ -357,27 +405,10 @@ class ApiSource:
|
||||
|
||||
def _list_dir(self, path: str) -> list[dict[str, object]]:
|
||||
endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/files"
|
||||
payload = self._request(endpoint, {"path": path, "recursive": "false"})
|
||||
return json.loads(payload.decode("utf-8")).get("entries", [])
|
||||
|
||||
def _request(self, endpoint: str, query: dict[str, object], timeout: int = 30) -> bytes:
|
||||
url = f"{self.base_url}{endpoint}?{urlencode(query)}"
|
||||
request = Request(url, headers={"Authorization": f"Bearer {self.token}", "Accept": "application/json"})
|
||||
last_error: Exception | None = None
|
||||
for attempt in range(3):
|
||||
try:
|
||||
with urlopen(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}")
|
||||
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:
|
||||
@@ -724,6 +755,24 @@ def normalize_cell(value: object) -> str:
|
||||
return str(value).strip()
|
||||
|
||||
|
||||
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="utf-8-sig", newline="") as file:
|
||||
writer = csv.writer(file)
|
||||
@@ -750,11 +799,9 @@ def build_parser() -> argparse.ArgumentParser:
|
||||
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("--lookback-days", type=int, default=int(os.getenv("INTERFERENCE_LOOKBACK_DAYS", "3")))
|
||||
parser.add_argument("--api-base", default=os.getenv("METRIX_API_BASE_URL", "http://172.17.0.1:8000"))
|
||||
parser.add_argument("--api-token", default=os.getenv("METRIX_API_TOKEN", ""))
|
||||
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 MySQL")
|
||||
parser.add_argument("--no-database", action="store_true", help="Generate files without writing the database")
|
||||
return parser
|
||||
|
||||
|
||||
@@ -762,8 +809,17 @@ 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")
|
||||
source: Source = LocalSource(args.source_dir) if args.source_dir else ApiSource(args.api_base, args.api_token, args.storage_id, args.root)
|
||||
store = None if args.no_database else MySQLSummaryStore()
|
||||
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
|
||||
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)
|
||||
return 0
|
||||
|
||||
|
||||
Reference in New Issue
Block a user