Files
silk/server-go/internal/service/iotdb.go
T

392 lines
10 KiB
Go

package service
import (
"bytes"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"strconv"
"strings"
"sync"
"time"
)
// HistoryRow 历史数据行
type HistoryRow struct {
TS time.Time `json:"ts"`
Value float64 `json:"value"`
}
// BucketRow 聚合桶数据
type BucketRow struct {
Bucket time.Time `json:"bucket"`
Avg float64 `json:"avg"`
Min float64 `json:"min"`
Max float64 `json:"max"`
}
// IoTDBService IoTDB 时序数据库服务,通过 HTTP REST API 查询
type IoTDBService struct {
baseURL string
user string
password string
database string
enabled bool
available bool
mu sync.RWMutex
createdTS sync.Map // 已创建的时间序列缓存
httpClient *http.Client
}
// NewIoTDBService 创建 IoTDB 服务
func NewIoTDBService(baseURL string) *IoTDBService {
return &IoTDBService{
baseURL: strings.TrimRight(baseURL, "/"),
user: "root",
password: "root",
database: "root.silk",
enabled: true,
httpClient: &http.Client{Timeout: 5 * time.Second},
}
}
// IsAvailable 返回 IoTDB 是否可用
func (s *IoTDBService) IsAvailable() bool {
s.mu.RLock()
defer s.mu.RUnlock()
return s.enabled && s.available
}
// Init 初始化 IoTDB(创建数据库)
func (s *IoTDBService) Init() error {
_, err := s.nonQuery(fmt.Sprintf("CREATE DATABASE %s", s.database))
if err != nil && !strings.Contains(strings.ToLower(err.Error()), "already exist") {
s.mu.Lock()
s.available = false
s.mu.Unlock()
return err
}
s.mu.Lock()
s.available = true
s.mu.Unlock()
slog.Info("IoTDB 连接成功", "url", s.baseURL, "database", s.database)
return nil
}
// InsertTelemetry 插入遥测数据
func (s *IoTDBService) InsertTelemetry(deviceKey, metric string, value float64, ts time.Time) bool {
if !s.IsAvailable() || deviceKey == "" || metric == "" {
return false
}
if err := s.ensureTimeseries(deviceKey, metric); err != nil {
slog.Warn("IoTDB ensureTimeseries 失败", "deviceKey", deviceKey, "metric", metric, "err", err)
s.markUnavailable()
return false
}
sql := fmt.Sprintf("INSERT INTO %s.%s(timestamp, %s) VALUES(%d, %g)",
s.database, s.devicePath(deviceKey), s.pathSegment(metric), ts.UnixMilli(), value)
if _, err := s.nonQuery(sql); err != nil {
slog.Warn("IoTDB insert 失败", "deviceKey", deviceKey, "metric", metric, "err", err)
s.markUnavailable()
return false
}
return true
}
// QueryHistory 查询历史数据
func (s *IoTDBService) QueryHistory(deviceKey, metric string, from, to time.Time, limit int) ([]HistoryRow, error) {
if !s.IsAvailable() {
return nil, fmt.Errorf("IoTDB unavailable")
}
if limit <= 0 || limit > 2000 {
limit = 2000
}
sql := fmt.Sprintf("SELECT %s FROM %s WHERE time >= %d AND time <= %d ORDER BY TIME DESC LIMIT %d",
s.pathSegment(metric), s.devicePath(deviceKey), from.UnixMilli(), to.UnixMilli(), limit)
resp, err := s.query(sql, limit)
if err != nil {
return nil, err
}
return s.parseHistoryRows(resp), nil
}
// ListMetrics 列出设备的所有指标
func (s *IoTDBService) ListMetrics(deviceKey string) ([]string, error) {
if !s.IsAvailable() {
return nil, fmt.Errorf("IoTDB unavailable")
}
sql := fmt.Sprintf("SHOW TIMESERIES %s.*", s.devicePath(deviceKey))
resp, err := s.query(sql, 2000)
if err != nil {
return nil, err
}
return s.extractMetrics(resp), nil
}
// QueryLatestAny 查询设备最新一条遥测
func (s *IoTDBService) QueryLatestAny(deviceKey string) (*HistoryRow, error) {
if !s.IsAvailable() {
return nil, fmt.Errorf("IoTDB unavailable")
}
metrics, err := s.ListMetrics(deviceKey)
if err != nil {
return nil, err
}
var picked *HistoryRow
for _, m := range metrics {
rows, err := s.QueryHistory(deviceKey, m, time.Time{}, time.Now(), 1)
if err != nil || len(rows) == 0 {
continue
}
if picked == nil || rows[0].TS.After(picked.TS) {
picked = &HistoryRow{TS: rows[0].TS, Value: rows[0].Value}
}
}
return picked, nil
}
// AggregateByBucket 按桶聚合
func (s *IoTDBService) AggregateByBucket(deviceKey, metric string, from, to time.Time, minutesBucket int) ([]BucketRow, error) {
if !s.IsAvailable() {
return nil, fmt.Errorf("IoTDB unavailable")
}
if minutesBucket <= 0 {
minutesBucket = 5
}
ms := s.pathSegment(metric)
sql := fmt.Sprintf("SELECT AVG(%s), MIN(%s), MAX(%s) FROM %s WHERE time >= %d AND time < %d GROUP BY ([%d, %d), %dm)",
ms, ms, ms, s.devicePath(deviceKey), from.UnixMilli(), to.UnixMilli(), from.UnixMilli(), to.UnixMilli(), minutesBucket)
resp, err := s.query(sql, 2000)
if err != nil {
slog.Warn("IoTDB aggregate 失败", "err", err)
return []BucketRow{}, nil
}
return s.parseAggregateRows(resp), nil
}
// --- 内部方法 ---
func (s *IoTDBService) markUnavailable() {
s.mu.Lock()
s.available = false
s.mu.Unlock()
}
func (s *IoTDBService) ensureTimeseries(deviceKey, metric string) error {
path := s.timeseriesPath(deviceKey, metric)
if _, ok := s.createdTS.Load(path); ok {
return nil
}
sql := fmt.Sprintf("CREATE TIMESERIES %s WITH DATATYPE=DOUBLE, ENCODING=GORILLA, COMPRESSOR=LZ4", path)
_, err := s.nonQuery(sql)
if err != nil && strings.Contains(strings.ToLower(err.Error()), "already exist") {
s.createdTS.Store(path, true)
return nil
}
if err == nil {
s.createdTS.Store(path, true)
}
return err
}
func (s *IoTDBService) nonQuery(sql string) (map[string]interface{}, error) {
return s.request("POST", "/rest/v2/nonQuery", map[string]interface{}{"sql": sql})
}
func (s *IoTDBService) query(sql string, rowLimit int) (map[string]interface{}, error) {
if rowLimit <= 0 || rowLimit > 2000 {
rowLimit = 2000
}
return s.request("POST", "/rest/v2/query", map[string]interface{}{"sql": sql, "row_limit": rowLimit})
}
func (s *IoTDBService) request(method, endpoint string, body interface{}) (map[string]interface{}, error) {
url := s.baseURL + endpoint
var bodyReader io.Reader
if body != nil {
data, err := json.Marshal(body)
if err != nil {
return nil, err
}
bodyReader = bytes.NewReader(data)
}
req, err := http.NewRequest(method, url, bodyReader)
if err != nil {
return nil, err
}
req.SetBasicAuth(s.user, s.password)
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
resp, err := s.httpClient.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
data, err := io.ReadAll(resp.Body)
if err != nil {
return nil, err
}
if resp.StatusCode >= 400 {
return nil, fmt.Errorf("IoTDB HTTP %d: %s", resp.StatusCode, string(data))
}
trimmed := strings.TrimSpace(string(data))
if trimmed == "" {
return map[string]interface{}{"code": float64(resp.StatusCode)}, nil
}
var result map[string]interface{}
if err := json.Unmarshal(data, &result); err != nil {
return nil, fmt.Errorf("IoTDB 响应解析失败: %w, body=%s", err, trimmed)
}
if code, ok := result["code"].(float64); ok && code != 0 && code != 200 {
msg, _ := result["message"].(string)
if msg == "" {
msg, _ = result["desc"].(string)
}
return nil, fmt.Errorf("IoTDB error code=%v: %s", code, msg)
}
return result, nil
}
func (s *IoTDBService) parseHistoryRows(resp map[string]interface{}) []HistoryRow {
var rows []HistoryRow
// 格式1: {timestamps: [...], values: [[...]]}
if timestamps, ok := resp["timestamps"].([]interface{}); ok {
if values, ok := resp["values"].([]interface{}); ok && len(values) > 0 {
if firstValues, ok := values[0].([]interface{}); ok {
for i, ts := range timestamps {
if i < len(firstValues) {
v := toFloat(firstValues[i])
if !isNaN(v) {
rows = append(rows, HistoryRow{TS: toTime(ts), Value: v})
}
}
}
return rows
}
}
}
// 格式2: {data: [[...]]}
if data, ok := resp["data"].([]interface{}); ok {
for _, row := range data {
if r, ok := row.([]interface{}); ok && len(r) >= 2 {
v := toFloat(r[1])
if !isNaN(v) {
rows = append(rows, HistoryRow{TS: toTime(r[0]), Value: v})
}
}
}
}
return rows
}
func (s *IoTDBService) parseAggregateRows(resp map[string]interface{}) []BucketRow {
var rows []BucketRow
if timestamps, ok := resp["timestamps"].([]interface{}); ok {
if values, ok := resp["values"].([]interface{}); ok && len(values) >= 3 {
avgs, _ := values[0].([]interface{})
mins, _ := values[1].([]interface{})
maxs, _ := values[2].([]interface{})
for i, ts := range timestamps {
row := BucketRow{Bucket: toTime(ts)}
if i < len(avgs) {
row.Avg = toFloat(avgs[i])
}
if i < len(mins) {
row.Min = toFloat(mins[i])
}
if i < len(maxs) {
row.Max = toFloat(maxs[i])
}
rows = append(rows, row)
}
return rows
}
}
if data, ok := resp["data"].([]interface{}); ok {
for _, row := range data {
if r, ok := row.([]interface{}); ok && len(r) >= 4 {
rows = append(rows, BucketRow{
Bucket: toTime(r[0]), Avg: toFloat(r[1]), Min: toFloat(r[2]), Max: toFloat(r[3]),
})
}
}
}
return rows
}
func (s *IoTDBService) extractMetrics(resp map[string]interface{}) []string {
prefix := s.database + ".telemetry."
seen := make(map[string]bool)
var metrics []string
extract := func(cell interface{}) {
str, ok := cell.(string)
if !ok || !strings.HasPrefix(str, prefix) {
return
}
// 提取最后一部分(反引号内的内容)作为 metric
parts := strings.Split(str, ".")
if len(parts) > 0 {
m := strings.Trim(parts[len(parts)-1], "`")
if m != "" && !seen[m] {
seen[m] = true
metrics = append(metrics, m)
}
}
}
if data, ok := resp["data"].([]interface{}); ok {
for _, row := range data {
if r, ok := row.([]interface{}); ok {
for _, cell := range r {
extract(cell)
}
}
}
}
return metrics
}
func (s *IoTDBService) devicePath(deviceKey string) string {
return fmt.Sprintf("%s.telemetry.%s", s.database, s.pathSegment(deviceKey))
}
func (s *IoTDBService) timeseriesPath(deviceKey, metric string) string {
return fmt.Sprintf("%s.%s", s.devicePath(deviceKey), s.pathSegment(metric))
}
func (s *IoTDBService) pathSegment(str string) string {
return "`" + strings.ReplaceAll(str, "`", "``") + "`"
}
// --- 辅助函数 ---
func toTime(v interface{}) time.Time {
switch t := v.(type) {
case float64:
return time.UnixMilli(int64(t))
case string:
if n, err := strconv.ParseInt(t, 10, 64); err == nil {
return time.UnixMilli(n)
}
if parsed, err := time.Parse(time.RFC3339, t); err == nil {
return parsed
}
}
return time.Now()
}
func toFloat(v interface{}) float64 {
switch f := v.(type) {
case float64:
return f
case string:
n, _ := strconv.ParseFloat(f, 64)
return n
}
return 0
}
func isNaN(v float64) bool { return v != v }