|
|
@@ -1,6 +1,7 @@
|
|
|
from __future__ import annotations
|
|
|
|
|
|
import re
|
|
|
+import json
|
|
|
from bisect import bisect_left
|
|
|
from collections import OrderedDict
|
|
|
from datetime import datetime, timedelta
|
|
|
@@ -46,8 +47,12 @@ PHASES = (
|
|
|
ANNOTATION_LABELS = ("正常", "异常")
|
|
|
|
|
|
OIL_PRESSURE_ALARM_TYPE = "润滑油压力低"
|
|
|
+CRUCIFORM_FAULT_ALARM_TYPE = "十字头故障"
|
|
|
+VALVE_FAULT_ALARM_TYPE = "气阀故障"
|
|
|
OIL_PRESSURE_DEVICE_PART = "润滑油"
|
|
|
OIL_PRESSURE_DEVICE_POINT = "压力"
|
|
|
+VIBRATION_DEVICE_POINT = "振动"
|
|
|
+OIL_PRESSURE_SCAN_DAYS = 14
|
|
|
|
|
|
# PKS 全场点位:机组号 -> pks_long_sample.import_batch_id
|
|
|
# 7号机=30、8号机=31、9号机=32。
|
|
|
@@ -132,6 +137,24 @@ def _stored_first_cycle(samples: np.ndarray) -> DetectedCycle | None:
|
|
|
)
|
|
|
|
|
|
|
|
|
+def _sample_first_cycle_360(samples: np.ndarray) -> list[tuple[int, float]]:
|
|
|
+ """Resample the stored first cycle into 360 integer-angle points.
|
|
|
+
|
|
|
+ Each integer angle 0..359 keeps the first raw sample whose crank angle is at
|
|
|
+ or after that integer value (the "取第一个值" rule from showPV).
|
|
|
+ """
|
|
|
+ if samples.ndim != 2 or samples.shape[1] < 2 or len(samples) < 2:
|
|
|
+ return []
|
|
|
+ count = len(samples)
|
|
|
+ angle = np.linspace(0.0, 360.0, count, endpoint=False)
|
|
|
+ indices = np.searchsorted(angle, np.arange(360, dtype=float), side="left")
|
|
|
+ indices = np.clip(indices, 0, count - 1)
|
|
|
+ return [
|
|
|
+ (int(deg), float(samples[int(idx), 1]))
|
|
|
+ for deg, idx in enumerate(indices)
|
|
|
+ ]
|
|
|
+
|
|
|
+
|
|
|
def _detect_source_cycles(
|
|
|
samples: np.ndarray,
|
|
|
single_cycle: bool,
|
|
|
@@ -176,6 +199,11 @@ def _alarm_dict(row: dict[str, Any]) -> dict[str, Any]:
|
|
|
"devicePoint": row["device_point"] or "",
|
|
|
"alarmType": row["alarm_type"] or "",
|
|
|
"alarmDes": row["alarm_des"] or "",
|
|
|
+ "status": int(row.get("status") or 0),
|
|
|
+ "alarmLevel": row.get("alarm_level"),
|
|
|
+ "scanStartTime": _time_string(row["scan_start_time"]) if row.get("scan_start_time") else "",
|
|
|
+ "scanEndTime": _time_string(row["scan_end_time"]) if row.get("scan_end_time") else "",
|
|
|
+ "alarmInfoJson": row.get("alarm_info_json") or "",
|
|
|
"alarmTimeStart": _time_string(row["alarm_time_start"]),
|
|
|
"alarmTimeEnd": _time_string(row["alarm_time_end"]),
|
|
|
}
|
|
|
@@ -761,7 +789,8 @@ class DataService:
|
|
|
cursor.execute(
|
|
|
"""
|
|
|
SELECT id, device_code, device_part, device_point, alarm_type,
|
|
|
- alarm_des, alarm_time_start, alarm_time_end
|
|
|
+ alarm_des, alarm_time_start, alarm_time_end, status,
|
|
|
+ alarm_level, scan_start_time, scan_end_time, alarm_info_json
|
|
|
FROM compressor_alarm
|
|
|
WHERE alarm_time_start <= %s AND alarm_time_end >= %s
|
|
|
ORDER BY alarm_time_start DESC, id DESC
|
|
|
@@ -796,9 +825,12 @@ class DataService:
|
|
|
raise ValueError("预警时间不能为空")
|
|
|
if hours <= 0:
|
|
|
raise ValueError("告警时长必须大于 0")
|
|
|
+ if alarm_type == VALVE_FAULT_ALARM_TYPE:
|
|
|
+ return self._scan_valve_fault(device_code, start, hours)
|
|
|
|
|
|
if alarm_type == OIL_PRESSURE_ALARM_TYPE:
|
|
|
- hours = 24.0
|
|
|
+ # `hours` controls the alarm validity window, not the analysis window.
|
|
|
+ # Oil pressure analysis always uses the previous 14 days below.
|
|
|
analysis = self._analyze_oil_pressure_alarm(device_code, start)
|
|
|
if analysis is None:
|
|
|
return {
|
|
|
@@ -807,16 +839,41 @@ class DataService:
|
|
|
"action": "no_alarm",
|
|
|
"id": 0,
|
|
|
}
|
|
|
- device_part = OIL_PRESSURE_DEVICE_PART
|
|
|
+ device_part = analysis["device_part"]
|
|
|
device_point = OIL_PRESSURE_DEVICE_POINT
|
|
|
- alarm_des = analysis
|
|
|
+ alarm_des = analysis["basis"]
|
|
|
+ status = 1
|
|
|
+ alarm_level = analysis["stage"]
|
|
|
+ scan_start = analysis["scan_start"]
|
|
|
+ scan_end = analysis["scan_end"]
|
|
|
+ alarm_info_json = json.dumps(analysis["info"], ensure_ascii=False, separators=(",", ":"))
|
|
|
+ elif alarm_type == CRUCIFORM_FAULT_ALARM_TYPE:
|
|
|
+ analysis = self._analyze_cruciform_vibration_alarm(device_code, start)
|
|
|
+ if analysis is None:
|
|
|
+ return {
|
|
|
+ "source": "database",
|
|
|
+ "notice": None,
|
|
|
+ "action": "no_alarm",
|
|
|
+ "id": 0,
|
|
|
+ }
|
|
|
+ device_part = analysis["device_part"]
|
|
|
+ device_point = VIBRATION_DEVICE_POINT
|
|
|
+ alarm_des = analysis["basis"]
|
|
|
status = 1
|
|
|
+ alarm_level = analysis["stage"]
|
|
|
+ scan_start = analysis["scan_start"]
|
|
|
+ scan_end = analysis["scan_end"]
|
|
|
+ alarm_info_json = json.dumps(analysis["info"], ensure_ascii=False, separators=(",", ":"))
|
|
|
else:
|
|
|
# Keep the existing endpoint contract for algorithms not yet implemented.
|
|
|
device_part = "测试组件"
|
|
|
device_point = "测试组件"
|
|
|
alarm_des = "测试数据"
|
|
|
status = 0
|
|
|
+ alarm_level = None
|
|
|
+ scan_start = None
|
|
|
+ scan_end = None
|
|
|
+ alarm_info_json = None
|
|
|
end = start + timedelta(hours=hours)
|
|
|
|
|
|
def database_query():
|
|
|
@@ -838,20 +895,24 @@ class DataService:
|
|
|
cursor.execute(
|
|
|
"""
|
|
|
UPDATE compressor_alarm
|
|
|
- SET alarm_time_end = %s, alarm_des = %s, status = %s
|
|
|
+ SET alarm_time_end = %s, alarm_des = %s, status = %s,
|
|
|
+ alarm_level = %s, scan_start_time = %s, scan_end_time = %s,
|
|
|
+ alarm_info_json = %s
|
|
|
WHERE id = %s
|
|
|
""",
|
|
|
- (end, alarm_des, status, row["id"]),
|
|
|
+ (end, alarm_des, status, alarm_level, scan_start, scan_end, alarm_info_json, 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)
|
|
|
+ alarm_des, alarm_time_start, alarm_time_end, status,
|
|
|
+ alarm_level, scan_start_time, scan_end_time, alarm_info_json)
|
|
|
+ VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
|
""",
|
|
|
- (device_code, device_part, device_point, alarm_type, alarm_des, start, end, status),
|
|
|
+ (device_code, device_part, device_point, alarm_type, alarm_des, start, end, status,
|
|
|
+ alarm_level, scan_start, scan_end, alarm_info_json),
|
|
|
)
|
|
|
return {"action": "inserted", "id": int(cursor.lastrowid)}
|
|
|
|
|
|
@@ -861,6 +922,226 @@ class DataService:
|
|
|
)
|
|
|
return {"source": source, "notice": self._source_notice(source), **result}
|
|
|
|
|
|
+ @staticmethod
|
|
|
+ def _valve_segment_values(values: list[float]) -> list[float]:
|
|
|
+ if not values:
|
|
|
+ return []
|
|
|
+ raw_average = sum(values) / len(values)
|
|
|
+ return [value for value in values if value >= raw_average * 0.5]
|
|
|
+
|
|
|
+ @classmethod
|
|
|
+ def _valve_segment_average(cls, values: list[float]) -> float | None:
|
|
|
+ filtered = cls._valve_segment_values(values)
|
|
|
+ return sum(filtered) / len(filtered) if filtered else None
|
|
|
+
|
|
|
+ @staticmethod
|
|
|
+ def _split_valve_segments(rows: list[dict[str, Any]]) -> list[list[dict[str, Any]]]:
|
|
|
+ segments: list[list[dict[str, Any]]] = []
|
|
|
+ current: list[dict[str, Any]] = []
|
|
|
+ for row in rows:
|
|
|
+ if current and row["sample_time"] - current[-1]["sample_time"] > timedelta(days=1):
|
|
|
+ segments.append(current)
|
|
|
+ current = []
|
|
|
+ current.append(row)
|
|
|
+ if current:
|
|
|
+ segments.append(current)
|
|
|
+ return segments
|
|
|
+
|
|
|
+ @staticmethod
|
|
|
+ def _previous_valve_segment(
|
|
|
+ cursor: Any, part: str, point: str, before: datetime, current_first: datetime
|
|
|
+ ) -> tuple[datetime, datetime] | None:
|
|
|
+ """Find the nearest preceding running segment lasting at least 24 hours."""
|
|
|
+ boundary = before
|
|
|
+ boundary_id = 1 << 63
|
|
|
+ newer = None
|
|
|
+ previous_end = None
|
|
|
+ previous_start = None
|
|
|
+ while True:
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ SELECT id, sample_time FROM wave_file
|
|
|
+ WHERE device_part = %s AND device_point = %s
|
|
|
+ AND measurement_type = %s AND rpm > 0
|
|
|
+ AND (sample_time < %s OR (sample_time = %s AND id < %s))
|
|
|
+ ORDER BY sample_time DESC, id DESC LIMIT 1000
|
|
|
+ """,
|
|
|
+ (part, point, "压力", boundary, boundary, boundary_id),
|
|
|
+ )
|
|
|
+ rows = cursor.fetchall()
|
|
|
+ if not rows:
|
|
|
+ break
|
|
|
+ for row in rows:
|
|
|
+ stamp = row["sample_time"]
|
|
|
+ if newer is None:
|
|
|
+ previous_end = stamp
|
|
|
+ previous_start = stamp
|
|
|
+ elif newer - stamp > timedelta(days=1):
|
|
|
+ if previous_end - previous_start >= timedelta(hours=24):
|
|
|
+ return previous_start, previous_end
|
|
|
+ previous_end = stamp
|
|
|
+ previous_start = stamp
|
|
|
+ else:
|
|
|
+ previous_start = stamp
|
|
|
+ newer = stamp
|
|
|
+ boundary = rows[-1]["sample_time"]
|
|
|
+ boundary_id = int(rows[-1]["id"])
|
|
|
+ if len(rows) < 1000:
|
|
|
+ break
|
|
|
+ if (
|
|
|
+ previous_start is not None
|
|
|
+ and previous_end is not None
|
|
|
+ and previous_end - previous_start >= timedelta(hours=24)
|
|
|
+ ):
|
|
|
+ return previous_start, previous_end
|
|
|
+ return None
|
|
|
+
|
|
|
+ def _scan_valve_fault(self, device_code: str, forecast_time: datetime, hours: float) -> dict[str, Any]:
|
|
|
+ unit_name = device_code if device_code.endswith("号机组") else f"{device_code.rstrip('#')}号机组"
|
|
|
+ scan_start = forecast_time - timedelta(days=14)
|
|
|
+ scan_end = forecast_time + timedelta(days=1)
|
|
|
+
|
|
|
+ def database_query():
|
|
|
+ with get_connection() as connection:
|
|
|
+ with connection.cursor() as cursor:
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ SELECT f.device_part, f.device_point, f.sample_time, pa.pressure
|
|
|
+ FROM wave_file f
|
|
|
+ LEFT JOIN statistic_pressure_angle pa
|
|
|
+ ON f.id = pa.wave_file_id AND pa.angle = 230
|
|
|
+ WHERE SUBSTRING(f.device_part, 1, 4) = %s
|
|
|
+ AND f.rpm > 0
|
|
|
+ AND f.measurement_type = %s
|
|
|
+ AND f.sample_time >= %s AND f.sample_time < %s
|
|
|
+ ORDER BY f.device_part, f.device_point, f.sample_time
|
|
|
+ """,
|
|
|
+ (unit_name, "压力", scan_start, scan_end),
|
|
|
+ )
|
|
|
+ current_rows = cursor.fetchall()
|
|
|
+ grouped: dict[tuple[str, str], list[dict[str, Any]]] = {}
|
|
|
+ for row in current_rows:
|
|
|
+ grouped.setdefault((row["device_part"], row["device_point"]), []).append(row)
|
|
|
+ results = []
|
|
|
+ for (part, point), running_rows in grouped.items():
|
|
|
+ segments = self._split_valve_segments(running_rows)
|
|
|
+ if not segments:
|
|
|
+ continue
|
|
|
+ current = segments[-1]
|
|
|
+ current_scan_start = current[0]["sample_time"]
|
|
|
+ current_values = self._valve_segment_values(
|
|
|
+ [float(row["pressure"]) for row in current if row["pressure"] is not None]
|
|
|
+ )
|
|
|
+ current_avg = sum(current_values) / len(current_values) if current_values else None
|
|
|
+ if current_avg is None or current_avg == 0:
|
|
|
+ continue
|
|
|
+ previous_range = None
|
|
|
+ previous_values: list[float] = []
|
|
|
+ eligible_previous = [
|
|
|
+ segment for segment in segments[:-1]
|
|
|
+ if segment[-1]["sample_time"] - segment[0]["sample_time"] >= timedelta(hours=24)
|
|
|
+ ]
|
|
|
+ if eligible_previous:
|
|
|
+ previous = eligible_previous[-1]
|
|
|
+ previous_start = previous[0]["sample_time"]
|
|
|
+ previous_end = previous[-1]["sample_time"]
|
|
|
+ previous_values = self._valve_segment_values(
|
|
|
+ [float(row["pressure"]) for row in previous if row["pressure"] is not None]
|
|
|
+ )
|
|
|
+ previous_range = (previous_start, previous_end)
|
|
|
+ else:
|
|
|
+ previous_range = self._previous_valve_segment(
|
|
|
+ cursor, part, point, current_scan_start, current_scan_start
|
|
|
+ )
|
|
|
+ previous_start = previous_end = None
|
|
|
+ previous_avg = None
|
|
|
+ if previous_range is not None:
|
|
|
+ previous_start, previous_end = previous_range
|
|
|
+ if not previous_values:
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ SELECT pa.pressure FROM wave_file f
|
|
|
+ INNER JOIN statistic_pressure_angle pa
|
|
|
+ ON f.id = pa.wave_file_id AND pa.angle = 230
|
|
|
+ WHERE f.device_part = %s AND f.device_point = %s
|
|
|
+ AND f.measurement_type = %s AND f.rpm > 0
|
|
|
+ AND f.sample_time >= %s AND f.sample_time <= %s
|
|
|
+ """,
|
|
|
+ (part, point, "压力", previous_start, previous_end),
|
|
|
+ )
|
|
|
+ previous_values = self._valve_segment_values(
|
|
|
+ [float(row["pressure"]) for row in cursor.fetchall() if row["pressure"] is not None]
|
|
|
+ )
|
|
|
+ previous_avg = sum(previous_values) / len(previous_values) if previous_values else None
|
|
|
+ else:
|
|
|
+ previous_values = []
|
|
|
+ is_fault = previous_avg is not None and previous_avg > current_avg * 1.05
|
|
|
+ if not is_fault:
|
|
|
+ continue
|
|
|
+ info = json.dumps(
|
|
|
+ {
|
|
|
+ "当前段平均值": current_avg,
|
|
|
+ "上一段平均值": previous_avg,
|
|
|
+ "当前段最高值": max(current_values) if current_values else None,
|
|
|
+ "当前段最低值": min(current_values) if current_values else None,
|
|
|
+ "上一段最高值": max(previous_values) if previous_values else None,
|
|
|
+ "上一段最低值": min(previous_values) if previous_values else None,
|
|
|
+ "上一段开始时间": _time_string(previous_start) if previous_start else None,
|
|
|
+ "上一段结束时间": _time_string(previous_end) if previous_end else None,
|
|
|
+ },
|
|
|
+ ensure_ascii=False,
|
|
|
+ )
|
|
|
+ description = (
|
|
|
+ f"230°压力上一段均值{previous_avg:.4f}高于当前段均值"
|
|
|
+ f"{current_avg:.4f}的105%"
|
|
|
+ )
|
|
|
+ status = 1
|
|
|
+ alarm_level = "1"
|
|
|
+ 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
|
|
|
+ LIMIT 1
|
|
|
+ """,
|
|
|
+ (device_code, part, point, VALVE_FAULT_ALARM_TYPE, forecast_time),
|
|
|
+ )
|
|
|
+ existing = cursor.fetchone()
|
|
|
+ if existing:
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ UPDATE compressor_alarm SET alarm_time_end = %s, alarm_des = %s,
|
|
|
+ alarm_level = %s, alarm_info_json = %s, status = %s,
|
|
|
+ scan_start_time = %s, scan_end_time = %s
|
|
|
+ WHERE id = %s
|
|
|
+ """,
|
|
|
+ (forecast_time + timedelta(hours=hours), description, alarm_level, info,
|
|
|
+ status, current_scan_start, forecast_time, existing["id"]),
|
|
|
+ )
|
|
|
+ results.append(("updated", int(existing["id"])))
|
|
|
+ else:
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ INSERT INTO compressor_alarm
|
|
|
+ (device_code, device_part, device_point, alarm_type, alarm_des,
|
|
|
+ alarm_level, alarm_info_json, alarm_time_start, alarm_time_end,
|
|
|
+ scan_start_time, scan_end_time, status)
|
|
|
+ VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
|
+ """,
|
|
|
+ (device_code, part, point, VALVE_FAULT_ALARM_TYPE, description, alarm_level, info,
|
|
|
+ forecast_time, forecast_time + timedelta(hours=hours), current_scan_start, forecast_time, status),
|
|
|
+ )
|
|
|
+ results.append(("inserted", int(cursor.lastrowid)))
|
|
|
+ return {
|
|
|
+ "action": results[0][0] if results else "no_alarm",
|
|
|
+ "id": results[0][1] if results else 0,
|
|
|
+ "inserted": sum(action == "inserted" for action, _ in results),
|
|
|
+ "updated": sum(action == "updated" for action, _ in results),
|
|
|
+ }
|
|
|
+
|
|
|
+ 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):
|
|
|
@@ -868,25 +1149,18 @@ class DataService:
|
|
|
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.
|
|
|
- """
|
|
|
+ def _analyze_oil_pressure_alarm(self, device_code: str, end: datetime) -> dict[str, Any] | None:
|
|
|
+ """Analyze oil-pressure stages in the 14 days before ``end``."""
|
|
|
unit = device_code.rstrip("#").strip()
|
|
|
- if unit not in UNIT_BATCH:
|
|
|
- raise ValueError("压缩机必须是 7#、8# 或 9#")
|
|
|
- start = end - timedelta(hours=24)
|
|
|
+ if not device_code.endswith("#") or not device_code[:-1].isdigit():
|
|
|
+ raise ValueError("压缩机编号格式必须为数字+#")
|
|
|
+ start = end - timedelta(days=OIL_PRESSURE_SCAN_DAYS)
|
|
|
|
|
|
with get_connection() as connection:
|
|
|
with connection.cursor() as cursor:
|
|
|
cursor.execute(
|
|
|
"""
|
|
|
- SELECT AlarmType1, AlarmType2, AlarmType3, AlarmType4,
|
|
|
+ SELECT ItemName, ItemDescription, AlarmType1, AlarmType2, AlarmType3, AlarmType4,
|
|
|
AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4
|
|
|
FROM site_point
|
|
|
WHERE ItemName = %s
|
|
|
@@ -901,12 +1175,23 @@ class DataService:
|
|
|
if low is None or low_low is None:
|
|
|
raise ValueError(f"{device_code} 缺少润滑油压力低压报警阈值")
|
|
|
|
|
|
+ point_descriptions = {}
|
|
|
+ for point in (5, 41):
|
|
|
+ cursor.execute(
|
|
|
+ "SELECT ItemName, ItemDescription FROM site_point WHERE ItemName = %s",
|
|
|
+ (f"YSJ{unit}_{point}",),
|
|
|
+ )
|
|
|
+ item = cursor.fetchone()
|
|
|
+ if item:
|
|
|
+ point_descriptions[f"YSJ_{point}"] = item.get("ItemDescription") or f"YSJ_{point}"
|
|
|
+
|
|
|
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
|
|
|
+ AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_avg,
|
|
|
+ GROUP_CONCAT(DISTINCT import_batch_id ORDER BY import_batch_id) AS batch_ids
|
|
|
FROM pks_long_sample
|
|
|
WHERE device_code = %s
|
|
|
AND sample_time >= %s AND sample_time < %s
|
|
|
@@ -920,61 +1205,455 @@ class DataService:
|
|
|
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})
|
|
|
+ rows.append({
|
|
|
+ "time": row["hour_start"], "value": value, "running": running,
|
|
|
+ "batch_ids": [int(value) for value in str(row["batch_ids"] or "").split(",") if value],
|
|
|
+ })
|
|
|
|
|
|
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
|
|
|
+ baseline = float(np.median([row["value"] for row in rows]))
|
|
|
+ candidates = []
|
|
|
+ history = []
|
|
|
+ for row in rows:
|
|
|
+ recent = history[-5:] + [row]
|
|
|
+ trend_values = [item["value"] for item in recent]
|
|
|
+ trend = float(np.polyfit(range(len(trend_values)), trend_values, 1)[0]) if len(trend_values) >= 3 else None
|
|
|
+ decline = 1
|
|
|
+ newer = row
|
|
|
+ for older in reversed(history):
|
|
|
+ if newer["time"] - older["time"] != timedelta(hours=1):
|
|
|
+ break
|
|
|
+ if newer["value"] > older["value"]:
|
|
|
+ break
|
|
|
+ decline += 1
|
|
|
+ newer = older
|
|
|
+ current = row["value"]
|
|
|
+ deviation = (baseline - current) / baseline if baseline else 0.0
|
|
|
+ downward = decline >= 3 or (trend is not None and trend < -0.0005)
|
|
|
+ if current <= low_low or (deviation >= 0.20 and downward):
|
|
|
+ phase = "严重异常"
|
|
|
+ elif current <= low or (deviation >= 0.10 and downward):
|
|
|
+ phase = "异常"
|
|
|
+ elif deviation >= 0.05 and downward:
|
|
|
+ phase = "轻微"
|
|
|
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
|
|
|
+ phase = "正常"
|
|
|
+ candidates.append({"row": recent[-1], "phase": phase, "trend": trend, "decline": decline, "deviation": deviation})
|
|
|
+ history.append(recent[-1])
|
|
|
|
|
|
- if not low_low_rows and not low_rows and not evolution_alarm:
|
|
|
+ abnormal = [item for item in candidates if item["phase"] != "正常"]
|
|
|
+ if not abnormal:
|
|
|
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)
|
|
|
+ stage_runs = []
|
|
|
+ current_run = []
|
|
|
+ for item in candidates:
|
|
|
+ if current_run and (
|
|
|
+ item["phase"] != current_run[-1]["phase"]
|
|
|
+ or item["row"]["time"] - current_run[-1]["row"]["time"] != timedelta(hours=1)
|
|
|
+ ):
|
|
|
+ if current_run[0]["phase"] != "正常":
|
|
|
+ stage_runs.append(current_run)
|
|
|
+ current_run = []
|
|
|
+ current_run.append(item)
|
|
|
+ if current_run and current_run[0]["phase"] != "正常":
|
|
|
+ stage_runs.append(current_run)
|
|
|
+ stage_rows = max(
|
|
|
+ stage_runs,
|
|
|
+ key=lambda run: (
|
|
|
+ {"轻微": 1, "异常": 2, "严重异常": 3}[run[0]["phase"]],
|
|
|
+ len(run),
|
|
|
+ ),
|
|
|
+ )
|
|
|
+ stage = stage_rows[0]["phase"]
|
|
|
+ first = stage_rows[0]
|
|
|
+ last = stage_rows[-1]
|
|
|
+ stage_values = [item["row"]["value"] for item in stage_rows]
|
|
|
+ basis = (
|
|
|
+ f"油压{first['row']['value']:.3f}~{last['row']['value']:.3f}MPa,"
|
|
|
+ f"最低{min(stage_values):.3f}MPa,连续{len(stage_rows)}个有效运行小时;"
|
|
|
+ f"周期基线油压{baseline:.3f}MPa,相对基线变化{last['deviation']:+.1%}"
|
|
|
+ )
|
|
|
+ batch_ids = sorted({batch for item in stage_rows for batch in item["row"]["batch_ids"]})
|
|
|
+ info = {
|
|
|
+ "机组": f"{unit}号机",
|
|
|
+ "批次": batch_ids[0] if len(batch_ids) == 1 else batch_ids,
|
|
|
+ "周期编号": f"{device_code}-P14",
|
|
|
+ "开始小时": first["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
|
|
|
+ "结束小时": last["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
|
|
|
+ "持续自然小时数": int((last["row"]["time"] - first["row"]["time"]).total_seconds() / 3600) + 1,
|
|
|
+ "有效运行小时数": len(stage_rows),
|
|
|
+ "周期基线油压": baseline,
|
|
|
+ "正常波动带": "",
|
|
|
+ "阶段开始小时油压": first["row"]["value"],
|
|
|
+ "阶段结束小时油压": last["row"]["value"],
|
|
|
+ "阶段最低小时油压": min(stage_values),
|
|
|
+ "阶段开始相对周期基线变化": first["deviation"],
|
|
|
+ "阶段结束相对周期基线变化": last["deviation"],
|
|
|
+ "阶段最大连续下降有效小时数": max(item["decline"] for item in stage_rows),
|
|
|
+ "阶段低报警小时数": sum(item["row"]["value"] <= low for item in stage_rows),
|
|
|
+ "阶段低低报警小时数": sum(item["row"]["value"] <= low_low for item in stage_rows),
|
|
|
+ "是否形成确认等级": "是" if len(stage_rows) >= 3 else "否",
|
|
|
+ }
|
|
|
+ return {
|
|
|
+ "stage": stage,
|
|
|
+ "basis": basis,
|
|
|
+ "device_part": "-".join(dict.fromkeys(point_descriptions.values())),
|
|
|
+ "scan_start": start,
|
|
|
+ "scan_end": end,
|
|
|
+ "info": info,
|
|
|
+ }
|
|
|
+
|
|
|
+ def _analyze_cruciform_vibration_alarm(self, device_code: str, end: datetime) -> dict[str, Any] | None:
|
|
|
+ """Build the strongest 14-day vibration stage for the crosshead alarm."""
|
|
|
+ unit = device_code.rstrip("#").strip()
|
|
|
+ if not device_code.endswith("#") or not device_code[:-1].isdigit():
|
|
|
+ raise ValueError("压缩机编号格式必须为数字+#")
|
|
|
+ start = end - timedelta(days=OIL_PRESSURE_SCAN_DAYS)
|
|
|
+
|
|
|
+ with get_connection() as connection:
|
|
|
+ with connection.cursor() as cursor:
|
|
|
+ configs = {}
|
|
|
+ descriptions = []
|
|
|
+ for point, side in ((10, "联轴器端"), (11, "链轮端"), (41, "运行状态")):
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ SELECT ItemDescription, AlarmType1, AlarmType2, AlarmType3, AlarmType4,
|
|
|
+ AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4
|
|
|
+ FROM site_point WHERE ItemName = %s
|
|
|
+ """,
|
|
|
+ (f"YSJ{unit}_{point}",),
|
|
|
+ )
|
|
|
+ item = cursor.fetchone()
|
|
|
+ if item is None:
|
|
|
+ raise ValueError(f"未找到 {device_code} 的 YSJ_{point} 点位配置")
|
|
|
+ descriptions.append(item.get("ItemDescription") or f"YSJ_{point}")
|
|
|
+ if point != 41:
|
|
|
+ high = self._alarm_limit(item, "PVHigh")
|
|
|
+ high_high = self._alarm_limit(item, "PVHighHigh")
|
|
|
+ if high is None or high_high is None:
|
|
|
+ raise ValueError(f"{device_code} YSJ_{point} 缺少振动报警阈值")
|
|
|
+ configs[side] = {"high": high, "high_high": high_high}
|
|
|
+
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ SELECT FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start,
|
|
|
+ SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples,
|
|
|
+ AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_avg,
|
|
|
+ MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_max,
|
|
|
+ AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_avg,
|
|
|
+ MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_max,
|
|
|
+ GROUP_CONCAT(DISTINCT import_batch_id ORDER BY import_batch_id) AS batch_ids
|
|
|
+ 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)
|
|
|
+ if running / 720 < 0.80:
|
|
|
+ continue
|
|
|
+ item = {
|
|
|
+ "time": row["hour_start"], "running": running,
|
|
|
+ "batch_ids": [int(value) for value in str(row["batch_ids"] or "").split(",") if value],
|
|
|
+ }
|
|
|
+ for key in ("coupling_avg", "coupling_max", "chain_avg", "chain_max"):
|
|
|
+ item[key] = _safe_float(row[key])
|
|
|
+ if item["coupling_max"] is not None or item["chain_max"] is not None:
|
|
|
+ rows.append(item)
|
|
|
+
|
|
|
+ if len(rows) < 3:
|
|
|
+ return None
|
|
|
+ signals = []
|
|
|
+ for index, row in enumerate(rows):
|
|
|
+ prior = rows[max(0, index - 336):index]
|
|
|
+ if len(prior) < 24:
|
|
|
+ continue
|
|
|
+ side_signals = []
|
|
|
+ for side, prefix in (("联轴器端", "coupling"), ("链轮端", "chain")):
|
|
|
+ current = row[f"{prefix}_max"]
|
|
|
+ values = [item[f"{prefix}_max"] for item in prior if item[f"{prefix}_max"] is not None]
|
|
|
+ averages = [item[f"{prefix}_avg"] for item in prior if item[f"{prefix}_avg"] is not None]
|
|
|
+ if current is None or len(values) < 24:
|
|
|
+ continue
|
|
|
+ baseline = float(np.median(values))
|
|
|
+ scale = max(float(np.median(np.abs(np.asarray(values) - baseline))) * 1.4826, abs(baseline) * 0.05, 1e-6)
|
|
|
+ peak_z = (current - baseline) / scale
|
|
|
+ average = row[f"{prefix}_avg"]
|
|
|
+ average_base = float(np.median(averages)) if averages else None
|
|
|
+ average_change = average / average_base - 1 if average is not None and average_base else 0.0
|
|
|
+ limit = configs[side]
|
|
|
+ severe = current >= limit["high_high"] or peak_z >= 8
|
|
|
+ abnormal = severe or current >= limit["high"] or peak_z >= 5 or average_change >= 0.10
|
|
|
+ mild = abnormal or peak_z >= 4 or average_change >= 0.05
|
|
|
+ if mild:
|
|
|
+ side_signals.append({
|
|
|
+ "side": side, "value": current, "baseline": baseline, "scale": scale,
|
|
|
+ "peak_z": peak_z, "average_change": average_change,
|
|
|
+ "level": "严重异常" if severe else "异常" if abnormal else "轻微",
|
|
|
+ })
|
|
|
+ if side_signals:
|
|
|
+ signals.append({"row": row, "signals": side_signals})
|
|
|
+
|
|
|
+ if not signals:
|
|
|
+ return None
|
|
|
+ runs = []
|
|
|
+ current_run = []
|
|
|
+ for signal in signals:
|
|
|
+ if current_run and signal["row"]["time"] - current_run[-1]["row"]["time"] > timedelta(hours=6):
|
|
|
+ runs.append(current_run)
|
|
|
+ current_run = []
|
|
|
+ current_run.append(signal)
|
|
|
+ if current_run:
|
|
|
+ runs.append(current_run)
|
|
|
+ run = max(
|
|
|
+ runs,
|
|
|
+ key=lambda values: (
|
|
|
+ max({"轻微": 1, "异常": 2, "严重异常": 3}[item["level"]] for signal in values for item in signal["signals"]),
|
|
|
+ len(values),
|
|
|
+ ),
|
|
|
+ )
|
|
|
+ all_signals = [item for signal in run for item in signal["signals"]]
|
|
|
+ strongest = max(all_signals, key=lambda item: ({"轻微": 1, "异常": 2, "严重异常": 3}[item["level"]], item["peak_z"]))
|
|
|
+ stage = strongest["level"]
|
|
|
+ first, last = run[0], run[-1]
|
|
|
+ batch_ids = sorted({batch for signal in run for batch in signal["row"]["batch_ids"]})
|
|
|
+ max_value = max(item["value"] for item in all_signals)
|
|
|
+ basis = (
|
|
|
+ f"{strongest['side']}振动峰值或均值偏离近期基线,阶段内{len(run)}个有效运行小时触发;"
|
|
|
+ f"最大振动值{max_value:.3f}mm/s,峰值高于近期基线{strongest['peak_z']:.2f}个稳健尺度"
|
|
|
+ f"(峰值/基线{max_value / max(strongest['baseline'], 1e-6):.2f}倍)"
|
|
|
+ )
|
|
|
+ info = {
|
|
|
+ "机组": f"{unit}号机", "批次": batch_ids[0] if len(batch_ids) == 1 else batch_ids,
|
|
|
+ "周期编号": f"{device_code}-V14", "开始小时": first["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
|
|
|
+ "结束小时": last["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
|
|
|
+ "持续自然小时数": int((last["row"]["time"] - first["row"]["time"]).total_seconds() / 3600) + 1,
|
|
|
+ "有效运行小时数": len(run), "近期基线振动": strongest["baseline"],
|
|
|
+ "正常波动带": f"{strongest['baseline'] - 2 * strongest['scale']:.3f}~{strongest['baseline'] + 2 * strongest['scale']:.3f} mm/s",
|
|
|
+ "阶段开始小时振动": max(item["value"] for item in first["signals"]),
|
|
|
+ "阶段结束小时振动": max(item["value"] for item in last["signals"]),
|
|
|
+ "阶段最大振动值": max_value,
|
|
|
+ "阶段开始相对近期基线变化": max(item["value"] / item["baseline"] - 1 for item in first["signals"]),
|
|
|
+ "阶段结束相对近期基线变化": max(item["value"] / item["baseline"] - 1 for item in last["signals"]),
|
|
|
+ "阶段最大连续异常有效小时数": len(run),
|
|
|
+ "阶段高报警小时数": sum(any(item["value"] >= configs[item["side"]]["high"] for item in signal["signals"]) for signal in run),
|
|
|
+ "阶段高高报警小时数": sum(any(item["value"] >= configs[item["side"]]["high_high"] for item in signal["signals"]) for signal in run),
|
|
|
+ "是否形成确认等级": "是" if stage != "轻微" else "否",
|
|
|
+ "最强异常侧": strongest["side"], "最强侧近期峰值基线": strongest["baseline"],
|
|
|
+ "最强侧近期稳健尺度": strongest["scale"], "阶段最大峰值/基线比例": max_value / max(strongest["baseline"], 1e-6),
|
|
|
+ }
|
|
|
+ return {
|
|
|
+ "stage": stage, "basis": basis, "device_part": "-".join(dict.fromkeys(descriptions)),
|
|
|
+ "scan_start": start, "scan_end": end, "info": info,
|
|
|
+ }
|
|
|
+
|
|
|
+ def sync_pressure_angle(self, device_part: str, device_points: list[str] | None = None) -> dict[str, Any]:
|
|
|
+ """按机组与部位批量生成 statistic_pressure_angle 采样数据。
|
|
|
+
|
|
|
+ 对每个压力 wave_file,若 statistic_pressure_angle 已存在至少一条记录则跳过,
|
|
|
+ 否则读 wave_sample_one(第一个周期)按“每个整数角度取第一个值”生成 360 个角度点。
|
|
|
+ """
|
|
|
+ points = device_points or []
|
|
|
+
|
|
|
+ def database_query():
|
|
|
+ with get_connection() as connection:
|
|
|
+ with connection.cursor() as cursor:
|
|
|
+ clauses = ["measurement_type = %s", "device_part = %s"]
|
|
|
+ params: list[Any] = ["压力", device_part]
|
|
|
+ if points:
|
|
|
+ placeholders = ", ".join(["%s"] * len(points))
|
|
|
+ clauses.append(f"device_point IN ({placeholders})")
|
|
|
+ params.extend(points)
|
|
|
+ cursor.execute(
|
|
|
+ f"SELECT id FROM wave_file WHERE {' AND '.join(clauses)} ORDER BY id",
|
|
|
+ params,
|
|
|
+ )
|
|
|
+ file_ids = [int(row["id"]) for row in cursor.fetchall()]
|
|
|
+ if not file_ids:
|
|
|
+ return {"total": 0, "created": 0, "skipped": 0}
|
|
|
+ placeholders = ", ".join(["%s"] * len(file_ids))
|
|
|
+ cursor.execute(
|
|
|
+ f"SELECT DISTINCT wave_file_id FROM statistic_pressure_angle WHERE wave_file_id IN ({placeholders})",
|
|
|
+ file_ids,
|
|
|
+ )
|
|
|
+ existing = {int(row["wave_file_id"]) for row in cursor.fetchall()}
|
|
|
+ created = 0
|
|
|
+ skipped = 0
|
|
|
+ for file_id in file_ids:
|
|
|
+ if file_id in existing:
|
|
|
+ skipped += 1
|
|
|
+ continue
|
|
|
+ cursor.execute(
|
|
|
+ "SELECT sample_index, signal_value, second_value FROM wave_sample_one "
|
|
|
+ "WHERE wave_file_id = %s ORDER BY sample_index",
|
|
|
+ (file_id,),
|
|
|
+ )
|
|
|
+ rows = cursor.fetchall()
|
|
|
+ if not rows:
|
|
|
+ skipped += 1
|
|
|
+ continue
|
|
|
+ samples = np.asarray(
|
|
|
+ [
|
|
|
+ (float(row["sample_index"]), float(row["signal_value"]),
|
|
|
+ float(row["second_value"]) if row["second_value"] is not None else 0.0)
|
|
|
+ for row in rows
|
|
|
+ ],
|
|
|
+ dtype=float,
|
|
|
+ )
|
|
|
+ sampled = _sample_first_cycle_360(samples)
|
|
|
+ if not sampled:
|
|
|
+ skipped += 1
|
|
|
+ continue
|
|
|
+ cursor.executemany(
|
|
|
+ "INSERT INTO statistic_pressure_angle (wave_file_id, period, angle, pressure) "
|
|
|
+ "VALUES (%s, %s, %s, %s)",
|
|
|
+ [(file_id, 1, angle, pressure) for angle, pressure in sampled],
|
|
|
+ )
|
|
|
+ created += 1
|
|
|
+ return {"total": len(file_ids), "created": created, "skipped": skipped}
|
|
|
+
|
|
|
+ result, source = self._run_with_fallback(
|
|
|
+ database_query,
|
|
|
+ lambda: {"total": 0, "created": 0, "skipped": 0},
|
|
|
+ )
|
|
|
+ return {"source": source, "notice": self._source_notice(source), **result}
|
|
|
+
|
|
|
+ def query_pressure_angle(
|
|
|
+ self,
|
|
|
+ device_part: str,
|
|
|
+ device_points: list[str] | None,
|
|
|
+ min_time: str | None,
|
|
|
+ max_time: str | None,
|
|
|
+ angles: list[int],
|
|
|
+ ) -> dict[str, Any]:
|
|
|
+ """按时间范围和角度列表查询 statistic_pressure_angle 曲线数据。"""
|
|
|
+ if not angles:
|
|
|
+ raise ValueError("请至少输入一个角度")
|
|
|
+ angles = [int(angle) for angle in angles if 0 <= int(angle) <= 360]
|
|
|
+
|
|
|
+ def database_query():
|
|
|
+ with get_connection() as connection:
|
|
|
+ with connection.cursor() as cursor:
|
|
|
+ clauses = ["wf.measurement_type = %s", "wf.device_part = %s"]
|
|
|
+ params: list[Any] = ["压力", device_part]
|
|
|
+ if device_points:
|
|
|
+ placeholders = ", ".join(["%s"] * len(device_points))
|
|
|
+ clauses.append(f"wf.device_point IN ({placeholders})")
|
|
|
+ params.extend(device_points)
|
|
|
+ start = _parse_time(min_time) if min_time else None
|
|
|
+ end = _parse_time(max_time) if max_time else None
|
|
|
+ if start:
|
|
|
+ clauses.append("wf.sample_time >= %s")
|
|
|
+ params.append(start)
|
|
|
+ if end:
|
|
|
+ clauses.append("wf.sample_time <= %s")
|
|
|
+ params.append(end)
|
|
|
+ angle_placeholders = ", ".join(["%s"] * len(angles))
|
|
|
+ clauses.append(f"spa.angle IN ({angle_placeholders})")
|
|
|
+ params.extend(angles)
|
|
|
+ cursor.execute(
|
|
|
+ f"""
|
|
|
+ SELECT wf.sample_time, spa.angle, spa.pressure
|
|
|
+ FROM statistic_pressure_angle spa
|
|
|
+ INNER JOIN wave_file wf ON wf.id = spa.wave_file_id
|
|
|
+ WHERE {' AND '.join(clauses)}
|
|
|
+ ORDER BY wf.sample_time, spa.angle
|
|
|
+ """,
|
|
|
+ params,
|
|
|
+ )
|
|
|
+ rows = cursor.fetchall()
|
|
|
+
|
|
|
+ stop_days = self._query_stop_days(
|
|
|
+ cursor, device_part, device_points, start, end,
|
|
|
+ )
|
|
|
+ series: dict[int, list[dict[str, Any]]] = {}
|
|
|
+ for row in rows:
|
|
|
+ angle = int(round(float(row["angle"])))
|
|
|
+ series.setdefault(angle, []).append({
|
|
|
+ "time": _time_string(row["sample_time"]),
|
|
|
+ "pressure": float(row["pressure"]),
|
|
|
+ })
|
|
|
+ for day in sorted(stop_days):
|
|
|
+ day_time = f"{day} 00:00:00"
|
|
|
+ for angle in angles:
|
|
|
+ series.setdefault(angle, []).append({
|
|
|
+ "time": day_time,
|
|
|
+ "pressure": 0.0,
|
|
|
+ })
|
|
|
+ for angle in angles:
|
|
|
+ series.setdefault(angle, []).sort(key=lambda point: point["time"])
|
|
|
+ return {
|
|
|
+ "angles": angles,
|
|
|
+ "series": [
|
|
|
+ {"angle": angle, "points": series.get(angle, [])}
|
|
|
+ for angle in angles
|
|
|
+ ],
|
|
|
+ }
|
|
|
+
|
|
|
+ result, source = self._run_with_fallback(
|
|
|
+ database_query,
|
|
|
+ lambda: {"angles": angles, "series": [{"angle": angle, "points": []} for angle in angles]},
|
|
|
+ )
|
|
|
+ return {"source": source, "notice": self._source_notice(source), **result}
|
|
|
+
|
|
|
+ @staticmethod
|
|
|
+ def _query_stop_days(
|
|
|
+ cursor: Any,
|
|
|
+ device_part: str,
|
|
|
+ device_points: list[str] | None,
|
|
|
+ start: datetime | None,
|
|
|
+ end: datetime | None,
|
|
|
+ ) -> set[str]:
|
|
|
+ """返回停机日集合:该机组在时间范围内 rpm=0 且当天无 rpm>0 记录的日期。"""
|
|
|
+ clauses = ["device_part = %s", "rpm = 0"]
|
|
|
+ params: list[Any] = [device_part]
|
|
|
+ if device_points:
|
|
|
+ placeholders = ", ".join(["%s"] * len(device_points))
|
|
|
+ clauses.append(f"device_point IN ({placeholders})")
|
|
|
+ params.extend(device_points)
|
|
|
+ if start:
|
|
|
+ clauses.append("sample_time >= %s")
|
|
|
+ params.append(start)
|
|
|
+ if end:
|
|
|
+ clauses.append("sample_time <= %s")
|
|
|
+ params.append(end)
|
|
|
+ cursor.execute(
|
|
|
+ f"""
|
|
|
+ SELECT DISTINCT DATE(sample_time) AS d
|
|
|
+ FROM wave_file
|
|
|
+ WHERE {' AND '.join(clauses)}
|
|
|
+ """,
|
|
|
+ params,
|
|
|
+ )
|
|
|
+ stop_days = {str(row["d"]) for row in cursor.fetchall()}
|
|
|
+
|
|
|
+ run_clauses = ["device_part = %s", "rpm > 0"]
|
|
|
+ run_params: list[Any] = [device_part]
|
|
|
+ if device_points:
|
|
|
+ placeholders = ", ".join(["%s"] * len(device_points))
|
|
|
+ run_clauses.append(f"device_point IN ({placeholders})")
|
|
|
+ run_params.extend(device_points)
|
|
|
+ if start:
|
|
|
+ run_clauses.append("sample_time >= %s")
|
|
|
+ run_params.append(start)
|
|
|
+ if end:
|
|
|
+ run_clauses.append("sample_time <= %s")
|
|
|
+ run_params.append(end)
|
|
|
+ cursor.execute(
|
|
|
+ f"""
|
|
|
+ SELECT DISTINCT DATE(sample_time) AS d
|
|
|
+ FROM wave_file
|
|
|
+ WHERE {' AND '.join(run_clauses)}
|
|
|
+ """,
|
|
|
+ run_params,
|
|
|
+ )
|
|
|
+ run_days = {str(row["d"]) for row in cursor.fetchall()}
|
|
|
+ return stop_days - run_days
|
|
|
|
|
|
@staticmethod
|
|
|
def _group_time_points(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
|
|
@@ -1142,6 +1821,34 @@ class DataService:
|
|
|
"notice": self._source_notice(source),
|
|
|
}
|
|
|
|
|
|
+ def alarm_compressor_options(self) -> dict[str, Any]:
|
|
|
+ """Return running compressor groups using device_part's machine prefix."""
|
|
|
+ def database_query():
|
|
|
+ with get_connection() as connection:
|
|
|
+ with connection.cursor() as cursor:
|
|
|
+ cursor.execute(
|
|
|
+ """
|
|
|
+ SELECT SUBSTRING(device_part, 1, 4) AS device_name,
|
|
|
+ COUNT(1) AS count_no
|
|
|
+ FROM wave_file
|
|
|
+ WHERE rpm > 0 AND device_part <> ''
|
|
|
+ GROUP BY SUBSTRING(device_part, 1, 4)
|
|
|
+ ORDER BY device_name
|
|
|
+ """,
|
|
|
+ )
|
|
|
+ return [
|
|
|
+ {
|
|
|
+ "deviceName": row["device_name"],
|
|
|
+ "deviceCode": f"{str(row['device_name'])[0]}#",
|
|
|
+ "count": int(row["count_no"]),
|
|
|
+ }
|
|
|
+ for row in cursor.fetchall()
|
|
|
+ if row["device_name"]
|
|
|
+ ]
|
|
|
+
|
|
|
+ result, source = self._run_with_fallback(database_query, lambda: [])
|
|
|
+ return {"source": source, "notice": self._source_notice(source), "items": result}
|
|
|
+
|
|
|
def wave_window(
|
|
|
self,
|
|
|
device_part: str,
|