Files
iAOP/core/edge-gateway/main.py
T

185 lines
7.1 KiB
Python
Raw Normal View History

# -*- coding: utf-8 -*-
"""iAOP 边缘采集网关 —— 主入口。
用法:
python main.py --config config/gateway.example.yaml --point-dict point_dict.csv [--dry-run]
流程(模板化封装,对齐 PRD 5.1 用户操作流程):
实施工程师导入 DCS 点表 CSV → 自动校验(缺失字段/量纲/重复点号)
→ 加载模板配置(协议/周期/背压阈值/Kafka)→ 网关启动只读采集
→ Kafka 流式上行 + spool 断点续传 → 实时健康度上报。
"""
from __future__ import annotations
import argparse
import logging
import sys
import time
from typing import List, Tuple
import yaml
# 允许直接以脚本方式运行(python main.py)时仍能解析包内模块
from point_dict import load_point_dict_csv, validate_point_dict_file
from point_dict.loader import PointDict
logger = logging.getLogger("edge_gateway.main")
def parse_args(argv: List[str]) -> argparse.Namespace:
parser = argparse.ArgumentParser(description="iAOP 边缘采集网关(模板化封装)")
parser.add_argument("--config", required=True, help="采集配置 YAML(模板参数化)")
parser.add_argument("--point-dict", required=True, help="点位字典 CSV(客户 DCS 点表)")
parser.add_argument("--dry-run", action="store_true",
help="仅加载配置并校验点位字典,不启动采集")
parser.add_argument("--rounds", type=int, default=0,
help="采集轮数上限(0=无限,调试用)")
parser.add_argument("--verbose", action="store_true", help="输出调试日志")
parser.add_argument("--log-json", action="store_true",
help="结构化 JSON 日志(NFR 9 章统一日志规范;生产模式建议开启)")
return parser.parse_args(argv)
def build_engine(config: dict, point_dict: PointDict):
"""按模板配置组装采集引擎(驱动路由 → spool → metrics → 引擎)。"""
from collector import CollectorEngine, HealthMetrics, SpoolStore
collector_cfg = config.get("collector", {})
spool = SpoolStore(
spool_dir=collector_cfg.get("spool_dir", "./spool"),
cache_limit_bytes=int(collector_cfg.get("cache_limit_bytes", 512 * 1024 * 1024)),
)
metrics = HealthMetrics()
# 组装驱动插槽:[(protocol, driver, device_prefixes)]
driver_slots: List[Tuple[str, object, List[str]]] = []
from drivers import from_template
for item in collector_cfg.get("drivers", []):
protocol = item.get("protocol")
if not protocol:
raise ValueError("collector.drivers[].protocol 必填")
driver = from_template(protocol, item.get("config", {}))
driver_slots.append((protocol, driver, item.get("device_prefixes", [])))
if not driver_slots:
# 模板配置未声明驱动时默认走模拟驱动(本地联调),避免空跑
from drivers import SimulatorDriver
driver_slots.append(("simulator", SimulatorDriver(), []))
engine = CollectorEngine(
point_dict=point_dict,
driver_slots=driver_slots,
spool=spool,
metrics=metrics,
interval_ms=int(collector_cfg.get("interval_ms", 1000)),
max_pending=int(collector_cfg.get("max_pending", 100_000)),
)
return engine, spool, metrics
def build_sink(config: dict, spool) -> object:
"""按模板配置组装 Kafka 上行通道(含 mTLS 传输安全,NFR 9 章)。"""
from upstream import KafkaSink
kafka_cfg = config.get("kafka", {})
return KafkaSink(
bootstrap_servers=kafka_cfg.get("bootstrap_servers", "127.0.0.1:9092"),
topic_prefix=kafka_cfg.get("topic_prefix", "iaop"),
spool=spool,
batch_size=int(kafka_cfg.get("batch_size", 500)),
security=kafka_cfg.get("security") or {},
)
def main(argv: List[str]) -> int:
args = parse_args(argv)
if args.log_json:
# NFR 9 章「可维护」:统一日志规范 —— 结构化 JSON,便于采集到统一日志平台
import json as _json
class JsonFormatter(logging.Formatter):
def format(self, record):
payload = {
"ts": self.formatTime(record, "%Y-%m-%dT%H:%M:%S%z"),
"level": record.levelname,
"logger": record.name,
"msg": record.getMessage(),
}
if record.exc_info:
payload["exc"] = self.formatException(record.exc_info)
return _json.dumps(payload, ensure_ascii=False)
handler = logging.StreamHandler()
handler.setFormatter(JsonFormatter())
logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO, handlers=[handler])
else:
logging.basicConfig(
level=logging.DEBUG if args.verbose else logging.INFO,
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
)
# 1) 加载模板配置
with open(args.config, "r", encoding="utf-8") as fh:
config = yaml.safe_load(fh) or {}
# 2) 加载 + 自动校验点位字典(缺失字段/量纲/重复点号)
report = validate_point_dict_file(args.point_dict)
logger.info("点位字典校验: %s", report.summary())
if not report.ok:
for issue in report.issues[:20]:
logger.error(" [%s] %s", issue.code, issue.message)
if len(report.issues) > 20:
logger.error(" ... 共 %d 条问题", len(report.issues))
print(f"点位字典校验失败:{report.summary()}")
return 2
point_dict = load_point_dict_csv(args.point_dict)
logger.info("点位字典加载完成:%d 个测点", len(point_dict))
if args.dry_run:
print("dry-run 通过:配置与点位字典校验 OK,未启动采集")
return 0
# 3) 组装并启动
engine, spool, metrics = build_engine(config, point_dict)
sink = build_sink(config, spool)
# 断点续传:启动时先重发上次未确认记录
pending = spool.pending_records(limit=10 ** 9)
if pending:
logger.info("检测到 %d 条未确认 spool 记录,启动续传", len(pending))
sink.publish(pending)
engine.start(sink=sink)
logger.info("采集网关已启动:%d 测点 @ %dms,Kafka=%s",
len(point_dict), engine.interval_ms, config.get("kafka", {}).get("bootstrap_servers"))
try:
rounds = 0
while True:
time.sleep(5)
rounds += 5
snap = metrics.snapshot()
logger.info("健康度: 轮次=%d 样本=%d 丢失率=%.4f%% P99=%.3fs 可用性=%.4f%%",
snap["total_rounds"], snap["total_samples"],
snap["loss_rate"] * 100, snap["p99_latency_sec"],
snap["availability"] * 100)
if args.rounds and rounds >= args.rounds:
break
except KeyboardInterrupt:
pass
finally:
engine.stop()
sink.close()
snap = metrics.snapshot()
print(f"SLA 达标(P99≤1.8s/丢失率≤0.02%/可用性≥99.8%): {metrics.meets_sla()}")
print(f"健康度快照: {snap}")
return 0
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))