diff --git a/core/data-bus/README.md b/core/data-bus/README.md index fb10717..101821b 100644 --- a/core/data-bus/README.md +++ b/core/data-bus/README.md @@ -10,6 +10,7 @@ topic / 时序库表 / 关系表 / 对象桶:换行业只改模板资产(模 | 文件 | 职责 | |------|------| | `templating.py` | 模板命名推导:Kafka topic + 分区策略、TDengine 超级表/子表、PostgreSQL schema/表、MinIO 桶/对象键 | +| `kafka_naming.py` | Kafka topic 命名/分区模板化组件(issue #28):配置外置、点位维度覆盖、按 (topic,partition) 生产路由 | | `tdengine_schema.py` | 时序库 schema 自动生成(超级表 + 每测点子表,**依据点位字典**)+ 批量 INSERT SQL | | `postgres_schema.py` | 关系库 schema(模板 / 模型 / 用户 / 权限)+ 角色授权语句 | | `batch_writer.py` | 批量写入缓冲:批量聚合(默认 5000 条/0.1s,对齐 PRD 5.2「5k/100ms」基线)、幂等去重、失败重试 —— **数据不丢不重** | diff --git a/core/data-bus/__init__.py b/core/data-bus/__init__.py index 588e060..7d51237 100644 --- a/core/data-bus/__init__.py +++ b/core/data-bus/__init__.py @@ -8,6 +8,7 @@ topic / 时序库表 / 关系表 / 对象桶,换行业只改模板资产(YAM 模块: - templating 模板命名推导(Kafka topic/分区、TDengine 表、PG schema、MinIO 桶); +- kafka_naming Kafka topic 命名/分区模板化组件(配置外置 + 点位维度覆盖 + 生产路由); - tdengine_schema 时序超级表/子表 DDL + 批量 INSERT SQL 生成; - postgres_schema 关系库 schema(模板/模型/用户/权限)+ 授权语句; - batch_writer 批量写入缓冲(StoreSink 抽象 / MemorySink / TdengineSink), @@ -24,6 +25,7 @@ from .batch_writer import ( StoreSink, TdengineSink, ) +from .kafka_naming import KafkaTopicNaming from .postgres_schema import generate_grant_ddl, generate_schema_ddl from .tdengine_schema import ( PointSpec, @@ -39,6 +41,7 @@ __all__ = [ "TemplateNaming", "sanitize", "sanitize_sql", + "KafkaTopicNaming", "PointSpec", "generate_supertable_ddl", "generate_subtable_ddls", diff --git a/core/data-bus/config/kafka.template.yaml b/core/data-bus/config/kafka.template.yaml new file mode 100644 index 0000000..03283cd --- /dev/null +++ b/core/data-bus/config/kafka.template.yaml @@ -0,0 +1,30 @@ +# -*- coding: utf-8 -*- +# 模板「Kafka 命名/分区」配置资产示例:ti-cl4(氯化车间/海绵钛,Template-Ti 一期)。 +# +# 说明(issue #28 / PRD 5.2「② 数据总线 + 时序库」): +# - topic 模板:`{topic_prefix}.{device_id}.points`(与 edge-gateway 上行一致); +# - 分区:hash(device_id) % num_partitions,单设备分区内有序; +# - per_point_topics:点位维度覆盖(点位前缀命中 → 独立 topic), +# 用于质量标签(LAB-*)、告警(ALM-*)等独立消费通道; +# - 换行业只改本文件(template / topic_prefix / 分区数),内核零改动。 +template: ti-cl4 +version: 1.0.0 + +kafka: + # topic 前缀(缺省 = template 名);建议形如 iaop.<模板> + topic_prefix: iaop.ti-cl4 + # 每 topic 分区数:按设备数 × 吞吐评估(600 点位 1Hz 建议 ≥ 12) + num_partitions: 12 + # 副本数:生产建议 ≥ 2(联调可 1) + replication_factor: 1 + # 生产者确认:all 保证不丢(配合 spool 断点续传) + acks: all + # 消息保留时长(小时):与数据保留策略(#33)联动 + retention_hours: 168 + + # 点位维度 topic 覆盖(可选):point_id 前缀命中 → 独立 topic + per_point_topics: + - point_id_prefix: "LAB-" + topic: "iaop.ti-cl4.quality.points" # 质量标签独立通道(LIMS 对接) + - point_id_prefix: "ALM-" + topic: "iaop.ti-cl4.alarm.points" # 告警独立通道 diff --git a/core/data-bus/kafka_naming.py b/core/data-bus/kafka_naming.py new file mode 100644 index 0000000..5ede13a --- /dev/null +++ b/core/data-bus/kafka_naming.py @@ -0,0 +1,136 @@ +# -*- coding: utf-8 -*- +"""Kafka topic 命名 / 分区模板化(按模板 + 点位维度)—— issue #28 / PRD 5.2。 + +在 #4(EPIC)templating 命名雏形之上,交付**完整模板化 Kafka 组件**: + +- **topic**:`{topic_prefix}.{device_id}.points`(与 edge-gateway 上行 + `upstream/kafka_sink.py` 完全一致,链路互通); +- **分区**:`hash(device_id) % num_partitions` 一致性哈希, + 保证单设备分区内严格有序(关键:趋势/告警按设备有序消费); +- **点位维度覆盖**:特定点位前缀可路由到独立 topic + (如质量标签 `LAB-*` → 质检专用 topic),模板配置驱动; +- **生产路由**:`routing(rows)` 按 (topic, partition) 批量分组, + 供 Kafka 生产者批量发送(对齐 PRD 5.2「5k/100ms」批量基线)。 + +全部参数来自模板配置资产 `config/kafka.template.yaml`:换行业只改配置, +内核零改动。 +""" +from __future__ import annotations + +import os +from typing import Dict, List, Optional, Sequence, Tuple + +from data_bus.templating import TemplateNaming, sanitize + +#: 默认配置资产路径(相对本模块) +DEFAULT_CONFIG_PATH = os.path.join( + os.path.dirname(os.path.abspath(__file__)), "config", "kafka.template.yaml") + + +class KafkaTopicNaming: + """Kafka topic/分区模板化命名(模板 + 点位维度)。""" + + def __init__( + self, + template: str, + topic_prefix: Optional[str] = None, + num_partitions: int = 12, + replication_factor: int = 1, + acks: str = "all", + retention_hours: int = 168, + per_point_topics: Optional[List[dict]] = None, + ) -> None: + """ + Args: + template: 行业模板名(如 ti-cl4 / resin); + topic_prefix: topic 前缀(缺省 = 模板名); + num_partitions: 每 topic 分区数(分区策略 hash(device_id)%n); + replication_factor: 副本数(生产建议 ≥ 2,联调可 1); + acks: 生产者确认级别(all 保证不丢); + retention_hours: 消息保留时长(小时); + per_point_topics: 点位维度 topic 覆盖规则 + `[{"point_id_prefix": "LAB-", "topic": "..."}]`。 + """ + self._tpl = TemplateNaming( + template, topic_prefix=topic_prefix, num_partitions=num_partitions) + self.num_partitions = self._tpl.num_partitions + self.replication_factor = max(1, int(replication_factor)) + self.acks = acks + self.retention_hours = max(1, int(retention_hours)) + # 点位前缀 → 覆盖 topic(前缀长优先匹配) + self._per_point: List[tuple] = sorted( + ((str(r["point_id_prefix"]), sanitize(r["topic"])) + for r in (per_point_topics or []) if r.get("point_id_prefix")), + key=lambda kv: len(kv[0]), reverse=True, + ) + + # ------------------------------------------------------------------ + @classmethod + def from_template_config(cls, path: str = DEFAULT_CONFIG_PATH) -> "KafkaTopicNaming": + """从模板配置资产加载(config/kafka.template.yaml)。""" + import yaml + with open(path, "r", encoding="utf-8") as fh: + raw = yaml.safe_load(fh) or {} + k = raw.get("kafka", {}) or {} + return cls( + template=str(raw.get("template", "default")), + topic_prefix=k.get("topic_prefix"), + num_partitions=int(k.get("num_partitions", 12)), + replication_factor=int(k.get("replication_factor", 1)), + acks=str(k.get("acks", "all")), + retention_hours=int(k.get("retention_hours", 168)), + per_point_topics=k.get("per_point_topics", []), + ) + + # ------------------------------------------------------------------ + def topic(self, device_id: str, point_id: Optional[str] = None) -> str: + """上行 topic:`{topic_prefix}.{device_id}.points`。 + + 与 edge-gateway `upstream/kafka_sink.py` 完全一致(device_id 原样, + 不做小写化——Kafka topic 允许大写,点位字典校验保证合法字符)。 + 点位维度覆盖:point_id 命中某前缀规则时路由到覆盖 topic + (如质量标签 LAB-* → 质检专用 topic),否则走设备默认 topic。 + """ + if point_id is not None: + for prefix, topic in self._per_point: + if point_id.startswith(prefix): + return topic + return f"{self._tpl.topic_prefix}.{device_id}.points" + + def partition(self, device_id: str) -> int: + """分区:device_id 一致性哈希(单设备分区内有序)。""" + return self._tpl.partition(device_id) + + def partitions(self, device_ids: Sequence[str]) -> dict: + """设备 → 分区映射(配置台预览)。""" + return self._tpl.partitions(list(device_ids), self.num_partitions) + + # ------------------------------------------------------------------ + def routing(self, rows: Sequence[dict]) -> Dict[Tuple[str, int], List[dict]]: + """按 (topic, partition) 批量分组样本行(Kafka 生产者路由)。 + + Args: + rows: 样本行 `[{"device_id", "point_id", "value", "ts", ...}]`; + 无 point_id 的行按设备维度路由。 + Returns: + {(topic, partition): [rows]} —— 每组可批量发送。 + """ + groups: Dict[Tuple[str, int], List[dict]] = {} + for row in rows: + device_id = str(row.get("device_id", "")) + point_id = row.get("point_id") + key = (self.topic(device_id, point_id), self.partition(device_id)) + groups.setdefault(key, []).append(row) + return groups + + def brief(self) -> dict: + """配置摘要(部署/巡检用)。""" + return { + "template": self._tpl.template, + "topic_prefix": self._tpl.topic_prefix, + "num_partitions": self.num_partitions, + "replication_factor": self.replication_factor, + "acks": self.acks, + "retention_hours": self.retention_hours, + "per_point_rules": [p for p, _ in self._per_point], + } diff --git a/core/data-bus/tests/test_kafka_naming.py b/core/data-bus/tests/test_kafka_naming.py new file mode 100644 index 0000000..86850e1 --- /dev/null +++ b/core/data-bus/tests/test_kafka_naming.py @@ -0,0 +1,137 @@ +# -*- coding: utf-8 -*- +"""Kafka topic 命名/分区模板化测试(issue #28)。 + +覆盖: +1. 模板配置资产加载(config/kafka.template.yaml,含点位维度覆盖规则); +2. topic 推导:`{topic_prefix}.{device_id}.points`(与 edge-gateway 上行一致); +3. 分区:device_id 一致性哈希(确定性、单设备分区内有序); +4. 点位维度覆盖:LAB-* 前缀 → 质检独立 topic,其余走设备 topic; +5. 生产路由:按 (topic, partition) 批量分组。 +""" +import os +import sys +import unittest + +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) +import _bootstrap # noqa: F401 + +from data_bus.kafka_naming import ( # noqa: E402 + DEFAULT_CONFIG_PATH, + KafkaTopicNaming, +) + +CONFIG = os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), + "config", "kafka.template.yaml", +) + + +def sample(device_id="CLF-01", point_id="CLF-01.TEMP", value=850.5, ts=1.0): + return {"device_id": device_id, "point_id": point_id, + "value": value, "ts": ts} + + +class TestConfigLoad(unittest.TestCase): + """模板配置资产加载与摘要。""" + + def test_from_template_config(self): + naming = KafkaTopicNaming.from_template_config(CONFIG) + self.assertEqual(naming._tpl.template, "ti-cl4") + self.assertEqual(naming.num_partitions, 12) + self.assertEqual(naming.acks, "all") + self.assertEqual(naming.retention_hours, 168) + brief = naming.brief() + self.assertIn("LAB-", brief["per_point_rules"]) + self.assertIn("ALM-", brief["per_point_rules"]) + + def test_default_config_path_exists(self): + self.assertTrue(os.path.isfile(DEFAULT_CONFIG_PATH)) + + +class TestTopicNaming(unittest.TestCase): + """topic 推导(与 edge-gateway 上行格式一致)。""" + + def setUp(self): + self.naming = KafkaTopicNaming.from_template_config(CONFIG) + + def test_device_topic_format(self): + # 与 upstream/kafka_sink.py 的 {prefix}.{device}.points 一致 + self.assertEqual(self.naming.topic("CLF-01"), + "iaop.ti-cl4.CLF-01.points") + self.assertEqual(self.naming.topic("S7-01"), + "iaop.ti-cl4.S7-01.points") + + def test_point_level_topic_override(self): + # LAB-* 前缀命中质量标签独立通道;普通点走设备 topic + self.assertEqual( + self.naming.topic("CLF-01", point_id="LAB-01.TI_PURITY"), + "iaop.ti-cl4.quality.points") + self.assertEqual( + self.naming.topic("CLF-01", point_id="ALM-01.TEMP_HI"), + "iaop.ti-cl4.alarm.points") + self.assertEqual( + self.naming.topic("CLF-01", point_id="CLF-01.TEMP"), + "iaop.ti-cl4.CLF-01.points") + + def test_sanitize_device_id(self): + # device_id 原样保留(与 edge-gateway 上行一致,点位字典保证合法) + self.assertEqual(self.naming.topic("CLF-01"), + "iaop.ti-cl4.CLF-01.points") + + +class TestPartition(unittest.TestCase): + """分区一致性哈希:确定性 + 单设备有序。""" + + def setUp(self): + self.naming = KafkaTopicNaming.from_template_config(CONFIG) + + def test_deterministic(self): + self.assertEqual(self.naming.partition("CLF-01"), + self.naming.partition("CLF-01")) + self.assertIn(self.naming.partition("CLF-01"), range(12)) + + def test_same_device_same_partition(self): + # 单设备所有点位落到同一分区 → 分区内有序 + p1 = self.naming.partition("CLF-01") + for pid in ("CLF-01.TEMP", "CLF-01.PRES", "CLF-01.FEED"): + self.assertEqual(self.naming.partition("CLF-01"), p1) + + def test_partitions_preview(self): + mapping = self.naming.partitions(["CLF-01", "S7-01", "E-01"]) + self.assertEqual(set(mapping), {"CLF-01", "S7-01", "E-01"}) + self.assertTrue(all(0 <= v < 12 for v in mapping.values())) + + +class TestRouting(unittest.TestCase): + """按 (topic, partition) 批量分组(生产路由)。""" + + def setUp(self): + self.naming = KafkaTopicNaming.from_template_config(CONFIG) + + def test_group_by_topic_and_partition(self): + rows = [ + sample("CLF-01", "CLF-01.TEMP"), + sample("CLF-01", "CLF-01.PRES"), + sample("CLF-01", "LAB-01.TI_PURITY"), # 质量标签 → 独立 topic + sample("S7-01", "S7-01.PUMP_A"), + ] + groups = self.naming.routing(rows) + # 3 个不同 (topic, partition) 组:CLF 默认 / LAB 质量 / S7 + self.assertEqual(len(groups), 3) + topics = {t for t, _ in groups} + self.assertIn("iaop.ti-cl4.quality.points", topics) + # CLF 默认 topic 内 2 条同一分区 + clf_key = ("iaop.ti-cl4.CLF-01.points", + self.naming.partition("CLF-01")) + self.assertEqual(len(groups[clf_key]), 2) + + def test_rows_without_point_id(self): + rows = [{"device_id": "CLF-01", "value": 1.0, "ts": 2.0}] + groups = self.naming.routing(rows) + self.assertEqual(len(groups), 1) + key = next(iter(groups)) + self.assertEqual(key[0], "iaop.ti-cl4.CLF-01.points") + + +if __name__ == "__main__": + unittest.main()