架构核对发现 2 处差距,本次补齐: - Kafka 上行支持 SSL/SASL_SSL 双向 mTLS(8.2 服务间 mTLS 边缘网关↔总线) - 网关日志支持结构化 JSON 输出(NFR 可维护:统一日志规范) - 配置示例与 README 验收口径同步更新
185 lines
7.1 KiB
Python
185 lines
7.1 KiB
Python
# -*- 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:]))
|