feat: 完成 issue #23 ① 采集配置外置(点位/采样率/协议参数化,去硬编码)
This commit is contained in:
@@ -62,7 +62,11 @@ python -m unittest discover -s tests
|
|||||||
## 模板化说明(换行业只改配置,不改代码)
|
## 模板化说明(换行业只改配置,不改代码)
|
||||||
|
|
||||||
- 点位范围 / 采样率 / 量纲 / 协议:由**点位字典 CSV** 驱动(模板 → 行业点位集)。
|
- 点位范围 / 采样率 / 量纲 / 协议:由**点位字典 CSV** 驱动(模板 → 行业点位集)。
|
||||||
协议列(可选)提供点位级协议覆盖;留空则按 `gateway.yaml` drivers 段前缀路由。
|
- 采样率(issue #23):点位字典 CSV 的 `sampleRate` 列逐点位控制采集频率,
|
||||||
|
引擎以 `gateway.yaml` 的 `interval_ms` 为基准 tick,`sampleRate > interval_ms`
|
||||||
|
的点位按比例降频(如 5000ms / 1000ms tick → 每 5 tick 采一次);
|
||||||
|
`sampleRate ≤ interval_ms` 的点位每 tick 采集。点位采样率完全外置、零硬编码。
|
||||||
|
- 协议列(可选)提供点位级协议覆盖;留空则按 `gateway.yaml` drivers 段前缀路由。
|
||||||
- 协议选型与连接参数:由 `gateway.yaml` 的 `drivers` 段驱动(模板 → 行业协议栈)。
|
- 协议选型与连接参数:由 `gateway.yaml` 的 `drivers` 段驱动(模板 → 行业协议栈)。
|
||||||
- Kafka topic 命名:`{template}.{device}.points`,随配置模板变化。
|
- Kafka topic 命名:`{template}.{device}.points`,随配置模板变化。
|
||||||
- 新增协议:实现 `drivers/base.py` 的 `Driver` 子类并注册即可,引擎与上行链路零改动。
|
- 新增协议:实现 `drivers/base.py` 的 `Driver` 子类并注册即可,引擎与上行链路零改动。
|
||||||
|
|||||||
@@ -2,7 +2,10 @@
|
|||||||
"""周期采集调度引擎 —— 只读采集 + 背压保护 + 健康度上报。
|
"""周期采集调度引擎 —— 只读采集 + 背压保护 + 健康度上报。
|
||||||
|
|
||||||
模板化封装要点:
|
模板化封装要点:
|
||||||
- 采集配置全部外置(gateway.yaml),引擎不关心具体协议;
|
- 采集配置全部外置(gateway.yaml + 点位字典 CSV),引擎不关心具体协议;
|
||||||
|
- 采样率外置(issue #23):每个点位按点位字典 CSV 的 sampleRate 列调度,
|
||||||
|
引擎以 gateway.yaml 的 interval_ms 为基准 tick,sampleRate > interval_ms 的
|
||||||
|
点位按比例降频读取(去硬编码:点位采样率完全由模板配置驱动);
|
||||||
- 点位按驱动实例的 device 匹配规则分组,一次 tick 内按驱动批量读取;
|
- 点位按驱动实例的 device 匹配规则分组,一次 tick 内按驱动批量读取;
|
||||||
- 严格只读:引擎只调用 Driver.read_points(),不存在任何控制指令路径;
|
- 严格只读:引擎只调用 Driver.read_points(),不存在任何控制指令路径;
|
||||||
- 背压保护:未确认(spool 待上行)记录超过阈值时丢弃新样本并计入丢失,
|
- 背压保护:未确认(spool 待上行)记录超过阈值时丢弃新样本并计入丢失,
|
||||||
@@ -55,6 +58,28 @@ class CollectorEngine:
|
|||||||
self._stop = threading.Event()
|
self._stop = threading.Event()
|
||||||
self._thread: Optional[threading.Thread] = None
|
self._thread: Optional[threading.Thread] = None
|
||||||
self._point_to_driver = self._build_routing()
|
self._point_to_driver = self._build_routing()
|
||||||
|
# 每点位采样调度(issue #23):tick 计数 + 各点位下次应采集的 tick
|
||||||
|
self._tick = 0
|
||||||
|
self._next_due_tick: Dict[str, int] = {p.point_id: 0 for p in point_dict.points}
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
def _ticks_per_sample(self, p: Point) -> int:
|
||||||
|
"""点位采样间隔(以基准 tick 计)。
|
||||||
|
|
||||||
|
sampleRate ≤ interval_ms 时每个 tick 都采;
|
||||||
|
sampleRate > interval_ms 时按比例降频(四舍五入,至少 1 tick)。
|
||||||
|
"""
|
||||||
|
if p.sample_rate <= 0:
|
||||||
|
return 1
|
||||||
|
return max(1, int(round(p.sample_rate / self.interval_ms)))
|
||||||
|
|
||||||
|
def _due_points(self) -> List[Point]:
|
||||||
|
"""本轮 tick 到期待采集的点位(按点位 sampleRate 调度)。"""
|
||||||
|
due: List[Point] = []
|
||||||
|
for p in self.point_dict.points:
|
||||||
|
if self._tick >= self._next_due_tick.get(p.point_id, 0):
|
||||||
|
due.append(p)
|
||||||
|
return due
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
def _build_routing(self) -> Dict[str, Driver]:
|
def _build_routing(self) -> Dict[str, Driver]:
|
||||||
@@ -92,13 +117,16 @@ class CollectorEngine:
|
|||||||
本轮成功读到的样本数。
|
本轮成功读到的样本数。
|
||||||
"""
|
"""
|
||||||
started = time.monotonic()
|
started = time.monotonic()
|
||||||
expected = len(self.point_dict.points)
|
# 本轮 tick 递增 + 取到期待集(按点位 sampleRate 调度,issue #23)
|
||||||
|
self._tick += 1
|
||||||
|
due = self._due_points()
|
||||||
|
expected = len(due)
|
||||||
got = 0
|
got = 0
|
||||||
samples: List[dict] = []
|
samples: List[dict] = []
|
||||||
|
|
||||||
# 1) 按驱动分组读取(只读)
|
# 1) 按驱动分组读取(只读)
|
||||||
by_driver: Dict[Driver, List[Point]] = {}
|
by_driver: Dict[Driver, List[Point]] = {}
|
||||||
for p in self.point_dict.points:
|
for p in due:
|
||||||
drv = self._point_to_driver.get(p.point_id)
|
drv = self._point_to_driver.get(p.point_id)
|
||||||
if drv is None:
|
if drv is None:
|
||||||
continue
|
continue
|
||||||
@@ -112,6 +140,8 @@ class CollectorEngine:
|
|||||||
self.metrics.record_round(time.monotonic() - started, expected, got, failed=True)
|
self.metrics.record_round(time.monotonic() - started, expected, got, failed=True)
|
||||||
return got
|
return got
|
||||||
for p in points:
|
for p in points:
|
||||||
|
# 无论本轮是否读到,都推进该点位采样调度(降频点位不连续空读)
|
||||||
|
self._next_due_tick[p.point_id] = self._tick + self._ticks_per_sample(p)
|
||||||
value = values.get(p.point_id)
|
value = values.get(p.point_id)
|
||||||
if value is None:
|
if value is None:
|
||||||
continue # 未读到 → 计入丢失
|
continue # 未读到 → 计入丢失
|
||||||
|
|||||||
@@ -1,7 +1,10 @@
|
|||||||
# iAOP 边缘采集网关 —— 模板配置示例(Template-Ti 一期:氯化车间/海绵钛)
|
# iAOP 边缘采集网关 —— 模板配置示例(Template-Ti 一期:氯化车间/海绵钛)
|
||||||
# 换行业模板时只需修改本文件 + 点位字典 CSV,网关代码零改动。
|
# 换行业模板时只需修改本文件 + 点位字典 CSV,网关代码零改动。
|
||||||
collector:
|
collector:
|
||||||
# 采集周期(毫秒):600 点位 1Hz 即 1000
|
# 采集周期(毫秒)基准 tick:600 点位 1Hz 即 1000。
|
||||||
|
# 点位级采样率由点位字典 CSV 的 sampleRate 列驱动(issue #23):
|
||||||
|
# sampleRate ≤ interval_ms 的点位每 tick 采集;> interval_ms 的点位按比例降频
|
||||||
|
# (如 sampleRate=5000 且 interval_ms=1000 → 每 5 tick 采一次),去硬编码。
|
||||||
interval_ms: 1000
|
interval_ms: 1000
|
||||||
# 背压保护:spool 待上行记录上限,超过则丢弃新样本并计入丢失率
|
# 背压保护:spool 待上行记录上限,超过则丢弃新样本并计入丢失率
|
||||||
max_pending: 100000
|
max_pending: 100000
|
||||||
|
|||||||
@@ -3,9 +3,9 @@ CLF-01,CLF-01.TEMP,炉温,℃,float,1000,true,ns=2;s=CLF.Temp,opcua
|
|||||||
CLF-01,CLF-01.PRES,炉压,kPa,float,1000,true,ns=2;s=CLF.Pres,opcua
|
CLF-01,CLF-01.PRES,炉压,kPa,float,1000,true,ns=2;s=CLF.Pres,opcua
|
||||||
CLF-01,CLF-01.FEED,进料量,t/h,float,1000,true,ns=2;s=CLF.Feed,opcua
|
CLF-01,CLF-01.FEED,进料量,t/h,float,1000,true,ns=2;s=CLF.Feed,opcua
|
||||||
CLF-02,CLF-02.TEMP,炉温2,℃,float,1000,true,ns=2;s=CLF2.Temp,
|
CLF-02,CLF-02.TEMP,炉温2,℃,float,1000,true,ns=2;s=CLF2.Temp,
|
||||||
CLF-02,CLF-02.RUN,运行状态,%,bool,1000,true,ns=2;s=CLF2.Run,
|
CLF-02,CLF-02.RUN,运行状态,%,bool,2000,true,ns=2;s=CLF2.Run,
|
||||||
S7-01,S7-01.PUMP_A,泵A频率,Hz,float,1000,true,DB100.0.0,s7
|
S7-01,S7-01.PUMP_A,泵A频率,Hz,float,1000,true,DB100.0.0,s7
|
||||||
S7-01,S7-01.PUMP_B,泵B频率,Hz,float,1000,true,DB100.4.0,s7
|
S7-01,S7-01.PUMP_B,泵B频率,Hz,float,1000,true,DB100.4.0,s7
|
||||||
W-01,W-01.WT,称重值,kg,float,1000,true,holding:40001,weighing
|
W-01,W-01.WT,称重值,kg,float,1000,true,holding:40001,weighing
|
||||||
E-01,E-01.PWR,电表功率,kW,float,1000,true,holding:40010,energy
|
E-01,E-01.PWR,电表功率,kW,float,1000,true,holding:40010,energy
|
||||||
E-01,E-01.KWH,电表累计,kWh,float,1000,true,holding:40012,energy
|
E-01,E-01.KWH,电表累计,kWh,float,5000,true,holding:40012,energy
|
||||||
|
|||||||
|
@@ -146,6 +146,52 @@ class EngineMetricsTest(unittest.TestCase):
|
|||||||
self.assertLessEqual(metrics.p99_latency(), 0.6)
|
self.assertLessEqual(metrics.p99_latency(), 0.6)
|
||||||
self.assertEqual(metrics.total_samples, 100)
|
self.assertEqual(metrics.total_samples, 100)
|
||||||
|
|
||||||
|
def test_per_point_sample_rate_schedule(self):
|
||||||
|
"""点位级 sampleRate 调度(issue #23):低频点位按比例降频,去硬编码。"""
|
||||||
|
# interval_ms=1000:A 点 1s 采样、B 点 5s 采样、C 点 500ms(小于 tick,退化为每 tick)
|
||||||
|
pd = PointDict([
|
||||||
|
Point(device_id="CLF-01", point_id="CLF-01.A", name="A", unit="℃",
|
||||||
|
data_type="float", sample_rate=1000, quality_code=True, row_number=2),
|
||||||
|
Point(device_id="CLF-01", point_id="CLF-01.B", name="B", unit="℃",
|
||||||
|
data_type="float", sample_rate=5000, quality_code=True, row_number=3),
|
||||||
|
Point(device_id="CLF-01", point_id="CLF-01.C", name="C", unit="℃",
|
||||||
|
data_type="float", sample_rate=500, quality_code=True, row_number=4),
|
||||||
|
])
|
||||||
|
spool = SpoolStore(os.path.join(self._tmp, "spool"))
|
||||||
|
metrics = HealthMetrics()
|
||||||
|
engine = CollectorEngine(
|
||||||
|
point_dict=pd,
|
||||||
|
driver_slots=[("simulator", SimulatorDriver(), [])],
|
||||||
|
spool=spool,
|
||||||
|
metrics=metrics,
|
||||||
|
interval_ms=1000,
|
||||||
|
)
|
||||||
|
# 记录每个点位在 10 个 tick 内被读取的次数(读驱动时计数)
|
||||||
|
read_counts = {"CLF-01.A": 0, "CLF-01.B": 0, "CLF-01.C": 0}
|
||||||
|
original = SimulatorDriver.read_points
|
||||||
|
|
||||||
|
def counting_read(self, points):
|
||||||
|
for p in points:
|
||||||
|
read_counts[p.point_id] += 1
|
||||||
|
return original(self, points)
|
||||||
|
|
||||||
|
SimulatorDriver.read_points = counting_read
|
||||||
|
try:
|
||||||
|
for _ in range(10):
|
||||||
|
engine.collect_once()
|
||||||
|
finally:
|
||||||
|
SimulatorDriver.read_points = original
|
||||||
|
# A:1000ms / 1000ms tick → 每 tick 都读 = 10 次
|
||||||
|
self.assertEqual(read_counts["CLF-01.A"], 10)
|
||||||
|
# B:5000ms / 1000ms tick → 每 5 tick 读一次 ≈ 2 次(tick1、tick6)
|
||||||
|
self.assertEqual(read_counts["CLF-01.B"], 2)
|
||||||
|
# C:500ms < tick → 每 tick 读 = 10 次(不做高于 tick 频率的超采样)
|
||||||
|
self.assertEqual(read_counts["CLF-01.C"], 10)
|
||||||
|
# 调度推进:低频点位不再被硬编码为每 tick 读取
|
||||||
|
self.assertEqual(engine._ticks_per_sample(pd.points[0]), 1)
|
||||||
|
self.assertEqual(engine._ticks_per_sample(pd.points[1]), 5)
|
||||||
|
self.assertEqual(engine._ticks_per_sample(pd.points[2]), 1)
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
Reference in New Issue
Block a user