2026-08-04 15:46:09 +08:00
|
|
|
|
# -*- 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
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-04 18:14:37 +08:00
|
|
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# 依据点位字典自动生成完整 schema(子任务 #29 交付入口)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
|
|
def specs_from_point_dict_rows(rows: List[dict]) -> List[PointSpec]:
|
|
|
|
|
|
"""把点位字典 CSV 行(edge-gateway 点位字典表头)转换为 PointSpec 列表。
|
|
|
|
|
|
|
|
|
|
|
|
与 `core/edge-gateway/point_dict/loader.py` 的 CSV 表头对齐:
|
|
|
|
|
|
``device_id / point_id / unit / dataType``(兼容小写 ``data_type``)。
|
|
|
|
|
|
其余列(name / sampleRate / qualityCode / opcNode / protocol)与
|
|
|
|
|
|
schema 生成无关,宽松忽略;非法行按 ``from_dict`` 宽松兜底,
|
|
|
|
|
|
不做静默丢弃(保持与点位数一致,便于外层校验定位)。
|
|
|
|
|
|
|
|
|
|
|
|
Args:
|
|
|
|
|
|
rows: 点位字典行(每行一个测点,dict 键为 CSV 表头)。
|
|
|
|
|
|
|
|
|
|
|
|
Returns:
|
|
|
|
|
|
与输入顺序一致的 PointSpec 列表。
|
|
|
|
|
|
"""
|
|
|
|
|
|
return [PointSpec.from_dict(raw) for raw in rows]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def generate_schema(
|
|
|
|
|
|
naming: TemplateNaming,
|
|
|
|
|
|
points: List[PointSpec],
|
|
|
|
|
|
retention_days: Optional[int] = None,
|
|
|
|
|
|
) -> List[str]:
|
|
|
|
|
|
"""依据点位字典一次性生成完整 TDengine schema(超级表 + 每测点子表)。
|
|
|
|
|
|
|
|
|
|
|
|
Args:
|
|
|
|
|
|
naming: 模板命名器(决定超级表/子表名)。
|
|
|
|
|
|
points: 点位字典(测点维度集合,可由 specs_from_point_dict_rows 得到)。
|
|
|
|
|
|
retention_days: 数据保留天数(KEEP,可选)。
|
|
|
|
|
|
|
|
|
|
|
|
Returns:
|
|
|
|
|
|
DDL 语句列表:第 0 条为超级表 CREATE STABLE,其余为子表
|
|
|
|
|
|
CREATE TABLE ... USING ... TAGS(...),可直接按序执行。
|
|
|
|
|
|
"""
|
|
|
|
|
|
return [generate_supertable_ddl(naming, retention_days=retention_days)] + generate_subtable_ddls(
|
|
|
|
|
|
naming, points
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-08-04 15:46:09 +08:00
|
|
|
|
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
|