feat: 完成 issue #28 ② Kafka topic 命名/分区模板化(按模板+点位维度)

This commit is contained in:
2026-08-05 00:35:19 +08:00
parent afee113901
commit 793dd0a3b8
5 changed files with 307 additions and 0 deletions
+1
View File
@@ -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」基线)、幂等去重、失败重试 —— **数据不丢不重** |
+3
View File
@@ -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",
+30
View File
@@ -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" # 告警独立通道
+136
View File
@@ -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],
}
+137
View File
@@ -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()