# -*- 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], }