|
@@ -45,6 +45,10 @@ PHASES = (
|
|
|
)
|
|
)
|
|
|
ANNOTATION_LABELS = ("正常", "异常")
|
|
ANNOTATION_LABELS = ("正常", "异常")
|
|
|
|
|
|
|
|
|
|
+OIL_PRESSURE_ALARM_TYPE = "润滑油压力低"
|
|
|
|
|
+OIL_PRESSURE_DEVICE_PART = "润滑油"
|
|
|
|
|
+OIL_PRESSURE_DEVICE_POINT = "压力"
|
|
|
|
|
+
|
|
|
# PKS 全场点位:机组号 -> pks_long_sample.import_batch_id
|
|
# PKS 全场点位:机组号 -> pks_long_sample.import_batch_id
|
|
|
# 7号机=30、8号机=31、9号机=32。
|
|
# 7号机=30、8号机=31、9号机=32。
|
|
|
UNIT_BATCH = {"7": 30, "8": 31, "9": 32}
|
|
UNIT_BATCH = {"7": 30, "8": 31, "9": 32}
|
|
@@ -164,6 +168,19 @@ def _annotation_dict(row: dict[str, Any]) -> dict[str, Any]:
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
+def _alarm_dict(row: dict[str, Any]) -> dict[str, Any]:
|
|
|
|
|
+ return {
|
|
|
|
|
+ "id": int(row["id"]),
|
|
|
|
|
+ "deviceCode": row["device_code"] or "",
|
|
|
|
|
+ "devicePart": row["device_part"] or "",
|
|
|
|
|
+ "devicePoint": row["device_point"] or "",
|
|
|
|
|
+ "alarmType": row["alarm_type"] or "",
|
|
|
|
|
+ "alarmDes": row["alarm_des"] or "",
|
|
|
|
|
+ "alarmTimeStart": _time_string(row["alarm_time_start"]),
|
|
|
|
|
+ "alarmTimeEnd": _time_string(row["alarm_time_end"]),
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
def _validate_device_points(
|
|
def _validate_device_points(
|
|
|
values: list[str] | tuple[str, ...] | None,
|
|
values: list[str] | tuple[str, ...] | None,
|
|
|
*,
|
|
*,
|
|
@@ -730,6 +747,235 @@ class DataService:
|
|
|
"notice": self._source_notice(source),
|
|
"notice": self._source_notice(source),
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ def list_alarms(self, current_time: str | None = None) -> dict[str, Any]:
|
|
|
|
|
+ """列出在给定时刻命中的告警。
|
|
|
|
|
+
|
|
|
|
|
+ 命中条件为 ``alarm_time_start <= 当前时间 <= alarm_time_end``,
|
|
|
|
|
+ 即告警时间段覆盖“告警当前时间”。
|
|
|
|
|
+ """
|
|
|
|
|
+ moment = _parse_time(current_time) or datetime.now()
|
|
|
|
|
+
|
|
|
|
|
+ def database_query():
|
|
|
|
|
+ with get_connection() as connection:
|
|
|
|
|
+ with connection.cursor() as cursor:
|
|
|
|
|
+ cursor.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ SELECT id, device_code, device_part, device_point, alarm_type,
|
|
|
|
|
+ alarm_des, alarm_time_start, alarm_time_end
|
|
|
|
|
+ FROM compressor_alarm
|
|
|
|
|
+ WHERE alarm_time_start <= %s AND alarm_time_end >= %s
|
|
|
|
|
+ ORDER BY alarm_time_start DESC, id DESC
|
|
|
|
|
+ """,
|
|
|
|
|
+ (moment, moment),
|
|
|
|
|
+ )
|
|
|
|
|
+ rows = cursor.fetchall()
|
|
|
|
|
+ return [_alarm_dict(row) for row in rows]
|
|
|
|
|
+
|
|
|
|
|
+ result, source = self._run_with_fallback(database_query, lambda: [])
|
|
|
|
|
+ return {
|
|
|
|
|
+ "source": source,
|
|
|
|
|
+ "currentTime": _time_string(moment),
|
|
|
|
|
+ "alarms": result,
|
|
|
|
|
+ "notice": self._source_notice(source),
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ def scan_alarm(self, payload: dict[str, Any]) -> dict[str, Any]:
|
|
|
|
|
+ """执行告警算法并按原有协议写入 compressor_alarm。"""
|
|
|
|
|
+ device_code = str(payload.get("device_code") or "").strip()
|
|
|
|
|
+ alarm_type = str(payload.get("alarm_type") or "").strip()
|
|
|
|
|
+ start = _parse_time(payload.get("forecast_time"))
|
|
|
|
|
+ try:
|
|
|
|
|
+ hours = float(payload.get("hours") or 0)
|
|
|
|
|
+ except (TypeError, ValueError):
|
|
|
|
|
+ raise ValueError("告警时长必须是数字") from None
|
|
|
|
|
+ if not device_code:
|
|
|
|
|
+ raise ValueError("压缩机不能为空")
|
|
|
|
|
+ if not alarm_type:
|
|
|
|
|
+ raise ValueError("算法不能为空")
|
|
|
|
|
+ if start is None:
|
|
|
|
|
+ raise ValueError("预警时间不能为空")
|
|
|
|
|
+ if hours <= 0:
|
|
|
|
|
+ raise ValueError("告警时长必须大于 0")
|
|
|
|
|
+
|
|
|
|
|
+ if alarm_type == OIL_PRESSURE_ALARM_TYPE:
|
|
|
|
|
+ hours = 24.0
|
|
|
|
|
+ analysis = self._analyze_oil_pressure_alarm(device_code, start)
|
|
|
|
|
+ if analysis is None:
|
|
|
|
|
+ return {
|
|
|
|
|
+ "source": "database",
|
|
|
|
|
+ "notice": None,
|
|
|
|
|
+ "action": "no_alarm",
|
|
|
|
|
+ "id": 0,
|
|
|
|
|
+ }
|
|
|
|
|
+ device_part = OIL_PRESSURE_DEVICE_PART
|
|
|
|
|
+ device_point = OIL_PRESSURE_DEVICE_POINT
|
|
|
|
|
+ alarm_des = analysis
|
|
|
|
|
+ status = 1
|
|
|
|
|
+ else:
|
|
|
|
|
+ # Keep the existing endpoint contract for algorithms not yet implemented.
|
|
|
|
|
+ device_part = "测试组件"
|
|
|
|
|
+ device_point = "测试组件"
|
|
|
|
|
+ alarm_des = "测试数据"
|
|
|
|
|
+ status = 0
|
|
|
|
|
+ end = start + timedelta(hours=hours)
|
|
|
|
|
+
|
|
|
|
|
+ def database_query():
|
|
|
|
|
+ with get_connection() as connection:
|
|
|
|
|
+ with connection.cursor() as cursor:
|
|
|
|
|
+ cursor.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ SELECT id FROM compressor_alarm
|
|
|
|
|
+ WHERE device_code = %s AND device_part = %s
|
|
|
|
|
+ AND device_point = %s AND alarm_type = %s
|
|
|
|
|
+ AND alarm_time_start <= %s AND alarm_time_end >= %s
|
|
|
|
|
+ ORDER BY id ASC
|
|
|
|
|
+ LIMIT 1
|
|
|
|
|
+ """,
|
|
|
|
|
+ (device_code, device_part, device_point, alarm_type, start, start),
|
|
|
|
|
+ )
|
|
|
|
|
+ row = cursor.fetchone()
|
|
|
|
|
+ if row is not None:
|
|
|
|
|
+ cursor.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ UPDATE compressor_alarm
|
|
|
|
|
+ SET alarm_time_end = %s, alarm_des = %s, status = %s
|
|
|
|
|
+ WHERE id = %s
|
|
|
|
|
+ """,
|
|
|
|
|
+ (end, alarm_des, status, row["id"]),
|
|
|
|
|
+ )
|
|
|
|
|
+ return {"action": "updated", "id": int(row["id"])}
|
|
|
|
|
+ cursor.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ INSERT INTO compressor_alarm
|
|
|
|
|
+ (device_code, device_part, device_point, alarm_type,
|
|
|
|
|
+ alarm_des, alarm_time_start, alarm_time_end, status)
|
|
|
|
|
+ VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
|
|
|
|
|
+ """,
|
|
|
|
|
+ (device_code, device_part, device_point, alarm_type, alarm_des, start, end, status),
|
|
|
|
|
+ )
|
|
|
|
|
+ return {"action": "inserted", "id": int(cursor.lastrowid)}
|
|
|
|
|
+
|
|
|
|
|
+ result, source = self._run_with_fallback(
|
|
|
|
|
+ database_query,
|
|
|
|
|
+ lambda: {"action": "demo", "id": 0},
|
|
|
|
|
+ )
|
|
|
|
|
+ return {"source": source, "notice": self._source_notice(source), **result}
|
|
|
|
|
+
|
|
|
|
|
+ @staticmethod
|
|
|
|
|
+ def _alarm_limit(item: dict[str, Any], alarm_type: str) -> float | None:
|
|
|
|
|
+ for index in range(1, 5):
|
|
|
|
|
+ if str(item.get(f"AlarmType{index}") or "").strip() == alarm_type:
|
|
|
|
|
+ return _safe_float(item.get(f"AlarmLimit{index}"))
|
|
|
|
|
+ return None
|
|
|
|
|
+
|
|
|
|
|
+ def _analyze_oil_pressure_alarm(self, device_code: str, end: datetime) -> str | None:
|
|
|
|
|
+ """Analyze YSJ_5 in the 24 hours before ``end``.
|
|
|
|
|
+
|
|
|
|
|
+ A low-pressure alarm requires either a configured low/low-low threshold
|
|
|
|
|
+ breach or a sustained decline of at least 10% from the 24-hour median.
|
|
|
|
|
+ The trend rule requires three consecutive non-rising valid hours or a
|
|
|
|
|
+ negative six-hour linear trend. Direct low-low threshold breaches remain
|
|
|
|
|
+ immediately reportable.
|
|
|
|
|
+ """
|
|
|
|
|
+ unit = device_code.rstrip("#").strip()
|
|
|
|
|
+ if unit not in UNIT_BATCH:
|
|
|
|
|
+ raise ValueError("压缩机必须是 7#、8# 或 9#")
|
|
|
|
|
+ start = end - timedelta(hours=24)
|
|
|
|
|
+
|
|
|
|
|
+ with get_connection() as connection:
|
|
|
|
|
+ with connection.cursor() as cursor:
|
|
|
|
|
+ cursor.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ SELECT AlarmType1, AlarmType2, AlarmType3, AlarmType4,
|
|
|
|
|
+ AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4
|
|
|
|
|
+ FROM site_point
|
|
|
|
|
+ WHERE ItemName = %s
|
|
|
|
|
+ """,
|
|
|
|
|
+ (f"YSJ{unit}_5",),
|
|
|
|
|
+ )
|
|
|
|
|
+ config = cursor.fetchone()
|
|
|
|
|
+ if config is None:
|
|
|
|
|
+ raise ValueError(f"未找到 {device_code} 的润滑油压力报警配置")
|
|
|
|
|
+ low = self._alarm_limit(config, "PVLow")
|
|
|
|
|
+ low_low = self._alarm_limit(config, "PVLowLow")
|
|
|
|
|
+ if low is None or low_low is None:
|
|
|
|
|
+ raise ValueError(f"{device_code} 缺少润滑油压力低压报警阈值")
|
|
|
|
|
+
|
|
|
|
|
+ cursor.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ SELECT FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start,
|
|
|
|
|
+ COUNT(*) AS samples,
|
|
|
|
|
+ SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples,
|
|
|
|
|
+ AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_avg
|
|
|
|
|
+ FROM pks_long_sample
|
|
|
|
|
+ WHERE device_code = %s
|
|
|
|
|
+ AND sample_time >= %s AND sample_time < %s
|
|
|
|
|
+ GROUP BY FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600)
|
|
|
|
|
+ ORDER BY hour_start
|
|
|
|
|
+ """,
|
|
|
|
|
+ (device_code, start, end),
|
|
|
|
|
+ )
|
|
|
|
|
+ rows = []
|
|
|
|
|
+ for row in cursor.fetchall():
|
|
|
|
|
+ running = int(row["running_samples"] or 0)
|
|
|
|
|
+ value = _safe_float(row["oil_avg"])
|
|
|
|
|
+ if running / 720 >= 0.80 and value is not None:
|
|
|
|
|
+ rows.append({"time": row["hour_start"], "value": value, "running": running})
|
|
|
|
|
+
|
|
|
|
|
+ if not rows:
|
|
|
|
|
+ return None
|
|
|
|
|
+ values = [row["value"] for row in rows]
|
|
|
|
|
+ baseline = float(np.median(values))
|
|
|
|
|
+ current = rows[-1]["value"]
|
|
|
|
|
+ relative_drop = (baseline - current) / baseline if baseline else 0.0
|
|
|
|
|
+ low_low_rows = [row for row in rows if row["value"] <= low_low]
|
|
|
|
|
+ low_rows = [row for row in rows if row["value"] <= low]
|
|
|
|
|
+ consecutive_low = 0
|
|
|
|
|
+ for row in reversed(rows):
|
|
|
|
|
+ if row["value"] <= low:
|
|
|
|
|
+ consecutive_low += 1
|
|
|
|
|
+ else:
|
|
|
|
|
+ break
|
|
|
|
|
+ consecutive_low_low = 0
|
|
|
|
|
+ for row in reversed(rows):
|
|
|
|
|
+ if row["value"] <= low_low:
|
|
|
|
|
+ consecutive_low_low += 1
|
|
|
|
|
+ else:
|
|
|
|
|
+ break
|
|
|
|
|
+ consecutive_decline = 1
|
|
|
|
|
+ for index in range(len(rows) - 1, 0, -1):
|
|
|
|
|
+ if rows[index]["time"] - rows[index - 1]["time"] != timedelta(hours=1):
|
|
|
|
|
+ break
|
|
|
|
|
+ if rows[index]["value"] > rows[index - 1]["value"]:
|
|
|
|
|
+ break
|
|
|
|
|
+ consecutive_decline += 1
|
|
|
|
|
+ trend_values = [row["value"] for row in rows[-6:]]
|
|
|
|
|
+ trend = float(np.polyfit(range(len(trend_values)), trend_values, 1)[0]) if len(trend_values) >= 3 else None
|
|
|
|
|
+ downward = consecutive_decline >= 3 or (trend is not None and trend < -0.0005)
|
|
|
|
|
+ evolution_alarm = relative_drop >= 0.10 and downward
|
|
|
|
|
+
|
|
|
|
|
+ if not low_low_rows and not low_rows and not evolution_alarm:
|
|
|
|
|
+ return None
|
|
|
|
|
+ reasons = [
|
|
|
|
|
+ f"前24小时有效运行数据{len(rows)}小时",
|
|
|
|
|
+ f"当前小时油压{current:.3f},24小时基线{baseline:.3f},相对下降{relative_drop:.1%}",
|
|
|
|
|
+ ]
|
|
|
|
|
+ if consecutive_low_low:
|
|
|
|
|
+ reasons.append(f"连续{consecutive_low_low}小时低于低低报警阈值{low_low:.3f}")
|
|
|
|
|
+ elif low_low_rows:
|
|
|
|
|
+ reasons.append(f"低低报警区间出现{len(low_low_rows)}小时")
|
|
|
|
|
+ elif consecutive_low:
|
|
|
|
|
+ reasons.append(f"连续{consecutive_low}小时低于低报警阈值{low:.3f}")
|
|
|
|
|
+ else:
|
|
|
|
|
+ reasons.append(f"低报警区间出现{len(low_rows)}小时")
|
|
|
|
|
+ if evolution_alarm:
|
|
|
|
|
+ reasons.append(
|
|
|
|
|
+ f"相对24小时基线下降至少10%,连续下降{consecutive_decline}小时,"
|
|
|
|
|
+ f"近6小时斜率{trend:.6f}" if trend is not None
|
|
|
|
|
+ else f"相对24小时基线下降至少10%,连续下降{consecutive_decline}小时"
|
|
|
|
|
+ )
|
|
|
|
|
+ return ";".join(reasons)
|
|
|
|
|
+
|
|
|
@staticmethod
|
|
@staticmethod
|
|
|
def _group_time_points(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
|
def _group_time_points(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
|
|
grouped: OrderedDict[Any, dict[str, Any]] = OrderedDict()
|
|
grouped: OrderedDict[Any, dict[str, Any]] = OrderedDict()
|