Files
iAOP/core/edge-gateway/main.py
yunmei 098125164d fix: 对齐技术架构补齐传输 mTLS 与结构化 JSON 日志(NFR 9 章)
架构核对发现 2 处差距,本次补齐:
- Kafka 上行支持 SSL/SASL_SSL 双向 mTLS(8.2 服务间 mTLS 边缘网关↔总线)
- 网关日志支持结构化 JSON 输出(NFR 可维护:统一日志规范)
- 配置示例与 README 验收口径同步更新
2026-08-04 15:46:09 +08:00

185 lines
7.1 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- 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:]))