# -*- coding: utf-8 -*- """TDengine 超级表 schema 自动生成 + 批量 INSERT —— 依据点位字典(子任务 #29)。 模型(TDengine 2.0/3.0 通用): - 一张超级表 `{tpl}_points`:列 `ts / value / quality`,标签 `device_id / unit / data_type`; - 点位字典每个测点自动生成一张子表 `{tpl}_pt_{point_id}`(点位维度, 字典 CSV 变更即重建子表集,内核零改动); - 批量写入:按子表聚合多行 `INSERT INTO {sub} VALUES (...),(...);`, 单条样本约 512B 时对齐 PRD 5.2「5k 条/100ms」基线。 不依赖 taospy:本模块只负责生成 DDL / 批量 SQL;真实落库由 `batch_writer.StoreSink` 的客户端实现(缺失时降级内存 sink 联调)。 """ from __future__ import annotations from dataclasses import dataclass from typing import Dict, List, Optional from .templating import TemplateNaming, sanitize_sql, _sql_str # dataType → TDengine 列类型(点位字典 schema.VALID_DATA_TYPES 子集) TD_VALUE_TYPE: Dict[str, str] = { "float": "DOUBLE", "int": "BIGINT", "bool": "BOOL", } DEFAULT_VALUE_TYPE = "DOUBLE" # 样本记录字段(与 edge-gateway spool / kafka_sink 保持一致) ROW_KEYS = ("device_id", "point_id", "value", "ts", "quality") @dataclass(frozen=True) class PointSpec: """点位字典中的最小维度信息(schema 生成所需)。""" device_id: str point_id: str unit: str = "" data_type: str = "float" @classmethod def from_dict(cls, raw: dict) -> "PointSpec": return cls( device_id=str(raw.get("device_id") or ""), point_id=str(raw.get("point_id") or ""), unit=str(raw.get("unit") or ""), data_type=str(raw.get("data_type") or raw.get("dataType") or "float"), ) def value_type(data_type: str) -> str: return TD_VALUE_TYPE.get((data_type or "").lower(), DEFAULT_VALUE_TYPE) def generate_supertable_ddl( naming: TemplateNaming, retention_days: Optional[int] = None, ) -> str: """生成超级表 DDL(自动建表,幂等 IF NOT EXISTS)。 Args: naming: 模板命名器(决定超级表名)。 retention_days: 数据保留天数(TDengine KEEP 表选项,可选)。 Returns: CREATE STABLE 语句。 """ ddl = ( f"CREATE STABLE IF NOT EXISTS {naming.stable()} (\n" " ts TIMESTAMP,\n" " value DOUBLE,\n" " quality TINYINT\n" ") TAGS (\n" " device_id NCHAR(64),\n" " unit NCHAR(16),\n" " data_type NCHAR(16)\n" ")" ) if retention_days and int(retention_days) > 0: ddl += f" KEEP({int(retention_days)})" return ddl + ";" def generate_subtable_ddls( naming: TemplateNaming, points: List[PointSpec], use_typed_value: bool = True, ) -> List[str]: """按点位字典为每个测点生成子表 DDL(依据点位字典自动生成)。 Args: naming: 模板命名器。 points: 点位字典(测点维度集合)。 use_typed_value: 是否按 dataType 派生 value 列类型(TDengine 3.x 支持列级类型;2.x 超级表统一 DOUBLE,传 False 时忽略)。 Returns: CREATE TABLE ... USING ... TAGS(...) 语句列表(顺序与 points 一致)。 """ if use_typed_value: # 说明:超级表 value 列按最宽类型 DOUBLE 建(TDengine 列类型在 # CREATE STABLE 时固定),子表继承,这里保留 data_type 标签供聚合。 pass out: List[str] = [] seen: set = set() for p in points: sub = naming.subtable(p.point_id) if sub in seen: continue seen.add(sub) tags = ", ".join( _sql_str(v) for v in (p.device_id, p.unit or "", (p.data_type or "float").lower()) ) out.append( f"CREATE TABLE IF NOT EXISTS {sub} USING {naming.stable()} " f"TAGS ({tags});" ) return out def build_batch_insert( naming: TemplateNaming, rows: List[dict], ) -> List[str]: """把样本批聚合为按子表分组的批量 INSERT 语句。 Args: naming: 模板命名器(决定子表名)。 rows: 样本列表,每条含 device_id/point_id/value/ts(/quality)。 Returns: SQL 语句列表:每子表一条 `INSERT INTO {sub} VALUES (...),(...);`。 空输入返回空列表。 """ grouped: Dict[str, List[dict]] = {} for row in rows: point_id = str(row.get("point_id") or "") sub = naming.subtable(point_id) grouped.setdefault(sub, []).append(row) statements: List[str] = [] for sub in sorted(grouped): tuples = [] for row in grouped[sub]: ts = row.get("ts") ts_ms = int(ts * 1000) if isinstance(ts, float) else int(ts) value = row.get("value") value_sql = "NULL" if value is None else repr(float(value)) quality = int(row.get("quality", 1)) tuples.append(f"({ts_ms}, {value_sql}, {quality})") statements.append( f"INSERT INTO {sub} VALUES " + ", ".join(tuples) + ";" ) return statements