data_service.py 104 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424
  1. from __future__ import annotations
  2. import re
  3. from bisect import bisect_left
  4. from collections import OrderedDict
  5. from datetime import datetime, timedelta
  6. from functools import lru_cache
  7. from time import monotonic
  8. from typing import Any, Callable
  9. import numpy as np
  10. from ..algorithms.cycles import (
  11. DetectedCycle,
  12. PULSES_PER_REVOLUTION,
  13. build_angle_vector,
  14. detect_cycles,
  15. downsample_indices,
  16. )
  17. from ..config import settings
  18. from ..db import get_connection
  19. MEASUREMENT_TYPES = ("位移", "加速度", "压力")
  20. MEASUREMENT_COLORS = {
  21. "压力": "#e4572e",
  22. "位移": "#1f7a8c",
  23. "加速度": "#7b61a8",
  24. }
  25. DEVICE_POINTS = ("压力盖侧", "压力轴侧", "活塞杆沉降", "十字头振动", "自由端振动", "驱动端振动")
  26. DEVICE_POINT_TO_TYPE = {
  27. "压力盖侧": "压力",
  28. "压力轴侧": "压力",
  29. "活塞杆沉降": "位移",
  30. "十字头振动": "加速度",
  31. "自由端振动": "加速度",
  32. "驱动端振动": "加速度",
  33. }
  34. PRIMARY_DEVICE_POINT = "压力盖侧"
  35. PHASES = (
  36. ("排气", "#7b1fa2", 0.0, 120.0),
  37. ("压缩", "#c62828", 120.0, 195.0),
  38. ("膨胀", "#1565c0", 195.0, 270.0),
  39. ("进气", "#2e7d32", 270.0, 330.0),
  40. )
  41. ANNOTATION_LABELS = ("正常", "异常")
  42. OIL_PRESSURE_ALARM_TYPE = "润滑油压力低"
  43. OIL_PRESSURE_DEVICE_PART = "润滑油"
  44. OIL_PRESSURE_DEVICE_POINT = "压力"
  45. # PKS 全场点位:机组号 -> pks_long_sample.import_batch_id
  46. # 7号机=30、8号机=31、9号机=32。
  47. UNIT_BATCH = {"7": 30, "8": 31, "9": 32}
  48. # 全场点位(PKS)每个时间点的曲线窗口:以采样时刻为中心的前后各 7.5 分钟。
  49. PKS_WINDOW_SECONDS = 15 * 60
  50. # 仅 PKS 时间点返回给时间条的最大锚点数:超过后按等步长均匀采样(首尾必保)。
  51. PKS_STRIP_MAX_ANCHORS = 5000
  52. _SITE_POINT_PATTERN = re.compile(r"^YSJ([789])_([1-9]|1[0-9]|2[0-9]|3[0-9]|4[0-1])$")
  53. # 不作为全场点位(PKS)下拉候选的序号(对应机组自己的运行状态/转速信号)。
  54. _EXCLUDED_SITE_INDEXES = frozenset({33, 34, 35, 41})
  55. # Extra samples fetched past cycle_end so detect_cycles can see the zero marker
  56. # that closes the first cycle (cycle_end is exclusive, the marker sits at it).
  57. FIRST_CYCLE_PAD = 64
  58. CYLINDER_BORE_MM = {
  59. "一缸": 360.0,
  60. "二缸": 490.0,
  61. "三缸": 390.0,
  62. "四缸": 490.0,
  63. "五缸": 390.0,
  64. "六缸": 490.0,
  65. }
  66. PISTON_STROKE_MM = 148.0
  67. CONNECTING_ROD_LENGTH_MM = 460.0
  68. CLEARANCE_VOLUME_L_BY_BORE = {
  69. 490.0: 1.51,
  70. 390.0: 0.74,
  71. 360.0: 0.62,
  72. }
  73. DEMO_POINT_NAME = "7号机组一缸压力盖侧"
  74. DEMO_DEVICE_PART = "7号机组一缸"
  75. DEMO_START = datetime(2026, 4, 12, 8, 0, 0)
  76. DEMO_POINT_COUNT = 72
  77. DEMO_SAMPLE_COUNT = 32768
  78. DEMO_REVOLUTION_SAMPLES = 800
  79. DEMO_ID_BY_TYPE = {name: 100000 + index * 1000 for index, name in enumerate(MEASUREMENT_TYPES)}
  80. def _time_string(value: Any) -> str:
  81. if isinstance(value, datetime):
  82. return value.strftime("%Y-%m-%d %H:%M:%S")
  83. return str(value)
  84. def _parse_time(value: str | None) -> datetime | None:
  85. if not value:
  86. return None
  87. return datetime.fromisoformat(value.replace("Z", "+00:00").replace("T", " "))
  88. def _safe_float(value: Any) -> float | None:
  89. if value is None:
  90. return None
  91. number = float(value)
  92. return number if np.isfinite(number) else None
  93. def _stored_first_cycle(samples: np.ndarray) -> DetectedCycle | None:
  94. """Build a cycle from wave_sample_one, which contains one cycle only."""
  95. if samples.ndim != 2 or samples.shape[1] < 3 or len(samples) < 2:
  96. return None
  97. start = 0
  98. end = len(samples)
  99. angle = np.linspace(0.0, 360.0, end - start, endpoint=False)
  100. trigger_offsets = tuple(
  101. int(offset)
  102. for offset in np.flatnonzero(
  103. np.nan_to_num(samples[:, 2], nan=-np.inf) >= 30.0,
  104. )
  105. if offset == 0 or samples[offset - 1, 2] < 30.0
  106. )
  107. return DetectedCycle(
  108. number=1,
  109. start_offset=start,
  110. end_offset=end,
  111. angle=angle,
  112. signal=samples[start:end, 1].copy(),
  113. trigger_offsets=trigger_offsets,
  114. )
  115. def _detect_source_cycles(
  116. samples: np.ndarray,
  117. single_cycle: bool,
  118. ) -> tuple[list[DetectedCycle], dict[str, Any]]:
  119. """Detect cycles: wave_sample_one holds exactly one stored cycle already."""
  120. if single_cycle:
  121. stored = _stored_first_cycle(samples)
  122. if stored is None:
  123. return [], {"completeCycleCount": 0}
  124. return [stored], {"completeCycleCount": 1, "storedFirstCycle": 1}
  125. return detect_cycles(samples)
  126. def _series_extent(items: list[dict[str, Any]]) -> tuple[float, float]:
  127. """Whole-window min/max over a series' finite raw values."""
  128. if not items:
  129. return (0.0, 1.0)
  130. values = np.fromiter((item["rawValue"] for item in items), dtype=float, count=len(items))
  131. values = values[np.isfinite(values)]
  132. if values.size == 0:
  133. return (0.0, 1.0)
  134. return (float(values.min()), float(values.max()))
  135. def _annotation_dict(row: dict[str, Any]) -> dict[str, Any]:
  136. return {
  137. "id": int(row["id"]),
  138. "waveFileId": int(row["wave_file_id"]),
  139. "label": row["label"],
  140. "periodStart": int(row["period_start"]),
  141. "periodEnd": int(row["period_end"]),
  142. "sampleIndexStart": int(row["sample_index_start"]),
  143. "sampleIndexEnd": int(row["sample_index_end"]),
  144. }
  145. def _alarm_dict(row: dict[str, Any]) -> dict[str, Any]:
  146. return {
  147. "id": int(row["id"]),
  148. "deviceCode": row["device_code"] or "",
  149. "devicePart": row["device_part"] or "",
  150. "devicePoint": row["device_point"] or "",
  151. "alarmType": row["alarm_type"] or "",
  152. "alarmDes": row["alarm_des"] or "",
  153. "alarmTimeStart": _time_string(row["alarm_time_start"]),
  154. "alarmTimeEnd": _time_string(row["alarm_time_end"]),
  155. }
  156. def _validate_device_points(
  157. values: list[str] | tuple[str, ...] | None,
  158. *,
  159. allow_empty: bool = False,
  160. ) -> list[str]:
  161. """校验波形点位。None 表示未传参 → 回退为全部点位;显式空列表仅在 allow_empty 时放行。"""
  162. selected = list(DEVICE_POINTS) if values is None else list(values)
  163. invalid = [value for value in selected if value not in DEVICE_POINTS]
  164. if invalid:
  165. raise ValueError(f"不支持的测试点位:{'、'.join(invalid)}")
  166. result = [value for value in DEVICE_POINTS if value in selected]
  167. if not result and not allow_empty:
  168. raise ValueError("请至少选择一个测试点位")
  169. return result
  170. def _primary_device_point(points: list[str]) -> str:
  171. """Prefer pressure cap, then pressure shaft, then any selected point."""
  172. for point in ("压力盖侧", "压力轴侧"):
  173. if point in points:
  174. return point
  175. return points[0]
  176. def _unit_number(device_part: str) -> str | None:
  177. """从 机组与部位 提取机组号(7/8/9),非 7/8/9 机组返回 None。"""
  178. match = re.match(r"^([789])号机组", device_part.strip())
  179. return match.group(1) if match else None
  180. def _site_point_column(item_name: str, unit: str | None) -> str | None:
  181. """校验全场点位名称并映射到 pks_long_sample 的列名(如 YSJ7_3 -> YSJ_3)。"""
  182. match = _SITE_POINT_PATTERN.match(item_name)
  183. if not match or match.group(1) != unit:
  184. return None
  185. return f"YSJ_{match.group(2)}"
  186. class DataService:
  187. # 数据库失败后的重试冷却时间(秒)。超过该时间后自动重连数据库,
  188. # 避免一次网络抖动就把服务永久锁死在演示数据模式。
  189. DB_RETRY_COOLDOWN = 30.0
  190. def __init__(self) -> None:
  191. self._db_failed = settings.demo_mode == "always"
  192. self._db_failed_at = monotonic() if self._db_failed else 0.0
  193. self._last_db_error = ""
  194. self._demo_annotations: dict[int, dict[str, Any]] = {}
  195. self._demo_annotation_seq = 1
  196. # 机组号 -> pks_long_sample 该批次 [min_time, max_time],进程内只查一次。
  197. self._pks_batch_bounds: dict[str, tuple[datetime, datetime]] = {}
  198. @property
  199. def source(self) -> str:
  200. return "demo" if self._db_failed else "database"
  201. @property
  202. def last_db_error(self) -> str:
  203. return self._last_db_error
  204. def _run_with_fallback(
  205. self,
  206. database_function: Callable[[], Any],
  207. demo_function: Callable[[], Any],
  208. ) -> tuple[Any, str]:
  209. if self._db_failed:
  210. if settings.demo_mode == "always":
  211. return demo_function(), "demo"
  212. if monotonic() - self._db_failed_at < self.DB_RETRY_COOLDOWN:
  213. return demo_function(), "demo"
  214. # 冷却结束,重新尝试数据库,数据库恢复后可自动切回真实数据。
  215. try:
  216. result = database_function()
  217. self._db_failed = False
  218. self._last_db_error = ""
  219. return result, "database"
  220. except Exception as error:
  221. if settings.demo_mode == "never":
  222. raise
  223. self._db_failed = True
  224. self._db_failed_at = monotonic()
  225. self._last_db_error = str(error)
  226. return demo_function(), "demo"
  227. def query_options(self) -> dict[str, Any]:
  228. def database_query():
  229. with get_connection() as connection:
  230. with connection.cursor() as cursor:
  231. cursor.execute(
  232. """
  233. SELECT device_part, device_point, measurement_type,
  234. MAX(sample_time) AS max_time,
  235. MIN(sample_time) AS min_time,
  236. COUNT(*) AS file_count
  237. FROM wave_file
  238. WHERE rpm > 0 AND device_part <> '' AND device_point <> ''
  239. GROUP BY device_part, device_point, measurement_type
  240. ORDER BY device_part, device_point
  241. """,
  242. )
  243. rows = cursor.fetchall()
  244. return [
  245. {
  246. "devicePart": row["device_part"],
  247. "devicePoint": row["device_point"],
  248. "measurementType": row["measurement_type"],
  249. "minTime": _time_string(row["min_time"]),
  250. "maxTime": _time_string(row["max_time"]),
  251. "fileCount": int(row["file_count"]),
  252. }
  253. for row in rows
  254. ]
  255. rows, source = self._run_with_fallback(database_query, self._demo_options)
  256. device_parts = list(dict.fromkeys(row["devicePart"] for row in rows))
  257. return {
  258. "source": source,
  259. "measurementTypes": list(MEASUREMENT_TYPES),
  260. "deviceParts": device_parts,
  261. "devicePoints": list(DEVICE_POINTS),
  262. "devicePointToType": dict(DEVICE_POINT_TO_TYPE),
  263. "options": rows,
  264. "notice": self._source_notice(source),
  265. }
  266. def abnormal_counts(self) -> dict[str, int]:
  267. """每个 point_name 的异常文件数量(tspluse_status > 0)。"""
  268. def database_query():
  269. with get_connection() as connection:
  270. with connection.cursor() as cursor:
  271. cursor.execute(
  272. """
  273. SELECT point_name, COUNT(*) AS cnt
  274. FROM wave_file
  275. WHERE tspluse_status > 0 AND point_name <> ''
  276. GROUP BY point_name
  277. """,
  278. )
  279. rows = cursor.fetchall()
  280. return {row["point_name"]: int(row["cnt"]) for row in rows}
  281. def demo_query():
  282. return {}
  283. result, _source = self._run_with_fallback(database_query, demo_query)
  284. return result
  285. def tspluse_ruler(self) -> dict[str, Any]:
  286. """压力部位 tspluse_status 的全局标尺,进入页面时只查询一次。"""
  287. def database_query():
  288. with get_connection() as connection:
  289. with connection.cursor() as cursor:
  290. cursor.execute(
  291. """
  292. SELECT MIN(tspluse_status) AS min_status,
  293. MAX(tspluse_status) AS max_status
  294. FROM wave_file
  295. WHERE rpm > 0 AND measurement_type = '压力'
  296. """,
  297. )
  298. row = cursor.fetchone()
  299. return {
  300. "min": int(row["min_status"]) if row and row["min_status"] is not None else 0,
  301. "max": int(row["max_status"]) if row and row["max_status"] is not None else 0,
  302. }
  303. def demo_query():
  304. return {"min": 0, "max": 0}
  305. result, source = self._run_with_fallback(database_query, demo_query)
  306. result["source"] = source
  307. return result
  308. def _pks_batch_range(self, unit: str) -> tuple[datetime, datetime] | None:
  309. """机组 pks 批次的时间覆盖范围(仅查一次并缓存)。"""
  310. cached = self._pks_batch_bounds.get(unit)
  311. if cached is not None:
  312. return cached
  313. batch = UNIT_BATCH[unit]
  314. with get_connection() as connection:
  315. with connection.cursor() as cursor:
  316. cursor.execute(
  317. "SELECT MIN(sample_time) AS lo, MAX(sample_time) AS hi "
  318. "FROM pks_long_sample WHERE import_batch_id = %s",
  319. (batch,),
  320. )
  321. row = cursor.fetchone()
  322. bounds = (row["lo"], row["hi"]) if row and row["lo"] is not None else None
  323. self._pks_batch_bounds[unit] = bounds
  324. return bounds
  325. def site_points(self, device_part: str) -> dict[str, Any]:
  326. """机组(7/8/9)的全场点位记录(site_point 表中 YSJ{机组号}_1..41)。
  327. 附带该机组 pks 批次的时间覆盖范围,供“仅 PKS 点位”组合自动赋值
  328. 开始/结束时间。
  329. """
  330. unit = _unit_number(device_part)
  331. def database_query():
  332. if unit is None:
  333. return {"items": [], "minTime": None, "maxTime": None}
  334. allowed = {
  335. f"YSJ{unit}_{n}"
  336. for n in range(1, 42)
  337. if n not in _EXCLUDED_SITE_INDEXES
  338. }
  339. with get_connection() as connection:
  340. with connection.cursor() as cursor:
  341. cursor.execute(
  342. "SELECT ItemName, ItemDescription FROM site_point WHERE ItemName LIKE %s",
  343. (f"YSJ{unit}\\_%",),
  344. )
  345. rows = cursor.fetchall()
  346. items = [
  347. {"itemName": row["ItemName"], "itemDescription": row["ItemDescription"] or ""}
  348. for row in rows
  349. if row["ItemName"] in allowed
  350. ]
  351. items.sort(key=lambda item: int(item["itemName"].rsplit("_", 1)[1]))
  352. bounds = self._pks_batch_range(unit)
  353. return {
  354. "items": items,
  355. "minTime": _time_string(bounds[0]) if bounds else None,
  356. "maxTime": _time_string(bounds[1]) if bounds else None,
  357. }
  358. def demo_query():
  359. return {"items": [], "minTime": None, "maxTime": None}
  360. result, source = self._run_with_fallback(database_query, demo_query)
  361. return {
  362. "source": source,
  363. "unit": unit,
  364. "items": result["items"],
  365. "minTime": result["minTime"],
  366. "maxTime": result["maxTime"],
  367. "notice": self._source_notice(source),
  368. }
  369. @staticmethod
  370. def _attach_site_values(points: list[dict[str, Any]], unit: str, site_points: list[str]) -> None:
  371. """把 pks_long_sample 的最近邻值挂到每个时间点的 siteValues 上。
  372. pks 数据 5 秒一条、wave_file 15 分钟一条。按用户口径做分钟/5秒级对齐:
  373. 把 wave 采样时刻四舍五入到最近的 5 秒格点,再用一次 ``sample_time IN (...)``
  374. 精确取数(结果行数 = 时间点数),避免把整段 pks 拉出来。
  375. """
  376. columns = [
  377. (item_name, column)
  378. for item_name in site_points
  379. if (column := _site_point_column(item_name, unit)) is not None
  380. ]
  381. if not columns or not points:
  382. for point in points:
  383. point.setdefault("siteValues", {})
  384. return
  385. rounded: list[datetime] = []
  386. for point in points:
  387. timestamp = datetime.strptime(point["sampleTime"], "%Y-%m-%d %H:%M:%S")
  388. rounded.append(datetime.fromtimestamp(round(timestamp.timestamp() / 5.0) * 5))
  389. chunk_size = 1000
  390. select_expr = ", ".join(f"`{column}`" for _item, column in columns)
  391. for index in range(0, len(rounded), chunk_size):
  392. chunk = rounded[index:index + chunk_size]
  393. placeholders = ", ".join(["%s"] * len(chunk))
  394. with get_connection() as connection:
  395. with connection.cursor() as cursor:
  396. cursor.execute(
  397. f"SELECT sample_time, {select_expr} FROM pks_long_sample "
  398. f"WHERE import_batch_id = %s AND sample_time IN ({placeholders})",
  399. (UNIT_BATCH[unit], *chunk),
  400. )
  401. rows = cursor.fetchall()
  402. by_time: dict[datetime, dict[str, Any]] = {
  403. row["sample_time"]: row for row in rows
  404. }
  405. for point_index in range(index, min(index + chunk_size, len(rounded))):
  406. target_time = rounded[point_index]
  407. row = by_time.get(target_time)
  408. if row is None:
  409. continue
  410. site_values = points[point_index].setdefault("siteValues", {})
  411. for item_name, column in columns:
  412. value = row[column]
  413. site_values[item_name] = float(value) if value is not None else None
  414. @staticmethod
  415. def _build_site_series(
  416. site_points: list[str],
  417. unit: str | None,
  418. points: list[dict[str, Any]],
  419. ) -> dict[str, Any]:
  420. """把每个时间点扩展为前后各 7.5 分钟的 pks 5s 曲线段。
  421. 每个选中点位(PKS)在该窗口内的每个时间点不再只画单值,而是以该
  422. 时间点采样时刻为中心,取 [t-7.5min, t+7.5min) 的原始 5s 数据映射
  423. 到该时间点在 x 轴占据的格子(索引 slot ~ slot+1)内连成一小段曲线。
  424. 相邻时间点若恰好间隔 15 分钟,则相邻窗口首尾衔接、无重叠。
  425. 数据按列做一次整段范围查询,再按时间点二分切段。
  426. """
  427. columns = [
  428. (item_name, column)
  429. for item_name in site_points
  430. if (column := _site_point_column(item_name, unit)) is not None
  431. ]
  432. centers: list[tuple[int, datetime]] = []
  433. for index, point in enumerate(points):
  434. try:
  435. timestamp = datetime.strptime(point["sampleTime"], "%Y-%m-%d %H:%M:%S")
  436. except (TypeError, ValueError):
  437. continue
  438. centers.append((index, timestamp))
  439. if not columns or not centers:
  440. return {"points": site_points, "series": []}
  441. half = timedelta(seconds=PKS_WINDOW_SECONDS // 2)
  442. query_start = min(timestamp for _index, timestamp in centers) - half
  443. query_end = max(timestamp for _index, timestamp in centers) + half
  444. select_expr = ", ".join(f"`{column}`" for _item, column in columns)
  445. with get_connection() as connection:
  446. with connection.cursor() as cursor:
  447. cursor.execute(
  448. f"SELECT sample_time, {select_expr} FROM pks_long_sample "
  449. "WHERE import_batch_id = %s AND sample_time >= %s AND sample_time < %s "
  450. "ORDER BY sample_time",
  451. (UNIT_BATCH[unit], query_start, query_end),
  452. )
  453. rows = cursor.fetchall()
  454. sample_times = [row["sample_time"] for row in rows]
  455. series = []
  456. for item_name, column in columns:
  457. values = [row[column] for row in rows]
  458. data: list[dict[str, Any]] = []
  459. for index, center in centers:
  460. low = bisect_left(sample_times, center - half)
  461. high = bisect_left(sample_times, center + half)
  462. for pos in range(low, high):
  463. raw_value = values[pos]
  464. if raw_value is None:
  465. continue
  466. value = float(raw_value)
  467. if not np.isfinite(value):
  468. continue
  469. x = index + 0.5 + (sample_times[pos] - center).total_seconds() / PKS_WINDOW_SECONDS
  470. data.append(
  471. {
  472. "value": [x, value],
  473. "x": x,
  474. "rawValue": value,
  475. "sampleTime": _time_string(sample_times[pos]),
  476. },
  477. )
  478. series.append({"itemName": item_name, "data": data})
  479. return {"points": site_points, "series": series}
  480. def time_points(
  481. self,
  482. device_part: str,
  483. device_points: list[str] | None,
  484. min_time: str | None,
  485. max_time: str | None,
  486. include_stopped: bool = False,
  487. min_status: int | None = None,
  488. status_filter: list[str] | None = None,
  489. site_points: list[str] | None = None,
  490. ) -> dict[str, Any]:
  491. if not device_part.strip():
  492. raise ValueError("机组与部位不能为空")
  493. status_filters = {value for value in (status_filter or [])}
  494. unknown = status_filters - {"abnormal", "no_cycle"}
  495. if unknown:
  496. raise ValueError(f"不支持的状态筛选:{'、'.join(sorted(unknown))}")
  497. start = _parse_time(min_time)
  498. end = _parse_time(max_time)
  499. if start and end and start > end:
  500. raise ValueError("开始时间不能晚于结束时间")
  501. unit = _unit_number(device_part)
  502. site_list = list(site_points or [])
  503. site_columns: list[tuple[str, str]] = []
  504. if unit:
  505. site_columns = [
  506. (item, column)
  507. for item in site_list
  508. if (column := _site_point_column(item, unit)) is not None
  509. ]
  510. if device_points is None or device_points == []:
  511. selected_points = (
  512. [] if (unit and site_columns) else _validate_device_points(None)
  513. )
  514. else:
  515. selected_points = _validate_device_points(device_points, allow_empty=True)
  516. pks_only = bool(unit and site_columns and not selected_points)
  517. if not selected_points and not pks_only:
  518. raise ValueError("请至少选择一个测试点位")
  519. point_names = [f"{device_part}{point}" for point in selected_points]
  520. if pks_only:
  521. points, source = self._run_with_fallback(
  522. lambda: self._pks_time_points(unit, site_columns, start, end, include_stopped),
  523. lambda: [],
  524. )
  525. return {
  526. "source": source,
  527. "devicePart": device_part,
  528. "devicePoints": [],
  529. "total": len(points),
  530. "points": points,
  531. "referencePoints": [],
  532. "notice": self._source_notice(source),
  533. }
  534. def database_query():
  535. point_placeholders = ", ".join(["%s"] * len(point_names))
  536. clauses = [
  537. f"point_name IN ({point_placeholders})",
  538. ]
  539. params: list[Any] = list(point_names)
  540. if not include_stopped:
  541. clauses.append("rpm > 0")
  542. if min_status is not None and min_status > 0 and "no_cycle" not in status_filters:
  543. clauses.append("tspluse_status >= %s")
  544. params.append(min_status)
  545. if "abnormal" in status_filters and "no_cycle" in status_filters:
  546. clauses.append("(tspluse_status > 0 OR tspluse_status = -1)")
  547. elif "abnormal" in status_filters:
  548. clauses.append("tspluse_status > 0")
  549. elif "no_cycle" in status_filters:
  550. clauses.append("tspluse_status = -1")
  551. if start:
  552. clauses.append("sample_time >= %s")
  553. params.append(start)
  554. if end:
  555. clauses.append("sample_time <= %s")
  556. params.append(end)
  557. with get_connection() as connection:
  558. with connection.cursor() as cursor:
  559. cursor.execute(
  560. f"""
  561. SELECT id, point_name, device_point, measurement_type,
  562. sample_time, sample_count, sample_frequency_hz,
  563. rpm, tspluse_status
  564. FROM wave_file
  565. WHERE {' AND '.join(clauses)}
  566. ORDER BY sample_time ASC, id ASC
  567. """,
  568. params,
  569. )
  570. rows = cursor.fetchall()
  571. reference_rows: list[dict[str, Any]] = []
  572. if status_filters and rows:
  573. primary = _primary_device_point(selected_points)
  574. reference_where = [
  575. "point_name = %s",
  576. "measurement_type = '压力'",
  577. "tspluse_status = 0",
  578. ]
  579. reference_params: list[Any] = [f"{device_part}{primary}"]
  580. if not include_stopped:
  581. reference_where.append("rpm > 0")
  582. if start:
  583. reference_where.append("sample_time >= %s")
  584. reference_params.append(start)
  585. if end:
  586. reference_where.append("sample_time <= %s")
  587. reference_params.append(end)
  588. reference_params.append(rows[0]["sample_time"])
  589. with get_connection() as connection:
  590. with connection.cursor() as cursor:
  591. cursor.execute(
  592. f"""
  593. SELECT id, point_name, device_point, measurement_type,
  594. sample_time, sample_count, sample_frequency_hz,
  595. rpm, tspluse_status
  596. FROM wave_file
  597. WHERE {' AND '.join(reference_where)}
  598. ORDER BY ABS(TIMESTAMPDIFF(SECOND, sample_time, %s)) ASC
  599. LIMIT 1
  600. """,
  601. reference_params,
  602. )
  603. reference_rows = cursor.fetchall()
  604. return self._group_time_points(rows), self._group_time_points(reference_rows)
  605. def demo_query():
  606. return self._demo_time_points(device_part, selected_points, start, end), []
  607. result, source = self._run_with_fallback(database_query, demo_query)
  608. points, reference = result
  609. # 全场点位(PKS)数据是辅助层:任意失败都静默跳过,不影响主查询。
  610. try:
  611. unit = _unit_number(device_part)
  612. if unit and site_points:
  613. self._attach_site_values(points, unit, site_points)
  614. except Exception:
  615. pass
  616. for point in points:
  617. point.setdefault("siteValues", {})
  618. return {
  619. "source": source,
  620. "devicePart": device_part,
  621. "devicePoints": selected_points,
  622. "total": len(points),
  623. "points": points,
  624. "referencePoints": reference,
  625. "notice": self._source_notice(source),
  626. }
  627. def faults(
  628. self,
  629. device_part: str,
  630. min_time: str | None,
  631. max_time: str | None,
  632. ) -> dict[str, Any]:
  633. """Return compressor_fault rows for the unit of ``device_part``.
  634. ``compressor_fault.unit_name`` ("7号机组") is matched against the unit
  635. prefix of ``wave_file.device_part`` ("7号机组一缸"). Only faults whose
  636. ``fault_date`` falls inside the requested time range are returned.
  637. """
  638. unit = _unit_number(device_part)
  639. start = _parse_time(min_time)
  640. end = _parse_time(max_time)
  641. def database_query():
  642. if unit is None:
  643. return []
  644. clauses = ["(unit_name = %s OR unit_name LIKE %s)"]
  645. params: list[Any] = [f"{unit}号机组", f"{unit}号机组%"]
  646. if start:
  647. clauses.append("fault_date >= DATE(%s)")
  648. params.append(start)
  649. if end:
  650. clauses.append("fault_date <= DATE(%s)")
  651. params.append(end)
  652. with get_connection() as connection:
  653. with connection.cursor() as cursor:
  654. cursor.execute(
  655. f"""
  656. SELECT id, unit_name, fault_date, fault_category
  657. FROM compressor_fault
  658. WHERE {' AND '.join(clauses)}
  659. ORDER BY fault_date ASC, id ASC
  660. """,
  661. params,
  662. )
  663. rows = cursor.fetchall()
  664. return [
  665. {
  666. "id": int(row["id"]),
  667. "unitName": row["unit_name"] or "",
  668. "faultDate": _time_string(row["fault_date"])[:10],
  669. "faultCategory": row["fault_category"] or "",
  670. }
  671. for row in rows
  672. ]
  673. result, source = self._run_with_fallback(database_query, lambda: [])
  674. return {
  675. "source": source,
  676. "unit": unit,
  677. "faults": result,
  678. "notice": self._source_notice(source),
  679. }
  680. def list_alarms(self, current_time: str | None = None) -> dict[str, Any]:
  681. """列出在给定时刻命中的告警。
  682. 命中条件为 ``alarm_time_start <= 当前时间 <= alarm_time_end``,
  683. 即告警时间段覆盖“告警当前时间”。
  684. """
  685. moment = _parse_time(current_time) or datetime.now()
  686. def database_query():
  687. with get_connection() as connection:
  688. with connection.cursor() as cursor:
  689. cursor.execute(
  690. """
  691. SELECT id, device_code, device_part, device_point, alarm_type,
  692. alarm_des, alarm_time_start, alarm_time_end
  693. FROM compressor_alarm
  694. WHERE alarm_time_start <= %s AND alarm_time_end >= %s
  695. ORDER BY alarm_time_start DESC, id DESC
  696. """,
  697. (moment, moment),
  698. )
  699. rows = cursor.fetchall()
  700. return [_alarm_dict(row) for row in rows]
  701. result, source = self._run_with_fallback(database_query, lambda: [])
  702. return {
  703. "source": source,
  704. "currentTime": _time_string(moment),
  705. "alarms": result,
  706. "notice": self._source_notice(source),
  707. }
  708. def scan_alarm(self, payload: dict[str, Any]) -> dict[str, Any]:
  709. """执行告警算法并按原有协议写入 compressor_alarm。"""
  710. device_code = str(payload.get("device_code") or "").strip()
  711. alarm_type = str(payload.get("alarm_type") or "").strip()
  712. start = _parse_time(payload.get("forecast_time"))
  713. try:
  714. hours = float(payload.get("hours") or 0)
  715. except (TypeError, ValueError):
  716. raise ValueError("告警时长必须是数字") from None
  717. if not device_code:
  718. raise ValueError("压缩机不能为空")
  719. if not alarm_type:
  720. raise ValueError("算法不能为空")
  721. if start is None:
  722. raise ValueError("预警时间不能为空")
  723. if hours <= 0:
  724. raise ValueError("告警时长必须大于 0")
  725. if alarm_type == OIL_PRESSURE_ALARM_TYPE:
  726. hours = 24.0
  727. analysis = self._analyze_oil_pressure_alarm(device_code, start)
  728. if analysis is None:
  729. return {
  730. "source": "database",
  731. "notice": None,
  732. "action": "no_alarm",
  733. "id": 0,
  734. }
  735. device_part = OIL_PRESSURE_DEVICE_PART
  736. device_point = OIL_PRESSURE_DEVICE_POINT
  737. alarm_des = analysis
  738. status = 1
  739. else:
  740. # Keep the existing endpoint contract for algorithms not yet implemented.
  741. device_part = "测试组件"
  742. device_point = "测试组件"
  743. alarm_des = "测试数据"
  744. status = 0
  745. end = start + timedelta(hours=hours)
  746. def database_query():
  747. with get_connection() as connection:
  748. with connection.cursor() as cursor:
  749. cursor.execute(
  750. """
  751. SELECT id FROM compressor_alarm
  752. WHERE device_code = %s AND device_part = %s
  753. AND device_point = %s AND alarm_type = %s
  754. AND alarm_time_start <= %s AND alarm_time_end >= %s
  755. ORDER BY id ASC
  756. LIMIT 1
  757. """,
  758. (device_code, device_part, device_point, alarm_type, start, start),
  759. )
  760. row = cursor.fetchone()
  761. if row is not None:
  762. cursor.execute(
  763. """
  764. UPDATE compressor_alarm
  765. SET alarm_time_end = %s, alarm_des = %s, status = %s
  766. WHERE id = %s
  767. """,
  768. (end, alarm_des, status, row["id"]),
  769. )
  770. return {"action": "updated", "id": int(row["id"])}
  771. cursor.execute(
  772. """
  773. INSERT INTO compressor_alarm
  774. (device_code, device_part, device_point, alarm_type,
  775. alarm_des, alarm_time_start, alarm_time_end, status)
  776. VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
  777. """,
  778. (device_code, device_part, device_point, alarm_type, alarm_des, start, end, status),
  779. )
  780. return {"action": "inserted", "id": int(cursor.lastrowid)}
  781. result, source = self._run_with_fallback(
  782. database_query,
  783. lambda: {"action": "demo", "id": 0},
  784. )
  785. return {"source": source, "notice": self._source_notice(source), **result}
  786. @staticmethod
  787. def _alarm_limit(item: dict[str, Any], alarm_type: str) -> float | None:
  788. for index in range(1, 5):
  789. if str(item.get(f"AlarmType{index}") or "").strip() == alarm_type:
  790. return _safe_float(item.get(f"AlarmLimit{index}"))
  791. return None
  792. def _analyze_oil_pressure_alarm(self, device_code: str, end: datetime) -> str | None:
  793. """Analyze YSJ_5 in the 24 hours before ``end``.
  794. A low-pressure alarm requires either a configured low/low-low threshold
  795. breach or a sustained decline of at least 10% from the 24-hour median.
  796. The trend rule requires three consecutive non-rising valid hours or a
  797. negative six-hour linear trend. Direct low-low threshold breaches remain
  798. immediately reportable.
  799. """
  800. unit = device_code.rstrip("#").strip()
  801. if unit not in UNIT_BATCH:
  802. raise ValueError("压缩机必须是 7#、8# 或 9#")
  803. start = end - timedelta(hours=24)
  804. with get_connection() as connection:
  805. with connection.cursor() as cursor:
  806. cursor.execute(
  807. """
  808. SELECT AlarmType1, AlarmType2, AlarmType3, AlarmType4,
  809. AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4
  810. FROM site_point
  811. WHERE ItemName = %s
  812. """,
  813. (f"YSJ{unit}_5",),
  814. )
  815. config = cursor.fetchone()
  816. if config is None:
  817. raise ValueError(f"未找到 {device_code} 的润滑油压力报警配置")
  818. low = self._alarm_limit(config, "PVLow")
  819. low_low = self._alarm_limit(config, "PVLowLow")
  820. if low is None or low_low is None:
  821. raise ValueError(f"{device_code} 缺少润滑油压力低压报警阈值")
  822. cursor.execute(
  823. """
  824. SELECT FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start,
  825. COUNT(*) AS samples,
  826. SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples,
  827. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_avg
  828. FROM pks_long_sample
  829. WHERE device_code = %s
  830. AND sample_time >= %s AND sample_time < %s
  831. GROUP BY FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600)
  832. ORDER BY hour_start
  833. """,
  834. (device_code, start, end),
  835. )
  836. rows = []
  837. for row in cursor.fetchall():
  838. running = int(row["running_samples"] or 0)
  839. value = _safe_float(row["oil_avg"])
  840. if running / 720 >= 0.80 and value is not None:
  841. rows.append({"time": row["hour_start"], "value": value, "running": running})
  842. if not rows:
  843. return None
  844. values = [row["value"] for row in rows]
  845. baseline = float(np.median(values))
  846. current = rows[-1]["value"]
  847. relative_drop = (baseline - current) / baseline if baseline else 0.0
  848. low_low_rows = [row for row in rows if row["value"] <= low_low]
  849. low_rows = [row for row in rows if row["value"] <= low]
  850. consecutive_low = 0
  851. for row in reversed(rows):
  852. if row["value"] <= low:
  853. consecutive_low += 1
  854. else:
  855. break
  856. consecutive_low_low = 0
  857. for row in reversed(rows):
  858. if row["value"] <= low_low:
  859. consecutive_low_low += 1
  860. else:
  861. break
  862. consecutive_decline = 1
  863. for index in range(len(rows) - 1, 0, -1):
  864. if rows[index]["time"] - rows[index - 1]["time"] != timedelta(hours=1):
  865. break
  866. if rows[index]["value"] > rows[index - 1]["value"]:
  867. break
  868. consecutive_decline += 1
  869. trend_values = [row["value"] for row in rows[-6:]]
  870. trend = float(np.polyfit(range(len(trend_values)), trend_values, 1)[0]) if len(trend_values) >= 3 else None
  871. downward = consecutive_decline >= 3 or (trend is not None and trend < -0.0005)
  872. evolution_alarm = relative_drop >= 0.10 and downward
  873. if not low_low_rows and not low_rows and not evolution_alarm:
  874. return None
  875. reasons = [
  876. f"前24小时有效运行数据{len(rows)}小时",
  877. f"当前小时油压{current:.3f},24小时基线{baseline:.3f},相对下降{relative_drop:.1%}",
  878. ]
  879. if consecutive_low_low:
  880. reasons.append(f"连续{consecutive_low_low}小时低于低低报警阈值{low_low:.3f}")
  881. elif low_low_rows:
  882. reasons.append(f"低低报警区间出现{len(low_low_rows)}小时")
  883. elif consecutive_low:
  884. reasons.append(f"连续{consecutive_low}小时低于低报警阈值{low:.3f}")
  885. else:
  886. reasons.append(f"低报警区间出现{len(low_rows)}小时")
  887. if evolution_alarm:
  888. reasons.append(
  889. f"相对24小时基线下降至少10%,连续下降{consecutive_decline}小时,"
  890. f"近6小时斜率{trend:.6f}" if trend is not None
  891. else f"相对24小时基线下降至少10%,连续下降{consecutive_decline}小时"
  892. )
  893. return ";".join(reasons)
  894. @staticmethod
  895. def _group_time_points(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
  896. grouped: OrderedDict[Any, dict[str, Any]] = OrderedDict()
  897. for row in rows:
  898. # Channels from one acquisition batch can be stamped a few seconds
  899. # apart (e.g. 10:30:02 vs 10:30:05), yet the acquisition cadence is
  900. # 15 minutes. Align on the minute so one batch is one time point.
  901. sample_time = row["sample_time"]
  902. key = str(sample_time)[:16]
  903. point = grouped.setdefault(
  904. key,
  905. {
  906. "sampleTime": _time_string(sample_time),
  907. "files": {},
  908. },
  909. )
  910. device_point = row["device_point"]
  911. point["files"].setdefault(device_point, {
  912. "id": int(row["id"]),
  913. "devicePoint": device_point,
  914. "measurementType": row["measurement_type"],
  915. "sampleCount": int(row["sample_count"] or 0),
  916. "sampleFrequencyHz": int(row["sample_frequency_hz"] or 0),
  917. "rpm": float(row["rpm"] or 0),
  918. "status": int(row.get("tspluse_status") or 0),
  919. })
  920. return [
  921. {"index": index, **point}
  922. for index, point in enumerate(grouped.values())
  923. ]
  924. def _pks_time_points(
  925. self,
  926. unit: str,
  927. columns: list[tuple[str, str]],
  928. start: datetime | None,
  929. end: datetime | None,
  930. include_stopped: bool,
  931. ) -> list[dict[str, Any]]:
  932. """PKS-only 时间点:按 15 分钟槽从 pks_long_sample 生成骨架。
  933. columns 为 [(item_name, pks 列名), ...],siteValues 直挂所选各列
  934. 原值。停机/运转只由该机组批次内的 YSJ_41 决定(>0=运转、=0=停机),
  935. 未勾选停机时仅保留运转锚点;勾选停机则停机锚点一并显示。
  936. """
  937. batch = UNIT_BATCH[unit]
  938. clauses = ["import_batch_id = %s"]
  939. params: list[Any] = [batch]
  940. if start:
  941. clauses.append("sample_time >= %s")
  942. params.append(start)
  943. if end:
  944. clauses.append("sample_time <= %s")
  945. params.append(end)
  946. with get_connection() as connection:
  947. with connection.cursor() as cursor:
  948. cursor.execute(
  949. f"""
  950. SELECT MIN(sample_time) AS anchor
  951. FROM pks_long_sample
  952. WHERE {' AND '.join(clauses)}
  953. GROUP BY FLOOR(UNIX_TIMESTAMP(sample_time) / 900)
  954. ORDER BY anchor
  955. """,
  956. params,
  957. )
  958. anchors = [row["anchor"] for row in cursor.fetchall()]
  959. if not anchors:
  960. return []
  961. if len(anchors) > PKS_STRIP_MAX_ANCHORS:
  962. step = (len(anchors) + PKS_STRIP_MAX_ANCHORS - 1) // PKS_STRIP_MAX_ANCHORS
  963. sampled = anchors[::step]
  964. if sampled[-1] != anchors[-1]:
  965. sampled.append(anchors[-1])
  966. anchors = sampled
  967. chunk_size = 1000
  968. select_expr = ", ".join(f"`{column}`" for _item, column in columns)
  969. values_by_time: dict[datetime, dict[str, float | None]] = {anchor: {} for anchor in anchors}
  970. running_by_time: dict[datetime, bool] = {}
  971. for index in range(0, len(anchors), chunk_size):
  972. chunk = anchors[index:index + chunk_size]
  973. placeholders = ", ".join(["%s"] * len(chunk))
  974. with get_connection() as connection:
  975. with connection.cursor() as cursor:
  976. cursor.execute(
  977. f"SELECT sample_time, {select_expr}, `YSJ_41` FROM pks_long_sample "
  978. f"WHERE import_batch_id = %s AND sample_time IN ({placeholders})",
  979. (batch, *chunk),
  980. )
  981. rows = cursor.fetchall()
  982. for row in rows:
  983. anchor_time = row["sample_time"]
  984. values = values_by_time[anchor_time]
  985. for item_name, column in columns:
  986. value = row[column]
  987. values[item_name] = float(value) if value is not None else None
  988. running_value = row["YSJ_41"]
  989. running_by_time[anchor_time] = bool(
  990. running_value is not None and running_value > 0
  991. )
  992. points: list[dict[str, Any]] = []
  993. for anchor in anchors:
  994. running = running_by_time.get(anchor, False)
  995. if not include_stopped and not running:
  996. continue
  997. points.append(
  998. {
  999. "sampleTime": _time_string(anchor),
  1000. "files": {},
  1001. "siteValues": values_by_time[anchor],
  1002. "machineRunning": running,
  1003. },
  1004. )
  1005. for index, point in enumerate(points):
  1006. point["index"] = index
  1007. return points
  1008. def _pks_only_wave_window(
  1009. self,
  1010. device_part: str,
  1011. points: list[dict[str, Any]],
  1012. site_points: list[str],
  1013. unit: str,
  1014. ) -> dict[str, Any]:
  1015. """PKS-only 波形窗口:不取 wave_file,仅构建所选 PKS 点位的微曲线。"""
  1016. def database_query():
  1017. return {
  1018. "points": points,
  1019. "siteSeries": self._build_site_series(site_points, unit, points),
  1020. }
  1021. def demo_query():
  1022. return {"points": points, "siteSeries": {"points": [], "series": []}}
  1023. result, source = self._run_with_fallback(database_query, demo_query)
  1024. return {
  1025. "devicePart": device_part,
  1026. "devicePoints": [],
  1027. "primaryPoint": None,
  1028. "points": points,
  1029. "xMin": 0,
  1030. "xMax": len(points),
  1031. "series": [],
  1032. "secondSeries": {
  1033. "name": "周期数据",
  1034. "color": "#f56c6c",
  1035. "sourceMeasurementType": None,
  1036. "sourceDevicePoint": None,
  1037. "data": [],
  1038. "finiteCount": 0,
  1039. "nonZeroCount": 0,
  1040. "min": None,
  1041. "max": None,
  1042. },
  1043. "angleSeries": {"color": "#d59b2b", "data": []},
  1044. "volumeSeries": {"color": "#4d9e6f", "data": [], "info": None},
  1045. "cycles": [],
  1046. "triggerXs": [],
  1047. "files": [],
  1048. "diagnostics": [],
  1049. "siteSeries": result["siteSeries"],
  1050. "source": source,
  1051. "notice": self._source_notice(source),
  1052. }
  1053. def wave_window(
  1054. self,
  1055. device_part: str,
  1056. device_points: list[str],
  1057. points: list[dict[str, Any]],
  1058. max_points: int,
  1059. no_sampling: bool = False,
  1060. first_cycle_only: bool = False,
  1061. site_points: list[str] | None = None,
  1062. ) -> dict[str, Any]:
  1063. selected_points = _validate_device_points(device_points, allow_empty=True)
  1064. if not points:
  1065. raise ValueError("至少选择一个时间点")
  1066. if len(points) > 200:
  1067. raise ValueError("单次最多预览 200 个时间点,请缩小时间窗口")
  1068. max_points = min(max(int(max_points), 256), 200000)
  1069. unit = _unit_number(device_part)
  1070. site_list = list(site_points or [])
  1071. site_columns: list[tuple[str, str]] = []
  1072. if unit:
  1073. site_columns = [
  1074. (item, column)
  1075. for item in site_list
  1076. if (column := _site_point_column(item, unit)) is not None
  1077. ]
  1078. pks_only = bool(unit and site_columns and not selected_points)
  1079. if not selected_points and not pks_only:
  1080. raise ValueError("请至少选择一个测试点位")
  1081. if pks_only:
  1082. return self._pks_only_wave_window(device_part, points, site_list, unit)
  1083. def database_query():
  1084. if first_cycle_only:
  1085. return self._build_first_cycle_window(
  1086. device_part,
  1087. selected_points,
  1088. points,
  1089. self._load_db_wave,
  1090. )
  1091. return self._build_wave_window(
  1092. device_part,
  1093. selected_points,
  1094. points,
  1095. max_points,
  1096. self._load_db_wave,
  1097. no_sampling,
  1098. )
  1099. def demo_query():
  1100. if first_cycle_only:
  1101. return self._build_first_cycle_window(
  1102. device_part,
  1103. selected_points,
  1104. points,
  1105. self._load_demo_wave,
  1106. )
  1107. return self._build_wave_window(
  1108. device_part,
  1109. selected_points,
  1110. points,
  1111. max_points,
  1112. self._load_demo_wave,
  1113. no_sampling,
  1114. )
  1115. result, source = self._run_with_fallback(database_query, demo_query)
  1116. # 全场点位(PKS)系列:每周期(PKS 5s)一小段、跨时间点连续的 15 分钟窗口曲线。
  1117. if first_cycle_only and site_points:
  1118. try:
  1119. unit = _unit_number(device_part)
  1120. result["siteSeries"] = (
  1121. self._build_site_series(site_points, unit, result["points"])
  1122. if unit is not None
  1123. else {"points": [], "series": []}
  1124. )
  1125. except Exception:
  1126. result["siteSeries"] = {"points": [], "series": []}
  1127. else:
  1128. result["siteSeries"] = {"points": [], "series": []}
  1129. result["source"] = source
  1130. result["notice"] = self._source_notice(source)
  1131. return result
  1132. def period_detail(self, wave_file_id: int, period_number: int) -> dict[str, Any]:
  1133. if wave_file_id <= 0 or period_number <= 0:
  1134. raise ValueError("wave_file_id 和周期编号必须为正整数")
  1135. if period_number != 1:
  1136. raise ValueError("当前数据仅保留第一个周期")
  1137. def database_query():
  1138. metadata, samples = self._load_db_wave(wave_file_id)
  1139. return self._build_period_detail(metadata, samples, period_number)
  1140. def demo_query():
  1141. metadata, samples = self._load_demo_wave(
  1142. wave_file_id,
  1143. "压力",
  1144. DEMO_POINT_NAME,
  1145. DEMO_START,
  1146. )
  1147. return self._build_period_detail(metadata, samples, period_number)
  1148. result, source = self._run_with_fallback(database_query, demo_query)
  1149. result["source"] = source
  1150. result["notice"] = self._source_notice(source)
  1151. return result
  1152. def annotation_config(self) -> dict[str, Any]:
  1153. return {
  1154. "source": self.source,
  1155. "annotationWidth": settings.annotation_width,
  1156. "notice": self._source_notice(self.source),
  1157. }
  1158. def list_annotations(self, wave_file_ids: list[int]) -> dict[str, Any]:
  1159. ids = sorted({int(value) for value in wave_file_ids if value})
  1160. if not ids:
  1161. return {
  1162. "source": self.source,
  1163. "annotations": [],
  1164. "notice": self._source_notice(self.source),
  1165. }
  1166. def database_query():
  1167. placeholders = ", ".join(["%s"] * len(ids))
  1168. with get_connection() as connection:
  1169. with connection.cursor() as cursor:
  1170. cursor.execute(
  1171. f"""
  1172. SELECT id, wave_file_id, label, period_start, period_end,
  1173. sample_index_start, sample_index_end
  1174. FROM wave_annotation
  1175. WHERE wave_file_id IN ({placeholders})
  1176. ORDER BY id ASC
  1177. """,
  1178. ids,
  1179. )
  1180. rows = cursor.fetchall()
  1181. return [_annotation_dict(row) for row in rows]
  1182. def demo_query():
  1183. return [
  1184. _annotation_dict(annotation)
  1185. for annotation in self._demo_annotations.values()
  1186. if annotation["wave_file_id"] in ids
  1187. ]
  1188. annotations, source = self._run_with_fallback(database_query, demo_query)
  1189. return {
  1190. "source": source,
  1191. "annotations": annotations,
  1192. "notice": self._source_notice(source),
  1193. }
  1194. def create_annotation(self, payload: dict[str, Any]) -> dict[str, Any]:
  1195. self._validate_annotation(payload)
  1196. wave_file_id = int(payload["wave_file_id"])
  1197. label = payload["label"]
  1198. period_start = int(payload["period_start"])
  1199. period_end = int(payload["period_end"])
  1200. sample_index_start = int(payload["sample_index_start"])
  1201. sample_index_end = int(payload["sample_index_end"])
  1202. def database_query():
  1203. with get_connection() as connection:
  1204. with connection.cursor() as cursor:
  1205. cursor.execute(
  1206. """
  1207. INSERT INTO wave_annotation
  1208. (wave_file_id, label, period_start, period_end,
  1209. sample_index_start, sample_index_end)
  1210. VALUES (%s, %s, %s, %s, %s, %s)
  1211. """,
  1212. (
  1213. wave_file_id,
  1214. label,
  1215. period_start,
  1216. period_end,
  1217. sample_index_start,
  1218. sample_index_end,
  1219. ),
  1220. )
  1221. annotation_id = cursor.lastrowid
  1222. cursor.execute(
  1223. """
  1224. SELECT id, wave_file_id, label, period_start, period_end,
  1225. sample_index_start, sample_index_end
  1226. FROM wave_annotation
  1227. WHERE id = %s
  1228. """,
  1229. (annotation_id,),
  1230. )
  1231. return _annotation_dict(cursor.fetchone())
  1232. def demo_query():
  1233. annotation_id = self._demo_annotation_seq
  1234. self._demo_annotation_seq += 1
  1235. annotation = {
  1236. "id": annotation_id,
  1237. "wave_file_id": wave_file_id,
  1238. "label": label,
  1239. "period_start": period_start,
  1240. "period_end": period_end,
  1241. "sample_index_start": sample_index_start,
  1242. "sample_index_end": sample_index_end,
  1243. }
  1244. self._demo_annotations[annotation_id] = annotation
  1245. return _annotation_dict(annotation)
  1246. result, source = self._run_with_fallback(database_query, demo_query)
  1247. result["source"] = source
  1248. result["notice"] = self._source_notice(source)
  1249. return result
  1250. def delete_annotation(self, annotation_id: int) -> dict[str, Any]:
  1251. if annotation_id <= 0:
  1252. raise ValueError("标注 id 必须为正整数")
  1253. def database_query():
  1254. with get_connection() as connection:
  1255. with connection.cursor() as cursor:
  1256. cursor.execute(
  1257. "DELETE FROM wave_annotation WHERE id = %s",
  1258. (annotation_id,),
  1259. )
  1260. return int(cursor.rowcount)
  1261. def demo_query():
  1262. if annotation_id not in self._demo_annotations:
  1263. return 0
  1264. del self._demo_annotations[annotation_id]
  1265. return 1
  1266. deleted, source = self._run_with_fallback(database_query, demo_query)
  1267. if not deleted:
  1268. raise ValueError(f"标注 id={annotation_id} 不存在")
  1269. return {
  1270. "deleted": annotation_id,
  1271. "source": source,
  1272. "notice": self._source_notice(source),
  1273. }
  1274. @staticmethod
  1275. def _validate_annotation(payload: dict[str, Any]) -> None:
  1276. label = payload.get("label")
  1277. if label not in ANNOTATION_LABELS:
  1278. raise ValueError("样本类型只能是 正常 或 异常")
  1279. for field in ("wave_file_id", "period_start", "period_end", "sample_index_start", "sample_index_end"):
  1280. if payload.get(field) is None:
  1281. raise ValueError(f"{field} 不能为空")
  1282. if int(payload["wave_file_id"]) <= 0:
  1283. raise ValueError("wave_file_id 必须为正整数")
  1284. if int(payload["period_start"]) <= 0 or int(payload["period_end"]) <= 0:
  1285. raise ValueError("周期编号必须为正整数")
  1286. if int(payload["period_start"]) > int(payload["period_end"]):
  1287. raise ValueError("起始周期不能大于结束周期")
  1288. if int(payload["sample_index_start"]) < 0 or int(payload["sample_index_end"]) < 0:
  1289. raise ValueError("采样点索引不能为负")
  1290. if int(payload["sample_index_start"]) > int(payload["sample_index_end"]):
  1291. raise ValueError("起始采样点不能大于结束采样点")
  1292. def health(self) -> dict[str, Any]:
  1293. return {
  1294. "status": "ok",
  1295. "source": self.source,
  1296. "databaseError": self._last_db_error or None,
  1297. }
  1298. def _source_notice(self, source: str) -> str | None:
  1299. if source == "demo":
  1300. if self._last_db_error:
  1301. return f"当前为演示数据:数据库暂不可用({self._last_db_error})"
  1302. return "当前为演示数据:可设置 DEMO_MODE=never 强制使用数据库"
  1303. return "已连接 MySQL 数据库"
  1304. def _build_wave_window(
  1305. self,
  1306. device_part: str,
  1307. device_points: list[str],
  1308. points: list[dict[str, Any]],
  1309. max_points: int,
  1310. loader: Callable[..., tuple[dict[str, Any], np.ndarray]],
  1311. no_sampling: bool = False,
  1312. ) -> dict[str, Any]:
  1313. series_data: dict[str, list[dict[str, Any]]] = {point: [] for point in device_points}
  1314. angle_data: list[dict[str, Any]] = []
  1315. volume_data: list[dict[str, Any]] = []
  1316. volume_info: dict[str, Any] | None = None
  1317. cycles: list[dict[str, Any]] = []
  1318. triggers: list[float] = []
  1319. files: list[dict[str, Any]] = []
  1320. diagnostics: list[dict[str, Any]] = []
  1321. second_series_data: list[dict[str, Any]] = []
  1322. second_finite_count = 0
  1323. second_non_zero_count = 0
  1324. second_min: float | None = None
  1325. second_max: float | None = None
  1326. primary = _primary_device_point(device_points)
  1327. primary_type = DEVICE_POINT_TO_TYPE[primary]
  1328. pressure_points = [point for point in device_points if DEVICE_POINT_TO_TYPE[point] == "压力"]
  1329. load_cache: dict[int, tuple[dict[str, Any], np.ndarray]] = {}
  1330. single_cycle = loader.__name__ != "_load_demo_wave"
  1331. for slot, point in enumerate(points):
  1332. target_per_file = max(256, int(np.ceil(max_points / max(len(points), 1))))
  1333. slot_files: dict[str, dict[str, Any]] = {}
  1334. for device_point in device_points:
  1335. file_info = point.get("files", {}).get(device_point)
  1336. if file_info:
  1337. slot_files[device_point] = file_info
  1338. # The per-slot "period source" supplies the 周期数据/体积/角度 series
  1339. # and the background cycle bands. Primary point is preferred; when a
  1340. # selected point's timestamp differs (e.g. 10:30:02 vs 10:30:05), a
  1341. # slot may only contain another device point, which is used instead
  1342. # so the period data and cycle bands do not disappear there.
  1343. source_point = primary if primary in slot_files else (next(iter(slot_files)) if slot_files else None)
  1344. source_samples: np.ndarray | None = None
  1345. source_detected: list[Any] = []
  1346. source_angle_full: np.ndarray | None = None
  1347. source_volume: np.ndarray | None = None
  1348. source_required: set[int] | None = None
  1349. source_id: int | None = None
  1350. source_type = ""
  1351. if source_point is not None:
  1352. source_info = slot_files[source_point]
  1353. source_id = int(source_info["id"] if isinstance(source_info, dict) else source_info)
  1354. source_type = DEVICE_POINT_TO_TYPE[source_point]
  1355. try:
  1356. source_meta, source_samples = self._load_for_window(
  1357. loader,
  1358. source_id,
  1359. source_type,
  1360. device_part + source_point,
  1361. point["sampleTime"],
  1362. load_cache,
  1363. )
  1364. except ValueError:
  1365. source_samples = None
  1366. if source_samples is not None and len(source_samples):
  1367. source_detected, source_diagnostic = _detect_source_cycles(source_samples, single_cycle)
  1368. source_angle_vector = build_angle_vector(len(source_samples), source_detected)
  1369. source_angle_full = np.full(len(source_samples), np.nan, dtype=float)
  1370. for detected_cycle in source_detected:
  1371. source_angle_full[detected_cycle.start_offset:detected_cycle.end_offset] = detected_cycle.angle
  1372. source_volume, current_volume_info = self._build_volume_vector(
  1373. len(source_samples),
  1374. source_detected,
  1375. device_part + primary,
  1376. )
  1377. if current_volume_info is not None:
  1378. volume_info = current_volume_info
  1379. source_indices = source_samples[:, 0].astype(np.int64)
  1380. source_span = max(len(source_samples), 1)
  1381. finite_second = source_samples[:, 2][np.isfinite(source_samples[:, 2])]
  1382. if len(finite_second):
  1383. second_finite_count += int(len(finite_second))
  1384. second_non_zero_count += int(np.count_nonzero(finite_second != 0))
  1385. current_min = float(np.min(finite_second))
  1386. current_max = float(np.max(finite_second))
  1387. second_min = current_min if second_min is None else min(second_min, current_min)
  1388. second_max = current_max if second_max is None else max(second_max, current_max)
  1389. source_required = {0, len(source_samples) - 1}
  1390. for cycle in source_detected:
  1391. source_required.update(
  1392. {
  1393. cycle.start_offset,
  1394. max(cycle.end_offset - 1, cycle.start_offset),
  1395. *cycle.trigger_offsets,
  1396. },
  1397. )
  1398. start_x = slot + cycle.start_offset / source_span
  1399. end_x = slot + cycle.end_offset / source_span
  1400. cycles.append(
  1401. {
  1402. "id": f"{source_id}:{cycle.number}",
  1403. "waveFileId": source_id,
  1404. "periodNo": cycle.number,
  1405. "pointIndex": slot,
  1406. "sampleTime": point["sampleTime"],
  1407. "startX": start_x,
  1408. "endX": end_x,
  1409. "startSampleIndex": int(source_indices[cycle.start_offset]),
  1410. "endSampleIndex": int(
  1411. source_indices[max(cycle.end_offset - 1, cycle.start_offset)],
  1412. ),
  1413. "sourceType": source_type,
  1414. "devicePoint": source_point,
  1415. "background": True,
  1416. },
  1417. )
  1418. for trigger_offset in sorted(source_required):
  1419. if any(
  1420. trigger_offset == run_offset
  1421. for cycle in source_detected
  1422. for run_offset in cycle.trigger_offsets
  1423. ):
  1424. triggers.append(slot + trigger_offset / source_span)
  1425. diagnostics.append(
  1426. {
  1427. "waveFileId": source_id,
  1428. "measurementType": source_type,
  1429. "devicePoint": source_point,
  1430. "sampleTime": point["sampleTime"],
  1431. **source_diagnostic,
  1432. },
  1433. )
  1434. source_base_index = int(source_indices[0])
  1435. second_chosen = downsample_indices(source_samples[:, 2], target_per_file, source_required)
  1436. for offset in second_chosen:
  1437. second_value = _safe_float(source_samples[offset, 2])
  1438. if second_value is None:
  1439. continue
  1440. x = slot + (int(source_indices[offset]) - source_base_index) / source_span
  1441. second_series_data.append(
  1442. {
  1443. "value": [x, second_value],
  1444. "x": x,
  1445. "rawValue": second_value,
  1446. "sampleIndex": int(source_indices[offset]),
  1447. "waveFileId": source_id,
  1448. "sampleTime": point["sampleTime"],
  1449. },
  1450. )
  1451. angle_chosen = downsample_indices(source_angle_vector, target_per_file, source_required)
  1452. for offset in angle_chosen:
  1453. value = _safe_float(source_angle_vector[offset])
  1454. if value is None:
  1455. continue
  1456. x = slot + (int(source_indices[offset]) - source_base_index) / source_span
  1457. angle_data.append(
  1458. {
  1459. "value": [x, value],
  1460. "x": x,
  1461. "angle": value,
  1462. "sampleIndex": int(source_indices[offset]),
  1463. "waveFileId": source_id,
  1464. "sampleTime": point["sampleTime"],
  1465. },
  1466. )
  1467. volume_chosen = downsample_indices(source_volume, target_per_file, source_required)
  1468. for offset in volume_chosen:
  1469. value = _safe_float(source_volume[offset])
  1470. if value is None:
  1471. continue
  1472. x = slot + (int(source_indices[offset]) - source_base_index) / source_span
  1473. volume_data.append(
  1474. {
  1475. "value": [x, value],
  1476. "x": x,
  1477. "volume": value,
  1478. "sampleIndex": int(source_indices[offset]),
  1479. "waveFileId": source_id,
  1480. "sampleTime": point["sampleTime"],
  1481. },
  1482. )
  1483. for device_point in device_points:
  1484. file_info = slot_files.get(device_point)
  1485. if not file_info:
  1486. continue
  1487. file_id = int(file_info["id"] if isinstance(file_info, dict) else file_info)
  1488. measurement_type = DEVICE_POINT_TO_TYPE[device_point]
  1489. try:
  1490. metadata, samples = self._load_for_window(
  1491. loader,
  1492. file_id,
  1493. measurement_type,
  1494. device_part + device_point,
  1495. point["sampleTime"],
  1496. load_cache,
  1497. )
  1498. except ValueError:
  1499. continue
  1500. sample_count = len(samples)
  1501. if not sample_count:
  1502. continue
  1503. sample_indices_for_file = samples[:, 0].astype(np.int64)
  1504. file_base_index = int(sample_indices_for_file[0])
  1505. file_span = max(sample_count, 1)
  1506. file_required = {0, sample_count - 1}
  1507. own_volume: np.ndarray | None = None
  1508. own_angle: np.ndarray | None = None
  1509. if device_point in pressure_points:
  1510. if source_point == device_point and source_samples is not None:
  1511. own_volume = source_volume
  1512. own_angle = source_angle_full
  1513. if source_required is not None:
  1514. file_required.update(source_required)
  1515. else:
  1516. detected_own, _ = _detect_source_cycles(samples, single_cycle)
  1517. for cycle in detected_own:
  1518. file_required.update(
  1519. {
  1520. cycle.start_offset,
  1521. max(cycle.end_offset - 1, cycle.start_offset),
  1522. *cycle.trigger_offsets,
  1523. },
  1524. )
  1525. cycles.append(
  1526. {
  1527. "id": f"{file_id}:{cycle.number}",
  1528. "waveFileId": file_id,
  1529. "periodNo": cycle.number,
  1530. "pointIndex": slot,
  1531. "sampleTime": point["sampleTime"],
  1532. "startX": slot + cycle.start_offset / file_span,
  1533. "endX": slot + cycle.end_offset / file_span,
  1534. "startSampleIndex": int(sample_indices_for_file[cycle.start_offset]),
  1535. "endSampleIndex": int(
  1536. sample_indices_for_file[max(cycle.end_offset - 1, cycle.start_offset)],
  1537. ),
  1538. "sourceType": measurement_type,
  1539. "devicePoint": device_point,
  1540. "background": False,
  1541. },
  1542. )
  1543. if detected_own:
  1544. angle_own = build_angle_vector(len(samples), detected_own)
  1545. full_own = np.full(len(samples), np.nan, dtype=float)
  1546. for detected_cycle in detected_own:
  1547. full_own[detected_cycle.start_offset:detected_cycle.end_offset] = detected_cycle.angle
  1548. own_volume, _ = self._build_volume_vector(
  1549. len(samples),
  1550. detected_own,
  1551. device_part + device_point,
  1552. )
  1553. own_angle = full_own
  1554. if no_sampling:
  1555. target_per_file_own = sample_count
  1556. else:
  1557. target_per_file_own = target_per_file
  1558. chosen = downsample_indices(samples[:, 1], target_per_file_own, file_required)
  1559. for offset in chosen:
  1560. x = slot + (int(sample_indices_for_file[offset]) - file_base_index) / file_span
  1561. raw_value = float(samples[offset, 1])
  1562. series_data[device_point].append(
  1563. {
  1564. "value": [x, raw_value],
  1565. "x": x,
  1566. "rawValue": raw_value,
  1567. "sampleIndex": int(sample_indices_for_file[offset]),
  1568. "waveFileId": file_id,
  1569. "sampleTime": point["sampleTime"],
  1570. "secondValue": _safe_float(samples[offset, 2]),
  1571. "volume": _safe_float(own_volume[offset]) if own_volume is not None else None,
  1572. "angle360": _safe_float(own_angle[offset]) if own_angle is not None else None,
  1573. },
  1574. )
  1575. files.append(
  1576. {
  1577. "id": file_id,
  1578. "pointIndex": slot,
  1579. "sampleTime": point["sampleTime"],
  1580. "devicePoint": device_point,
  1581. "measurementType": measurement_type,
  1582. "sampleCount": int(metadata.get("sample_count") or sample_count),
  1583. "sampleFrequencyHz": int(metadata.get("sample_frequency_hz") or 0),
  1584. "pointName": str(metadata.get("point_name") or device_part + device_point),
  1585. "rpm": float(metadata.get("rpm") or 0),
  1586. "status": int(metadata.get("tspluse_status") or 0),
  1587. "fileName": str(metadata.get("file_name") or ""),
  1588. },
  1589. )
  1590. load_cache.clear()
  1591. for device_point in series_data:
  1592. series_data[device_point].sort(key=lambda item: item["x"])
  1593. second_series_data.sort(key=lambda item: item["x"])
  1594. angle_data.sort(key=lambda item: item["x"])
  1595. volume_data.sort(key=lambda item: item["x"])
  1596. extents = {device_point: _series_extent(series_data[device_point]) for device_point in device_points}
  1597. return {
  1598. "devicePart": device_part,
  1599. "devicePoints": device_points,
  1600. "primaryPoint": primary,
  1601. "points": points,
  1602. "xMin": 0,
  1603. "xMax": len(points),
  1604. "series": [
  1605. {
  1606. "devicePoint": device_point,
  1607. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  1608. "color": MEASUREMENT_COLORS[DEVICE_POINT_TO_TYPE[device_point]],
  1609. "data": series_data[device_point],
  1610. "min": extents[device_point][0],
  1611. "max": extents[device_point][1],
  1612. }
  1613. for device_point in device_points
  1614. ],
  1615. "secondSeries": {
  1616. "name": "周期数据",
  1617. "color": "#f56c6c",
  1618. "sourceMeasurementType": primary_type,
  1619. "sourceDevicePoint": primary,
  1620. "data": second_series_data,
  1621. "finiteCount": second_finite_count,
  1622. "nonZeroCount": second_non_zero_count,
  1623. "min": second_min,
  1624. "max": second_max,
  1625. },
  1626. "angleSeries": {
  1627. "color": "#d59b2b",
  1628. "data": angle_data,
  1629. },
  1630. "volumeSeries": {
  1631. "color": "#4d9e6f",
  1632. "data": volume_data,
  1633. "info": volume_info,
  1634. },
  1635. "cycles": cycles,
  1636. "triggerXs": sorted(set(triggers)),
  1637. "files": files,
  1638. "diagnostics": diagnostics,
  1639. }
  1640. def _build_first_cycle_window(
  1641. self,
  1642. device_part: str,
  1643. device_points: list[str],
  1644. points: list[dict[str, Any]],
  1645. loader: Callable[..., tuple[dict[str, Any], np.ndarray]],
  1646. ) -> dict[str, Any]:
  1647. """Build a slot-layout window where each file shows only its first cycle.
  1648. The primary point is selected in pressure-cap, pressure-shaft, then
  1649. selected-point order. It is the point used for the first-cycle slice;
  1650. files without recorded bounds still appear in the file list.
  1651. The slice for every primary file is fetched in a single JOIN query (range
  1652. scan on the ``(wave_file_id, sample_index)`` primary key). Other selected
  1653. device points contribute no curve but still appear in the file list. The
  1654. curve is continuous: all cycle samples are kept and each file spans its
  1655. own x slot.
  1656. """
  1657. primary = _primary_device_point(device_points)
  1658. primary_type = DEVICE_POINT_TO_TYPE[primary]
  1659. series_data: dict[str, list[dict[str, Any]]] = {point: [] for point in device_points}
  1660. second_series_data: list[dict[str, Any]] = []
  1661. angle_data: list[dict[str, Any]] = []
  1662. volume_data: list[dict[str, Any]] = []
  1663. volume_info: dict[str, Any] | None = None
  1664. cycles: list[dict[str, Any]] = []
  1665. triggers: list[float] = []
  1666. files: list[dict[str, Any]] = []
  1667. diagnostics: list[dict[str, Any]] = []
  1668. second_finite_count = 0
  1669. second_non_zero_count = 0
  1670. second_min: float | None = None
  1671. second_max: float | None = None
  1672. slot_sources: list[tuple[int, dict[str, Any], str, int]] = []
  1673. all_file_ids: set[int] = set()
  1674. for slot, point in enumerate(points):
  1675. for device_point in device_points:
  1676. file_info = point.get("files", {}).get(device_point)
  1677. if not file_info:
  1678. continue
  1679. file_id = int(file_info["id"] if isinstance(file_info, dict) else file_info)
  1680. slot_sources.append((slot, point, device_point, file_id))
  1681. all_file_ids.add(file_id)
  1682. if not slot_sources:
  1683. return self._assemble_first_cycle_window(
  1684. device_part, device_points, primary, points, series_data, second_series_data,
  1685. angle_data, volume_data, volume_info,
  1686. second_finite_count, second_non_zero_count, second_min, second_max,
  1687. cycles, triggers, files, diagnostics, 0,
  1688. )
  1689. is_demo = loader.__name__ == "_load_demo_wave"
  1690. metas: dict[int, dict[str, Any]] = {}
  1691. # file_id -> (cycle_start, cycle_end, padded_slice[offset, signal, second])
  1692. slices: dict[int, tuple[int, int, np.ndarray]] = {}
  1693. if is_demo:
  1694. for _slot, point, device_point, source_id in slot_sources:
  1695. metadata, samples = self._load_for_window(
  1696. loader, source_id, DEVICE_POINT_TO_TYPE[device_point],
  1697. device_part + device_point, point["sampleTime"], {},
  1698. )
  1699. metas[source_id] = metadata
  1700. detected, _ = detect_cycles(samples)
  1701. if detected:
  1702. cycle = detected[0]
  1703. start_si = int(samples[cycle.start_offset, 0])
  1704. end_si = int(samples[cycle.end_offset, 0])
  1705. pad_end = min(cycle.end_offset + FIRST_CYCLE_PAD, len(samples))
  1706. slices[source_id] = (start_si, end_si, samples[cycle.start_offset:pad_end])
  1707. else:
  1708. source_ids = [source_id for _, _, _, source_id in slot_sources]
  1709. with get_connection() as connection:
  1710. with connection.cursor() as cursor:
  1711. placeholders = ", ".join(["%s"] * len(all_file_ids))
  1712. cursor.execute(
  1713. f"""
  1714. SELECT id, point_name, measurement_type, sample_time,
  1715. sample_count, sample_frequency_hz, rpm, tspluse_status,
  1716. file_name, cycle_start, cycle_end
  1717. FROM wave_file
  1718. WHERE id IN ({placeholders})
  1719. """,
  1720. tuple(all_file_ids),
  1721. )
  1722. for row in cursor.fetchall():
  1723. metas[int(row["id"])] = row
  1724. source_placeholders = ", ".join(["%s"] * len(source_ids))
  1725. cursor.execute(
  1726. f"""
  1727. SELECT ws.wave_file_id,
  1728. ws.sample_index,
  1729. CAST(ws.signal_value AS FLOAT) AS sig,
  1730. CAST(ws.second_value AS FLOAT) AS sec
  1731. FROM wave_sample_one ws
  1732. WHERE ws.wave_file_id IN ({source_placeholders})
  1733. ORDER BY ws.wave_file_id, ws.sample_index
  1734. """,
  1735. tuple(source_ids),
  1736. )
  1737. raw: dict[int, list[tuple[float, float, float]]] = {}
  1738. for row in cursor.fetchall():
  1739. raw.setdefault(int(row["wave_file_id"]), []).append(
  1740. (
  1741. float(row["sample_index"]),
  1742. float(row["sig"]),
  1743. float(row["sec"]) if row["sec"] is not None else float("nan"),
  1744. ),
  1745. )
  1746. for file_id, rows in raw.items():
  1747. meta = metas.get(file_id)
  1748. if meta is None or not rows:
  1749. continue
  1750. # wave_sample_one 只保留第一个周期,返回的数组本身就是该周期。
  1751. slices[file_id] = (
  1752. int(rows[0][0]),
  1753. int(rows[-1][0]) + 1,
  1754. np.asarray(rows, dtype=float),
  1755. )
  1756. # Per-slot period source: primary preferred, else the first selected
  1757. # point with a file at that slot (timestamps may differ by seconds).
  1758. source_by_slot: dict[int, str] = {}
  1759. for _slot, _point, device_point, _file_id in slot_sources:
  1760. if _slot not in source_by_slot:
  1761. source_by_slot[_slot] = device_point
  1762. for _slot, _point, device_point, _file_id in slot_sources:
  1763. if device_point == primary:
  1764. source_by_slot[_slot] = primary
  1765. for slot, point, device_point, source_id in slot_sources:
  1766. sample_time = point["sampleTime"]
  1767. is_slot_source = device_point == source_by_slot.get(slot, device_point)
  1768. cycle_bounds = slices.get(source_id)
  1769. cycle_start = 0
  1770. cycle_end = 0
  1771. slice_arr: np.ndarray | None = None
  1772. if cycle_bounds is not None:
  1773. cycle_start, cycle_end, slice_arr = cycle_bounds
  1774. cycle_len = len(slice_arr) if slice_arr is not None else 0
  1775. has_cycle = slice_arr is not None and 1 < cycle_len
  1776. if has_cycle:
  1777. angle: np.ndarray | None = None
  1778. volume: np.ndarray | None = None
  1779. point_volume_info: dict[str, Any] | None = None
  1780. stored_cycle = _stored_first_cycle(slice_arr)
  1781. detected = [stored_cycle] if stored_cycle is not None else []
  1782. if detected:
  1783. cycle = detected[0]
  1784. angle = cycle.angle
  1785. volume, point_volume_info = self._build_volume_vector(
  1786. len(slice_arr),
  1787. [cycle],
  1788. device_part + primary,
  1789. )
  1790. if is_slot_source and point_volume_info is not None:
  1791. volume_info = point_volume_info
  1792. for offset in range(cycle_len):
  1793. sample_index = int(slice_arr[offset, 0])
  1794. raw_value = float(slice_arr[offset, 1])
  1795. second = _safe_float(slice_arr[offset, 2])
  1796. angle360 = _safe_float(angle[offset]) if angle is not None and offset < len(angle) else None
  1797. volume_value = _safe_float(volume[offset]) if volume is not None and offset < len(volume) else None
  1798. x = slot + (sample_index - cycle_start) / cycle_len
  1799. series_data[device_point].append(
  1800. {
  1801. "value": [x, raw_value],
  1802. "x": x,
  1803. "rawValue": raw_value,
  1804. "sampleIndex": sample_index,
  1805. "waveFileId": source_id,
  1806. "sampleTime": sample_time,
  1807. "secondValue": second,
  1808. "volume": volume_value,
  1809. "angle360": angle360,
  1810. },
  1811. )
  1812. if is_slot_source:
  1813. if volume_value is not None:
  1814. volume_data.append(
  1815. {
  1816. "value": [x, volume_value],
  1817. "x": x,
  1818. "volume": volume_value,
  1819. "sampleIndex": sample_index,
  1820. "waveFileId": source_id,
  1821. "sampleTime": sample_time,
  1822. },
  1823. )
  1824. if angle360 is not None:
  1825. display_angle = angle360 if angle360 <= 180.0 else 360.0 - angle360
  1826. angle_data.append(
  1827. {
  1828. "value": [x, display_angle],
  1829. "x": x,
  1830. "angle": display_angle,
  1831. "sampleIndex": sample_index,
  1832. "waveFileId": source_id,
  1833. "sampleTime": sample_time,
  1834. },
  1835. )
  1836. if second is not None:
  1837. second_finite_count += 1
  1838. if second != 0:
  1839. second_non_zero_count += 1
  1840. second_min = second if second_min is None else min(second_min, second)
  1841. second_max = second if second_max is None else max(second_max, second)
  1842. second_series_data.append(
  1843. {
  1844. "value": [x, second],
  1845. "x": x,
  1846. "rawValue": second,
  1847. "sampleIndex": sample_index,
  1848. "waveFileId": source_id,
  1849. "sampleTime": sample_time,
  1850. },
  1851. )
  1852. if offset > 0:
  1853. prev = _safe_float(slice_arr[offset - 1, 2])
  1854. if second is not None and (prev is None or prev < 30) and second >= 30:
  1855. triggers.append(x)
  1856. cycles.append(
  1857. {
  1858. "id": f"{source_id}:1",
  1859. "waveFileId": source_id,
  1860. "periodNo": 1,
  1861. "pointIndex": slot,
  1862. "sampleTime": sample_time,
  1863. "startX": float(slot),
  1864. "endX": float(slot + 1),
  1865. "startSampleIndex": cycle_start,
  1866. "endSampleIndex": cycle_end - 1,
  1867. "sourceType": DEVICE_POINT_TO_TYPE[device_point],
  1868. "devicePoint": device_point,
  1869. "background": is_slot_source,
  1870. },
  1871. )
  1872. diagnostics.append(
  1873. {
  1874. "waveFileId": source_id,
  1875. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  1876. "devicePoint": device_point,
  1877. "sampleTime": sample_time,
  1878. "cycleStart": cycle_start,
  1879. "cycleEnd": cycle_end,
  1880. "cycleSampleCount": cycle_len,
  1881. },
  1882. )
  1883. for slot, point in enumerate(points):
  1884. for device_point in device_points:
  1885. file_info = point.get("files", {}).get(device_point)
  1886. if not file_info:
  1887. continue
  1888. file_id = int(file_info["id"] if isinstance(file_info, dict) else file_info)
  1889. metadata = metas.get(file_id)
  1890. cycle_start = metadata.get("cycle_start") if metadata else None
  1891. cycle_end = metadata.get("cycle_end") if metadata else None
  1892. files.append(
  1893. {
  1894. "id": file_id,
  1895. "pointIndex": slot,
  1896. "sampleTime": point["sampleTime"],
  1897. "devicePoint": device_point,
  1898. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  1899. "sampleCount": int(
  1900. (metadata.get("sample_count") if metadata else file_info.get("sampleCount") or 0) or 0,
  1901. ),
  1902. "sampleFrequencyHz": int(
  1903. (metadata.get("sample_frequency_hz") if metadata else file_info.get("sampleFrequencyHz") or 0) or 0,
  1904. ),
  1905. "pointName": str(
  1906. (metadata.get("point_name") if metadata else device_part + device_point)
  1907. or device_part + device_point,
  1908. ),
  1909. "rpm": float((metadata.get("rpm") if metadata else file_info.get("rpm") or 0) or 0),
  1910. "status": int(
  1911. (metadata.get("tspluse_status") if metadata else file_info.get("status") or 0) or 0,
  1912. ),
  1913. "fileName": str((metadata.get("file_name") if metadata else "") or ""),
  1914. "cycleStart": int(cycle_start) if cycle_start is not None else None,
  1915. "cycleEnd": int(cycle_end) if cycle_end is not None else None,
  1916. },
  1917. )
  1918. return self._assemble_first_cycle_window(
  1919. device_part, device_points, primary, points, series_data, second_series_data,
  1920. angle_data, volume_data, volume_info,
  1921. second_finite_count, second_non_zero_count, second_min, second_max,
  1922. cycles, triggers, files, diagnostics,
  1923. max(len(slot_sources) - len(cycles), 0),
  1924. )
  1925. @staticmethod
  1926. def _assemble_first_cycle_window(
  1927. device_part: str,
  1928. device_points: list[str],
  1929. primary: str,
  1930. points: list[dict[str, Any]],
  1931. series_data: dict[str, list[dict[str, Any]]],
  1932. second_series_data: list[dict[str, Any]],
  1933. angle_data: list[dict[str, Any]],
  1934. volume_data: list[dict[str, Any]],
  1935. volume_info: dict[str, Any] | None,
  1936. second_finite_count: int,
  1937. second_non_zero_count: int,
  1938. second_min: float | None,
  1939. second_max: float | None,
  1940. cycles: list[dict[str, Any]],
  1941. triggers: list[float],
  1942. files: list[dict[str, Any]],
  1943. diagnostics: list[dict[str, Any]],
  1944. missing_cycles: int = 0,
  1945. ) -> dict[str, Any]:
  1946. for device_point in series_data:
  1947. series_data[device_point].sort(key=lambda item: item["x"])
  1948. second_series_data.sort(key=lambda item: item["x"])
  1949. angle_data.sort(key=lambda item: item["x"])
  1950. volume_data.sort(key=lambda item: item["x"])
  1951. extents = {device_point: _series_extent(series_data[device_point]) for device_point in device_points}
  1952. return {
  1953. "devicePart": device_part,
  1954. "devicePoints": device_points,
  1955. "primaryPoint": primary,
  1956. "points": points,
  1957. "xMin": 0,
  1958. "xMax": len(points),
  1959. "series": [
  1960. {
  1961. "devicePoint": device_point,
  1962. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  1963. "color": MEASUREMENT_COLORS[DEVICE_POINT_TO_TYPE[device_point]],
  1964. "data": series_data[device_point],
  1965. "min": extents[device_point][0],
  1966. "max": extents[device_point][1],
  1967. }
  1968. for device_point in device_points
  1969. ],
  1970. "secondSeries": {
  1971. "name": "周期数据",
  1972. "color": "#f56c6c",
  1973. "sourceMeasurementType": DEVICE_POINT_TO_TYPE[primary],
  1974. "sourceDevicePoint": primary,
  1975. "data": second_series_data,
  1976. "finiteCount": second_finite_count,
  1977. "nonZeroCount": second_non_zero_count,
  1978. "min": second_min,
  1979. "max": second_max,
  1980. },
  1981. "angleSeries": {
  1982. "color": "#d59b2b",
  1983. "data": angle_data,
  1984. },
  1985. "volumeSeries": {
  1986. "color": "#4d9e6f",
  1987. "data": volume_data,
  1988. "info": volume_info,
  1989. },
  1990. "cycles": cycles,
  1991. "triggerXs": sorted(set(triggers)),
  1992. "files": files,
  1993. "diagnostics": diagnostics,
  1994. "firstCycleMode": True,
  1995. "firstCycleNotice": (
  1996. f"窗口内 {missing_cycles} 个文件没有可用的单周期采样数据(wave_sample_one),"
  1997. "多为停机(rpm=0)文件,无法绘制首周期曲线。"
  1998. if missing_cycles > 0
  1999. else None
  2000. ),
  2001. }
  2002. @staticmethod
  2003. def _build_volume_vector(
  2004. sample_count: int,
  2005. cycles: list[DetectedCycle],
  2006. point_name: str,
  2007. ) -> tuple[np.ndarray, dict[str, Any] | None]:
  2008. volume = np.full(sample_count, np.nan, dtype=float)
  2009. cylinder_name = next(
  2010. (name for name in CYLINDER_BORE_MM if name in point_name),
  2011. None,
  2012. )
  2013. if cylinder_name is None or not cycles:
  2014. return volume, None
  2015. bore_mm = CYLINDER_BORE_MM[cylinder_name]
  2016. clearance = CLEARANCE_VOLUME_L_BY_BORE[bore_mm]
  2017. crank_radius = PISTON_STROKE_MM / 2.0
  2018. piston_area = np.pi * (bore_mm / 2.0) ** 2
  2019. for cycle in cycles:
  2020. angle_rad = np.deg2rad(cycle.angle)
  2021. travel = (
  2022. crank_radius * (1.0 - np.cos(angle_rad))
  2023. + CONNECTING_ROD_LENGTH_MM
  2024. - np.sqrt(
  2025. CONNECTING_ROD_LENGTH_MM**2
  2026. - (crank_radius * np.sin(angle_rad)) ** 2,
  2027. )
  2028. )
  2029. volume[cycle.start_offset : cycle.end_offset] = (
  2030. clearance + piston_area * travel / 1_000_000.0
  2031. )
  2032. finite = volume[np.isfinite(volume)]
  2033. return volume, {
  2034. "cylinder": cylinder_name,
  2035. "boreMm": bore_mm,
  2036. "clearanceVolumeL": clearance,
  2037. "minVolumeL": float(np.min(finite)) if len(finite) else None,
  2038. "maxVolumeL": float(np.max(finite)) if len(finite) else None,
  2039. }
  2040. @staticmethod
  2041. def _build_period_detail(
  2042. metadata: dict[str, Any],
  2043. samples: np.ndarray,
  2044. period_number: int,
  2045. ) -> dict[str, Any]:
  2046. stored_cycle = _stored_first_cycle(samples) if metadata.get("single_cycle") else None
  2047. if stored_cycle is not None:
  2048. detected = [stored_cycle]
  2049. diagnostics = {"completeCycleCount": 1, "storedFirstCycle": 1}
  2050. else:
  2051. detected, diagnostics = detect_cycles(samples)
  2052. cycle = next((item for item in detected if item.number == period_number), None)
  2053. if cycle is None:
  2054. raise ValueError(f"没有找到周期 {period_number}")
  2055. angles360 = np.linspace(0.0, 359.0, 360)
  2056. pressure = np.interp(angles360, cycle.angle, cycle.signal)
  2057. display_angles = np.where(angles360 <= 180.0, angles360, 360.0 - angles360)
  2058. point_name = str(metadata.get("point_name") or "")
  2059. cylinder_name = next(
  2060. (name for name in CYLINDER_BORE_MM if name in point_name),
  2061. None,
  2062. )
  2063. volume: np.ndarray | None = None
  2064. volume_info: dict[str, Any] | None = None
  2065. if cylinder_name:
  2066. bore_mm = CYLINDER_BORE_MM[cylinder_name]
  2067. clearance_volume = CLEARANCE_VOLUME_L_BY_BORE[bore_mm]
  2068. angle_rad = np.deg2rad(angles360)
  2069. crank_radius = PISTON_STROKE_MM / 2.0
  2070. piston_travel = (
  2071. crank_radius * (1.0 - np.cos(angle_rad))
  2072. + CONNECTING_ROD_LENGTH_MM
  2073. - np.sqrt(
  2074. CONNECTING_ROD_LENGTH_MM**2
  2075. - (crank_radius * np.sin(angle_rad)) ** 2,
  2076. )
  2077. )
  2078. piston_area = np.pi * (bore_mm / 2.0) ** 2
  2079. volume = clearance_volume + piston_area * piston_travel / 1_000_000.0
  2080. volume_info = {
  2081. "cylinder": cylinder_name,
  2082. "boreMm": bore_mm,
  2083. "clearanceVolumeL": clearance_volume,
  2084. "minVolumeL": float(np.min(volume)),
  2085. "maxVolumeL": float(np.max(volume)),
  2086. }
  2087. start_index = int(samples[cycle.start_offset, 0])
  2088. end_offset = min(cycle.end_offset, len(samples) - 1)
  2089. end_index = int(samples[max(cycle.end_offset - 1, cycle.start_offset), 0])
  2090. return {
  2091. "waveFile": {
  2092. "id": int(metadata["id"]),
  2093. "pointName": point_name,
  2094. "measurementType": metadata.get("measurement_type"),
  2095. "sampleTime": _time_string(metadata.get("sample_time")),
  2096. "sampleFrequencyHz": int(metadata.get("sample_frequency_hz") or 0),
  2097. "sampleCount": int(metadata.get("sample_count") or len(samples)),
  2098. },
  2099. "period": {
  2100. "periodNo": cycle.number,
  2101. "startSampleIndex": start_index,
  2102. "endSampleIndex": end_index,
  2103. "sampleCount": int(cycle.end_offset - cycle.start_offset),
  2104. "triggerSampleIndices": [
  2105. int(samples[offset, 0])
  2106. for offset in cycle.trigger_offsets
  2107. if 0 <= offset < len(samples)
  2108. ],
  2109. },
  2110. "angles": display_angles.tolist(),
  2111. "angles360": angles360.tolist(),
  2112. "pressure": pressure.tolist(),
  2113. "volume": volume.tolist() if volume is not None else None,
  2114. "volumeInfo": volume_info,
  2115. "phases": [
  2116. {
  2117. "name": name,
  2118. "color": color,
  2119. "start": start,
  2120. "end": end,
  2121. }
  2122. for name, color, start, end in PHASES
  2123. ],
  2124. "diagnostics": diagnostics,
  2125. }
  2126. @staticmethod
  2127. def _load_for_window(
  2128. loader: Callable[..., tuple[dict[str, Any], np.ndarray]],
  2129. file_id: int,
  2130. measurement_type: str,
  2131. point_name: str,
  2132. sample_time: str,
  2133. load_cache: dict[int, tuple[dict[str, Any], np.ndarray]],
  2134. ) -> tuple[dict[str, Any], np.ndarray]:
  2135. if file_id not in load_cache:
  2136. if loader.__name__ == "_load_demo_wave":
  2137. load_cache[file_id] = loader(file_id, measurement_type, point_name, sample_time)
  2138. else:
  2139. load_cache[file_id] = loader(file_id)
  2140. return load_cache[file_id]
  2141. def _load_db_wave(self, file_id: int) -> tuple[dict[str, Any], np.ndarray]:
  2142. with get_connection() as connection:
  2143. with connection.cursor() as cursor:
  2144. cursor.execute(
  2145. """
  2146. SELECT id, point_name, measurement_type, sample_frequency_hz,
  2147. sample_count, sample_time, rpm, file_name, tspluse_status,
  2148. cycle_start, cycle_end
  2149. FROM wave_file
  2150. WHERE id = %s
  2151. """,
  2152. (file_id,),
  2153. )
  2154. metadata = cursor.fetchone()
  2155. if metadata is None:
  2156. raise ValueError(f"wave_file.id={file_id} 不存在")
  2157. # wave_sample_one 只保留第一个周期,直接读取该文件全部样本。
  2158. cursor.execute(
  2159. """
  2160. SELECT sample_index, signal_value, second_value
  2161. FROM wave_sample_one
  2162. WHERE wave_file_id = %s
  2163. ORDER BY sample_index ASC
  2164. """,
  2165. (file_id,),
  2166. )
  2167. rows = cursor.fetchall()
  2168. if not rows:
  2169. raise ValueError(f"wave_file.id={file_id} 没有采样数据")
  2170. metadata["single_cycle"] = True
  2171. samples = np.asarray(
  2172. [
  2173. (
  2174. float(row["sample_index"]),
  2175. float(row["signal_value"]),
  2176. float(row["second_value"]) if row["second_value"] is not None else np.nan,
  2177. )
  2178. for row in rows
  2179. ],
  2180. dtype=float,
  2181. )
  2182. return metadata, samples
  2183. @staticmethod
  2184. @lru_cache(maxsize=24)
  2185. def _demo_samples(file_id: int, measurement_type: str) -> np.ndarray:
  2186. count = DEMO_SAMPLE_COUNT
  2187. index = np.arange(count, dtype=float)
  2188. revolution = DEMO_REVOLUTION_SAMPLES
  2189. phase = (index % revolution) / revolution * 2 * np.pi
  2190. second = np.zeros(count, dtype=float)
  2191. for revolution_start in range(0, count, revolution):
  2192. for pulse in range(PULSES_PER_REVOLUTION):
  2193. pulse_start = revolution_start + int(round(pulse * revolution / PULSES_PER_REVOLUTION))
  2194. width = 22 if pulse == 0 else 8
  2195. pulse_end = min(count, pulse_start + width)
  2196. second[pulse_start:pulse_end] = 40.0
  2197. variation = (file_id % 17) / 17.0
  2198. if measurement_type == "压力":
  2199. signal = (
  2200. 4.2
  2201. + 1.8 * np.sin(phase - 0.4)
  2202. + 0.55 * np.sin(2 * phase + variation)
  2203. + 0.22 * np.sin(7 * phase)
  2204. )
  2205. signal += 0.2 * np.maximum(np.sin(phase - 0.2), 0) ** 5
  2206. elif measurement_type == "位移":
  2207. signal = 0.5 + 0.18 * np.cos(phase) + 0.035 * np.sin(3 * phase + variation)
  2208. else:
  2209. signal = 0.15 * np.sin(phase * 2 + variation) + 0.04 * np.sin(11 * phase)
  2210. signal += 0.018 * np.cos(index / 37.0)
  2211. return np.column_stack((index, signal, second))
  2212. def _load_demo_wave(
  2213. self,
  2214. file_id: int,
  2215. measurement_type: str,
  2216. point_name: str,
  2217. sample_time: str,
  2218. ) -> tuple[dict[str, Any], np.ndarray]:
  2219. samples = self._demo_samples(file_id, measurement_type)
  2220. metadata = {
  2221. "id": file_id,
  2222. "point_name": point_name,
  2223. "measurement_type": measurement_type,
  2224. "sample_frequency_hz": 25600,
  2225. "sample_count": len(samples),
  2226. "sample_time": sample_time,
  2227. "rpm": 998.0,
  2228. "file_name": f"demo-{file_id}.dat",
  2229. }
  2230. return metadata, samples
  2231. @staticmethod
  2232. def _demo_options() -> list[dict[str, Any]]:
  2233. end = DEMO_START + timedelta(minutes=5 * (DEMO_POINT_COUNT - 1))
  2234. return [
  2235. {
  2236. "devicePart": DEMO_DEVICE_PART,
  2237. "devicePoint": device_point,
  2238. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  2239. "minTime": _time_string(DEMO_START),
  2240. "maxTime": _time_string(end),
  2241. "fileCount": DEMO_POINT_COUNT,
  2242. }
  2243. for device_point in DEVICE_POINTS
  2244. ]
  2245. @staticmethod
  2246. def _demo_time_points(
  2247. device_part: str,
  2248. device_points: list[str],
  2249. start: datetime | None,
  2250. end: datetime | None,
  2251. ) -> list[dict[str, Any]]:
  2252. points = []
  2253. for index in range(DEMO_POINT_COUNT):
  2254. timestamp = DEMO_START + timedelta(minutes=5 * index)
  2255. if start and timestamp < start:
  2256. continue
  2257. if end and timestamp > end:
  2258. continue
  2259. files = {}
  2260. for device_point in device_points:
  2261. measurement_type = DEVICE_POINT_TO_TYPE[device_point]
  2262. files[device_point] = {
  2263. "id": DEMO_ID_BY_TYPE[measurement_type] + index,
  2264. "devicePoint": device_point,
  2265. "measurementType": measurement_type,
  2266. "sampleCount": DEMO_SAMPLE_COUNT,
  2267. "sampleFrequencyHz": 25600,
  2268. "rpm": 998.0,
  2269. }
  2270. points.append(
  2271. {
  2272. "index": len(points),
  2273. "sampleTime": _time_string(timestamp),
  2274. "files": files,
  2275. },
  2276. )
  2277. return points
  2278. data_service = DataService()