data_service.py 140 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681268226832684268526862687268826892690269126922693269426952696269726982699270027012702270327042705270627072708270927102711271227132714271527162717271827192720272127222723272427252726272727282729273027312732273327342735273627372738273927402741274227432744274527462747274827492750275127522753275427552756275727582759276027612762276327642765276627672768276927702771277227732774277527762777277827792780278127822783278427852786278727882789279027912792279327942795279627972798279928002801280228032804280528062807280828092810281128122813281428152816281728182819282028212822282328242825282628272828282928302831283228332834283528362837283828392840284128422843284428452846284728482849285028512852285328542855285628572858285928602861286228632864286528662867286828692870287128722873287428752876287728782879288028812882288328842885288628872888288928902891289228932894289528962897289828992900290129022903290429052906290729082909291029112912291329142915291629172918291929202921292229232924292529262927292829292930293129322933293429352936293729382939294029412942294329442945294629472948294929502951295229532954295529562957295829592960296129622963296429652966296729682969297029712972297329742975297629772978297929802981298229832984298529862987298829892990299129922993299429952996299729982999300030013002300330043005300630073008300930103011301230133014301530163017301830193020302130223023302430253026302730283029303030313032303330343035303630373038303930403041304230433044304530463047304830493050305130523053305430553056305730583059306030613062306330643065306630673068306930703071307230733074307530763077307830793080308130823083308430853086308730883089309030913092309330943095309630973098309931003101310231033104310531063107310831093110311131123113311431153116311731183119312031213122312331243125312631273128312931303131
  1. from __future__ import annotations
  2. import re
  3. import json
  4. from bisect import bisect_left
  5. from collections import OrderedDict
  6. from datetime import datetime, timedelta
  7. from functools import lru_cache
  8. from time import monotonic
  9. from typing import Any, Callable
  10. import numpy as np
  11. from ..algorithms.cycles import (
  12. DetectedCycle,
  13. PULSES_PER_REVOLUTION,
  14. build_angle_vector,
  15. detect_cycles,
  16. downsample_indices,
  17. )
  18. from ..config import settings
  19. from ..db import get_connection
  20. MEASUREMENT_TYPES = ("位移", "加速度", "压力")
  21. MEASUREMENT_COLORS = {
  22. "压力": "#e4572e",
  23. "位移": "#1f7a8c",
  24. "加速度": "#7b61a8",
  25. }
  26. DEVICE_POINTS = ("压力盖侧", "压力轴侧", "活塞杆沉降", "十字头振动", "自由端振动", "驱动端振动")
  27. DEVICE_POINT_TO_TYPE = {
  28. "压力盖侧": "压力",
  29. "压力轴侧": "压力",
  30. "活塞杆沉降": "位移",
  31. "十字头振动": "加速度",
  32. "自由端振动": "加速度",
  33. "驱动端振动": "加速度",
  34. }
  35. PRIMARY_DEVICE_POINT = "压力盖侧"
  36. PHASES = (
  37. ("排气", "#7b1fa2", 0.0, 120.0),
  38. ("压缩", "#c62828", 120.0, 195.0),
  39. ("膨胀", "#1565c0", 195.0, 270.0),
  40. ("进气", "#2e7d32", 270.0, 330.0),
  41. )
  42. ANNOTATION_LABELS = ("正常", "异常")
  43. OIL_PRESSURE_ALARM_TYPE = "润滑油压力低"
  44. CRUCIFORM_FAULT_ALARM_TYPE = "十字头故障"
  45. VALVE_FAULT_ALARM_TYPE = "气阀故障"
  46. OIL_PRESSURE_DEVICE_PART = "润滑油"
  47. OIL_PRESSURE_DEVICE_POINT = "压力"
  48. VIBRATION_DEVICE_POINT = "振动"
  49. OIL_PRESSURE_SCAN_DAYS = 14
  50. # PKS 全场点位:机组号 -> pks_long_sample.import_batch_id
  51. # 7号机=30、8号机=31、9号机=32。
  52. UNIT_BATCH = {"7": 30, "8": 31, "9": 32}
  53. # 全场点位(PKS)每个时间点的曲线窗口:以采样时刻为中心的前后各 7.5 分钟。
  54. PKS_WINDOW_SECONDS = 15 * 60
  55. # 仅 PKS 时间点返回给时间条的最大锚点数:超过后按等步长均匀采样(首尾必保)。
  56. PKS_STRIP_MAX_ANCHORS = 5000
  57. _SITE_POINT_PATTERN = re.compile(r"^YSJ([789])_([1-9]|1[0-9]|2[0-9]|3[0-9]|4[0-1])$")
  58. # 不作为全场点位(PKS)下拉候选的序号(对应机组自己的运行状态/转速信号)。
  59. _EXCLUDED_SITE_INDEXES = frozenset({33, 34, 35, 41})
  60. # Extra samples fetched past cycle_end so detect_cycles can see the zero marker
  61. # that closes the first cycle (cycle_end is exclusive, the marker sits at it).
  62. FIRST_CYCLE_PAD = 64
  63. CYLINDER_BORE_MM = {
  64. "一缸": 360.0,
  65. "二缸": 490.0,
  66. "三缸": 390.0,
  67. "四缸": 490.0,
  68. "五缸": 390.0,
  69. "六缸": 490.0,
  70. }
  71. PISTON_STROKE_MM = 148.0
  72. CONNECTING_ROD_LENGTH_MM = 460.0
  73. CLEARANCE_VOLUME_L_BY_BORE = {
  74. 490.0: 1.51,
  75. 390.0: 0.74,
  76. 360.0: 0.62,
  77. }
  78. DEMO_POINT_NAME = "7号机组一缸压力盖侧"
  79. DEMO_DEVICE_PART = "7号机组一缸"
  80. DEMO_START = datetime(2026, 4, 12, 8, 0, 0)
  81. DEMO_POINT_COUNT = 72
  82. DEMO_SAMPLE_COUNT = 32768
  83. DEMO_REVOLUTION_SAMPLES = 800
  84. DEMO_ID_BY_TYPE = {name: 100000 + index * 1000 for index, name in enumerate(MEASUREMENT_TYPES)}
  85. def _time_string(value: Any) -> str:
  86. if isinstance(value, datetime):
  87. return value.strftime("%Y-%m-%d %H:%M:%S")
  88. return str(value)
  89. def _parse_time(value: str | None) -> datetime | None:
  90. if not value:
  91. return None
  92. return datetime.fromisoformat(value.replace("Z", "+00:00").replace("T", " "))
  93. def _safe_float(value: Any) -> float | None:
  94. if value is None:
  95. return None
  96. number = float(value)
  97. return number if np.isfinite(number) else None
  98. def _stored_first_cycle(samples: np.ndarray) -> DetectedCycle | None:
  99. """Build a cycle from wave_sample_one, which contains one cycle only."""
  100. if samples.ndim != 2 or samples.shape[1] < 3 or len(samples) < 2:
  101. return None
  102. start = 0
  103. end = len(samples)
  104. angle = np.linspace(0.0, 360.0, end - start, endpoint=False)
  105. trigger_offsets = tuple(
  106. int(offset)
  107. for offset in np.flatnonzero(
  108. np.nan_to_num(samples[:, 2], nan=-np.inf) >= 30.0,
  109. )
  110. if offset == 0 or samples[offset - 1, 2] < 30.0
  111. )
  112. return DetectedCycle(
  113. number=1,
  114. start_offset=start,
  115. end_offset=end,
  116. angle=angle,
  117. signal=samples[start:end, 1].copy(),
  118. trigger_offsets=trigger_offsets,
  119. )
  120. def _sample_first_cycle_360(samples: np.ndarray) -> list[tuple[int, float]]:
  121. """Resample the stored first cycle into 360 integer-angle points.
  122. Each integer angle 0..359 keeps the first raw sample whose crank angle is at
  123. or after that integer value (the "取第一个值" rule from showPV).
  124. """
  125. if samples.ndim != 2 or samples.shape[1] < 2 or len(samples) < 2:
  126. return []
  127. count = len(samples)
  128. angle = np.linspace(0.0, 360.0, count, endpoint=False)
  129. indices = np.searchsorted(angle, np.arange(360, dtype=float), side="left")
  130. indices = np.clip(indices, 0, count - 1)
  131. return [
  132. (int(deg), float(samples[int(idx), 1]))
  133. for deg, idx in enumerate(indices)
  134. ]
  135. def _detect_source_cycles(
  136. samples: np.ndarray,
  137. single_cycle: bool,
  138. ) -> tuple[list[DetectedCycle], dict[str, Any]]:
  139. """Detect cycles: wave_sample_one holds exactly one stored cycle already."""
  140. if single_cycle:
  141. stored = _stored_first_cycle(samples)
  142. if stored is None:
  143. return [], {"completeCycleCount": 0}
  144. return [stored], {"completeCycleCount": 1, "storedFirstCycle": 1}
  145. return detect_cycles(samples)
  146. def _series_extent(items: list[dict[str, Any]]) -> tuple[float, float]:
  147. """Whole-window min/max over a series' finite raw values."""
  148. if not items:
  149. return (0.0, 1.0)
  150. values = np.fromiter((item["rawValue"] for item in items), dtype=float, count=len(items))
  151. values = values[np.isfinite(values)]
  152. if values.size == 0:
  153. return (0.0, 1.0)
  154. return (float(values.min()), float(values.max()))
  155. def _annotation_dict(row: dict[str, Any]) -> dict[str, Any]:
  156. return {
  157. "id": int(row["id"]),
  158. "waveFileId": int(row["wave_file_id"]),
  159. "label": row["label"],
  160. "periodStart": int(row["period_start"]),
  161. "periodEnd": int(row["period_end"]),
  162. "sampleIndexStart": int(row["sample_index_start"]),
  163. "sampleIndexEnd": int(row["sample_index_end"]),
  164. }
  165. def _alarm_dict(row: dict[str, Any]) -> dict[str, Any]:
  166. return {
  167. "id": int(row["id"]),
  168. "deviceCode": row["device_code"] or "",
  169. "devicePart": row["device_part"] or "",
  170. "devicePoint": row["device_point"] or "",
  171. "alarmType": row["alarm_type"] or "",
  172. "alarmDes": row["alarm_des"] or "",
  173. "status": int(row.get("status") or 0),
  174. "alarmLevel": row.get("alarm_level"),
  175. "scanStartTime": _time_string(row["scan_start_time"]) if row.get("scan_start_time") else "",
  176. "scanEndTime": _time_string(row["scan_end_time"]) if row.get("scan_end_time") else "",
  177. "alarmInfoJson": row.get("alarm_info_json") or "",
  178. "alarmTimeStart": _time_string(row["alarm_time_start"]),
  179. "alarmTimeEnd": _time_string(row["alarm_time_end"]),
  180. }
  181. def _validate_device_points(
  182. values: list[str] | tuple[str, ...] | None,
  183. *,
  184. allow_empty: bool = False,
  185. ) -> list[str]:
  186. """校验波形点位。None 表示未传参 → 回退为全部点位;显式空列表仅在 allow_empty 时放行。"""
  187. selected = list(DEVICE_POINTS) if values is None else list(values)
  188. invalid = [value for value in selected if value not in DEVICE_POINTS]
  189. if invalid:
  190. raise ValueError(f"不支持的测试点位:{'、'.join(invalid)}")
  191. result = [value for value in DEVICE_POINTS if value in selected]
  192. if not result and not allow_empty:
  193. raise ValueError("请至少选择一个测试点位")
  194. return result
  195. def _primary_device_point(points: list[str]) -> str:
  196. """Prefer pressure cap, then pressure shaft, then any selected point."""
  197. for point in ("压力盖侧", "压力轴侧"):
  198. if point in points:
  199. return point
  200. return points[0]
  201. def _unit_number(device_part: str) -> str | None:
  202. """从 机组与部位 提取机组号(7/8/9),非 7/8/9 机组返回 None。"""
  203. match = re.match(r"^([789])号机组", device_part.strip())
  204. return match.group(1) if match else None
  205. def _site_point_column(item_name: str, unit: str | None) -> str | None:
  206. """校验全场点位名称并映射到 pks_long_sample 的列名(如 YSJ7_3 -> YSJ_3)。"""
  207. match = _SITE_POINT_PATTERN.match(item_name)
  208. if not match or match.group(1) != unit:
  209. return None
  210. return f"YSJ_{match.group(2)}"
  211. class DataService:
  212. # 数据库失败后的重试冷却时间(秒)。超过该时间后自动重连数据库,
  213. # 避免一次网络抖动就把服务永久锁死在演示数据模式。
  214. DB_RETRY_COOLDOWN = 30.0
  215. def __init__(self) -> None:
  216. self._db_failed = settings.demo_mode == "always"
  217. self._db_failed_at = monotonic() if self._db_failed else 0.0
  218. self._last_db_error = ""
  219. self._demo_annotations: dict[int, dict[str, Any]] = {}
  220. self._demo_annotation_seq = 1
  221. # 机组号 -> pks_long_sample 该批次 [min_time, max_time],进程内只查一次。
  222. self._pks_batch_bounds: dict[str, tuple[datetime, datetime]] = {}
  223. @property
  224. def source(self) -> str:
  225. return "demo" if self._db_failed else "database"
  226. @property
  227. def last_db_error(self) -> str:
  228. return self._last_db_error
  229. def _run_with_fallback(
  230. self,
  231. database_function: Callable[[], Any],
  232. demo_function: Callable[[], Any],
  233. ) -> tuple[Any, str]:
  234. if self._db_failed:
  235. if settings.demo_mode == "always":
  236. return demo_function(), "demo"
  237. if monotonic() - self._db_failed_at < self.DB_RETRY_COOLDOWN:
  238. return demo_function(), "demo"
  239. # 冷却结束,重新尝试数据库,数据库恢复后可自动切回真实数据。
  240. try:
  241. result = database_function()
  242. self._db_failed = False
  243. self._last_db_error = ""
  244. return result, "database"
  245. except Exception as error:
  246. if settings.demo_mode == "never":
  247. raise
  248. self._db_failed = True
  249. self._db_failed_at = monotonic()
  250. self._last_db_error = str(error)
  251. return demo_function(), "demo"
  252. def query_options(self) -> dict[str, Any]:
  253. def database_query():
  254. with get_connection() as connection:
  255. with connection.cursor() as cursor:
  256. cursor.execute(
  257. """
  258. SELECT device_part, device_point, measurement_type,
  259. MAX(sample_time) AS max_time,
  260. MIN(sample_time) AS min_time,
  261. COUNT(*) AS file_count
  262. FROM wave_file
  263. WHERE rpm > 0 AND device_part <> '' AND device_point <> ''
  264. GROUP BY device_part, device_point, measurement_type
  265. ORDER BY device_part, device_point
  266. """,
  267. )
  268. rows = cursor.fetchall()
  269. return [
  270. {
  271. "devicePart": row["device_part"],
  272. "devicePoint": row["device_point"],
  273. "measurementType": row["measurement_type"],
  274. "minTime": _time_string(row["min_time"]),
  275. "maxTime": _time_string(row["max_time"]),
  276. "fileCount": int(row["file_count"]),
  277. }
  278. for row in rows
  279. ]
  280. rows, source = self._run_with_fallback(database_query, self._demo_options)
  281. device_parts = list(dict.fromkeys(row["devicePart"] for row in rows))
  282. return {
  283. "source": source,
  284. "measurementTypes": list(MEASUREMENT_TYPES),
  285. "deviceParts": device_parts,
  286. "devicePoints": list(DEVICE_POINTS),
  287. "devicePointToType": dict(DEVICE_POINT_TO_TYPE),
  288. "options": rows,
  289. "notice": self._source_notice(source),
  290. }
  291. def abnormal_counts(self) -> dict[str, int]:
  292. """每个 point_name 的异常文件数量(tspluse_status > 0)。"""
  293. def database_query():
  294. with get_connection() as connection:
  295. with connection.cursor() as cursor:
  296. cursor.execute(
  297. """
  298. SELECT point_name, COUNT(*) AS cnt
  299. FROM wave_file
  300. WHERE tspluse_status > 0 AND point_name <> ''
  301. GROUP BY point_name
  302. """,
  303. )
  304. rows = cursor.fetchall()
  305. return {row["point_name"]: int(row["cnt"]) for row in rows}
  306. def demo_query():
  307. return {}
  308. result, _source = self._run_with_fallback(database_query, demo_query)
  309. return result
  310. def tspluse_ruler(self) -> dict[str, Any]:
  311. """压力部位 tspluse_status 的全局标尺,进入页面时只查询一次。"""
  312. def database_query():
  313. with get_connection() as connection:
  314. with connection.cursor() as cursor:
  315. cursor.execute(
  316. """
  317. SELECT MIN(tspluse_status) AS min_status,
  318. MAX(tspluse_status) AS max_status
  319. FROM wave_file
  320. WHERE rpm > 0 AND measurement_type = '压力'
  321. """,
  322. )
  323. row = cursor.fetchone()
  324. return {
  325. "min": int(row["min_status"]) if row and row["min_status"] is not None else 0,
  326. "max": int(row["max_status"]) if row and row["max_status"] is not None else 0,
  327. }
  328. def demo_query():
  329. return {"min": 0, "max": 0}
  330. result, source = self._run_with_fallback(database_query, demo_query)
  331. result["source"] = source
  332. return result
  333. def _pks_batch_range(self, unit: str) -> tuple[datetime, datetime] | None:
  334. """机组 pks 批次的时间覆盖范围(仅查一次并缓存)。"""
  335. cached = self._pks_batch_bounds.get(unit)
  336. if cached is not None:
  337. return cached
  338. batch = UNIT_BATCH[unit]
  339. with get_connection() as connection:
  340. with connection.cursor() as cursor:
  341. cursor.execute(
  342. "SELECT MIN(sample_time) AS lo, MAX(sample_time) AS hi "
  343. "FROM pks_long_sample WHERE import_batch_id = %s",
  344. (batch,),
  345. )
  346. row = cursor.fetchone()
  347. bounds = (row["lo"], row["hi"]) if row and row["lo"] is not None else None
  348. self._pks_batch_bounds[unit] = bounds
  349. return bounds
  350. def site_points(self, device_part: str) -> dict[str, Any]:
  351. """机组(7/8/9)的全场点位记录(site_point 表中 YSJ{机组号}_1..41)。
  352. 附带该机组 pks 批次的时间覆盖范围,供“仅 PKS 点位”组合自动赋值
  353. 开始/结束时间。
  354. """
  355. unit = _unit_number(device_part)
  356. def database_query():
  357. if unit is None:
  358. return {"items": [], "minTime": None, "maxTime": None}
  359. allowed = {
  360. f"YSJ{unit}_{n}"
  361. for n in range(1, 42)
  362. if n not in _EXCLUDED_SITE_INDEXES
  363. }
  364. with get_connection() as connection:
  365. with connection.cursor() as cursor:
  366. cursor.execute(
  367. "SELECT ItemName, ItemDescription FROM site_point WHERE ItemName LIKE %s",
  368. (f"YSJ{unit}\\_%",),
  369. )
  370. rows = cursor.fetchall()
  371. items = [
  372. {"itemName": row["ItemName"], "itemDescription": row["ItemDescription"] or ""}
  373. for row in rows
  374. if row["ItemName"] in allowed
  375. ]
  376. items.sort(key=lambda item: int(item["itemName"].rsplit("_", 1)[1]))
  377. bounds = self._pks_batch_range(unit)
  378. return {
  379. "items": items,
  380. "minTime": _time_string(bounds[0]) if bounds else None,
  381. "maxTime": _time_string(bounds[1]) if bounds else None,
  382. }
  383. def demo_query():
  384. return {"items": [], "minTime": None, "maxTime": None}
  385. result, source = self._run_with_fallback(database_query, demo_query)
  386. return {
  387. "source": source,
  388. "unit": unit,
  389. "items": result["items"],
  390. "minTime": result["minTime"],
  391. "maxTime": result["maxTime"],
  392. "notice": self._source_notice(source),
  393. }
  394. @staticmethod
  395. def _attach_site_values(points: list[dict[str, Any]], unit: str, site_points: list[str]) -> None:
  396. """把 pks_long_sample 的最近邻值挂到每个时间点的 siteValues 上。
  397. pks 数据 5 秒一条、wave_file 15 分钟一条。按用户口径做分钟/5秒级对齐:
  398. 把 wave 采样时刻四舍五入到最近的 5 秒格点,再用一次 ``sample_time IN (...)``
  399. 精确取数(结果行数 = 时间点数),避免把整段 pks 拉出来。
  400. """
  401. columns = [
  402. (item_name, column)
  403. for item_name in site_points
  404. if (column := _site_point_column(item_name, unit)) is not None
  405. ]
  406. if not columns or not points:
  407. for point in points:
  408. point.setdefault("siteValues", {})
  409. return
  410. rounded: list[datetime] = []
  411. for point in points:
  412. timestamp = datetime.strptime(point["sampleTime"], "%Y-%m-%d %H:%M:%S")
  413. rounded.append(datetime.fromtimestamp(round(timestamp.timestamp() / 5.0) * 5))
  414. chunk_size = 1000
  415. select_expr = ", ".join(f"`{column}`" for _item, column in columns)
  416. for index in range(0, len(rounded), chunk_size):
  417. chunk = rounded[index:index + chunk_size]
  418. placeholders = ", ".join(["%s"] * len(chunk))
  419. with get_connection() as connection:
  420. with connection.cursor() as cursor:
  421. cursor.execute(
  422. f"SELECT sample_time, {select_expr} FROM pks_long_sample "
  423. f"WHERE import_batch_id = %s AND sample_time IN ({placeholders})",
  424. (UNIT_BATCH[unit], *chunk),
  425. )
  426. rows = cursor.fetchall()
  427. by_time: dict[datetime, dict[str, Any]] = {
  428. row["sample_time"]: row for row in rows
  429. }
  430. for point_index in range(index, min(index + chunk_size, len(rounded))):
  431. target_time = rounded[point_index]
  432. row = by_time.get(target_time)
  433. if row is None:
  434. continue
  435. site_values = points[point_index].setdefault("siteValues", {})
  436. for item_name, column in columns:
  437. value = row[column]
  438. site_values[item_name] = float(value) if value is not None else None
  439. @staticmethod
  440. def _build_site_series(
  441. site_points: list[str],
  442. unit: str | None,
  443. points: list[dict[str, Any]],
  444. ) -> dict[str, Any]:
  445. """把每个时间点扩展为前后各 7.5 分钟的 pks 5s 曲线段。
  446. 每个选中点位(PKS)在该窗口内的每个时间点不再只画单值,而是以该
  447. 时间点采样时刻为中心,取 [t-7.5min, t+7.5min) 的原始 5s 数据映射
  448. 到该时间点在 x 轴占据的格子(索引 slot ~ slot+1)内连成一小段曲线。
  449. 相邻时间点若恰好间隔 15 分钟,则相邻窗口首尾衔接、无重叠。
  450. 数据按列做一次整段范围查询,再按时间点二分切段。
  451. """
  452. columns = [
  453. (item_name, column)
  454. for item_name in site_points
  455. if (column := _site_point_column(item_name, unit)) is not None
  456. ]
  457. centers: list[tuple[int, datetime]] = []
  458. for index, point in enumerate(points):
  459. try:
  460. timestamp = datetime.strptime(point["sampleTime"], "%Y-%m-%d %H:%M:%S")
  461. except (TypeError, ValueError):
  462. continue
  463. centers.append((index, timestamp))
  464. if not columns or not centers:
  465. return {"points": site_points, "series": []}
  466. half = timedelta(seconds=PKS_WINDOW_SECONDS // 2)
  467. query_start = min(timestamp for _index, timestamp in centers) - half
  468. query_end = max(timestamp for _index, timestamp in centers) + half
  469. select_expr = ", ".join(f"`{column}`" for _item, column in columns)
  470. with get_connection() as connection:
  471. with connection.cursor() as cursor:
  472. cursor.execute(
  473. f"SELECT sample_time, {select_expr} FROM pks_long_sample "
  474. "WHERE import_batch_id = %s AND sample_time >= %s AND sample_time < %s "
  475. "ORDER BY sample_time",
  476. (UNIT_BATCH[unit], query_start, query_end),
  477. )
  478. rows = cursor.fetchall()
  479. sample_times = [row["sample_time"] for row in rows]
  480. series = []
  481. for item_name, column in columns:
  482. values = [row[column] for row in rows]
  483. data: list[dict[str, Any]] = []
  484. for index, center in centers:
  485. low = bisect_left(sample_times, center - half)
  486. high = bisect_left(sample_times, center + half)
  487. for pos in range(low, high):
  488. raw_value = values[pos]
  489. if raw_value is None:
  490. continue
  491. value = float(raw_value)
  492. if not np.isfinite(value):
  493. continue
  494. x = index + 0.5 + (sample_times[pos] - center).total_seconds() / PKS_WINDOW_SECONDS
  495. data.append(
  496. {
  497. "value": [x, value],
  498. "x": x,
  499. "rawValue": value,
  500. "sampleTime": _time_string(sample_times[pos]),
  501. },
  502. )
  503. series.append({"itemName": item_name, "data": data})
  504. return {"points": site_points, "series": series}
  505. def time_points(
  506. self,
  507. device_part: str,
  508. device_points: list[str] | None,
  509. min_time: str | None,
  510. max_time: str | None,
  511. include_stopped: bool = False,
  512. min_status: int | None = None,
  513. status_filter: list[str] | None = None,
  514. site_points: list[str] | None = None,
  515. ) -> dict[str, Any]:
  516. if not device_part.strip():
  517. raise ValueError("机组与部位不能为空")
  518. status_filters = {value for value in (status_filter or [])}
  519. unknown = status_filters - {"abnormal", "no_cycle"}
  520. if unknown:
  521. raise ValueError(f"不支持的状态筛选:{'、'.join(sorted(unknown))}")
  522. start = _parse_time(min_time)
  523. end = _parse_time(max_time)
  524. if start and end and start > end:
  525. raise ValueError("开始时间不能晚于结束时间")
  526. unit = _unit_number(device_part)
  527. site_list = list(site_points or [])
  528. site_columns: list[tuple[str, str]] = []
  529. if unit:
  530. site_columns = [
  531. (item, column)
  532. for item in site_list
  533. if (column := _site_point_column(item, unit)) is not None
  534. ]
  535. if device_points is None or device_points == []:
  536. selected_points = (
  537. [] if (unit and site_columns) else _validate_device_points(None)
  538. )
  539. else:
  540. selected_points = _validate_device_points(device_points, allow_empty=True)
  541. pks_only = bool(unit and site_columns and not selected_points)
  542. if not selected_points and not pks_only:
  543. raise ValueError("请至少选择一个测试点位")
  544. point_names = [f"{device_part}{point}" for point in selected_points]
  545. if pks_only:
  546. points, source = self._run_with_fallback(
  547. lambda: self._pks_time_points(unit, site_columns, start, end, include_stopped),
  548. lambda: [],
  549. )
  550. return {
  551. "source": source,
  552. "devicePart": device_part,
  553. "devicePoints": [],
  554. "total": len(points),
  555. "points": points,
  556. "referencePoints": [],
  557. "notice": self._source_notice(source),
  558. }
  559. def database_query():
  560. point_placeholders = ", ".join(["%s"] * len(point_names))
  561. clauses = [
  562. f"point_name IN ({point_placeholders})",
  563. ]
  564. params: list[Any] = list(point_names)
  565. if not include_stopped:
  566. clauses.append("rpm > 0")
  567. if min_status is not None and min_status > 0 and "no_cycle" not in status_filters:
  568. clauses.append("tspluse_status >= %s")
  569. params.append(min_status)
  570. if "abnormal" in status_filters and "no_cycle" in status_filters:
  571. clauses.append("(tspluse_status > 0 OR tspluse_status = -1)")
  572. elif "abnormal" in status_filters:
  573. clauses.append("tspluse_status > 0")
  574. elif "no_cycle" in status_filters:
  575. clauses.append("tspluse_status = -1")
  576. if start:
  577. clauses.append("sample_time >= %s")
  578. params.append(start)
  579. if end:
  580. clauses.append("sample_time <= %s")
  581. params.append(end)
  582. with get_connection() as connection:
  583. with connection.cursor() as cursor:
  584. cursor.execute(
  585. f"""
  586. SELECT id, point_name, device_point, measurement_type,
  587. sample_time, sample_count, sample_frequency_hz,
  588. rpm, tspluse_status
  589. FROM wave_file
  590. WHERE {' AND '.join(clauses)}
  591. ORDER BY sample_time ASC, id ASC
  592. """,
  593. params,
  594. )
  595. rows = cursor.fetchall()
  596. reference_rows: list[dict[str, Any]] = []
  597. if status_filters and rows:
  598. primary = _primary_device_point(selected_points)
  599. reference_where = [
  600. "point_name = %s",
  601. "measurement_type = '压力'",
  602. "tspluse_status = 0",
  603. ]
  604. reference_params: list[Any] = [f"{device_part}{primary}"]
  605. if not include_stopped:
  606. reference_where.append("rpm > 0")
  607. if start:
  608. reference_where.append("sample_time >= %s")
  609. reference_params.append(start)
  610. if end:
  611. reference_where.append("sample_time <= %s")
  612. reference_params.append(end)
  613. reference_params.append(rows[0]["sample_time"])
  614. with get_connection() as connection:
  615. with connection.cursor() as cursor:
  616. cursor.execute(
  617. f"""
  618. SELECT id, point_name, device_point, measurement_type,
  619. sample_time, sample_count, sample_frequency_hz,
  620. rpm, tspluse_status
  621. FROM wave_file
  622. WHERE {' AND '.join(reference_where)}
  623. ORDER BY ABS(TIMESTAMPDIFF(SECOND, sample_time, %s)) ASC
  624. LIMIT 1
  625. """,
  626. reference_params,
  627. )
  628. reference_rows = cursor.fetchall()
  629. return self._group_time_points(rows), self._group_time_points(reference_rows)
  630. def demo_query():
  631. return self._demo_time_points(device_part, selected_points, start, end), []
  632. result, source = self._run_with_fallback(database_query, demo_query)
  633. points, reference = result
  634. # 全场点位(PKS)数据是辅助层:任意失败都静默跳过,不影响主查询。
  635. try:
  636. unit = _unit_number(device_part)
  637. if unit and site_points:
  638. self._attach_site_values(points, unit, site_points)
  639. except Exception:
  640. pass
  641. for point in points:
  642. point.setdefault("siteValues", {})
  643. return {
  644. "source": source,
  645. "devicePart": device_part,
  646. "devicePoints": selected_points,
  647. "total": len(points),
  648. "points": points,
  649. "referencePoints": reference,
  650. "notice": self._source_notice(source),
  651. }
  652. def faults(
  653. self,
  654. device_part: str,
  655. min_time: str | None,
  656. max_time: str | None,
  657. ) -> dict[str, Any]:
  658. """Return compressor_fault rows for the unit of ``device_part``.
  659. ``compressor_fault.unit_name`` ("7号机组") is matched against the unit
  660. prefix of ``wave_file.device_part`` ("7号机组一缸"). Only faults whose
  661. ``fault_date`` falls inside the requested time range are returned.
  662. """
  663. unit = _unit_number(device_part)
  664. start = _parse_time(min_time)
  665. end = _parse_time(max_time)
  666. def database_query():
  667. if unit is None:
  668. return []
  669. clauses = ["(unit_name = %s OR unit_name LIKE %s)"]
  670. params: list[Any] = [f"{unit}号机组", f"{unit}号机组%"]
  671. if start:
  672. clauses.append("fault_date >= DATE(%s)")
  673. params.append(start)
  674. if end:
  675. clauses.append("fault_date <= DATE(%s)")
  676. params.append(end)
  677. with get_connection() as connection:
  678. with connection.cursor() as cursor:
  679. cursor.execute(
  680. f"""
  681. SELECT id, unit_name, fault_date, fault_category
  682. FROM compressor_fault
  683. WHERE {' AND '.join(clauses)}
  684. ORDER BY fault_date ASC, id ASC
  685. """,
  686. params,
  687. )
  688. rows = cursor.fetchall()
  689. return [
  690. {
  691. "id": int(row["id"]),
  692. "unitName": row["unit_name"] or "",
  693. "faultDate": _time_string(row["fault_date"])[:10],
  694. "faultCategory": row["fault_category"] or "",
  695. }
  696. for row in rows
  697. ]
  698. result, source = self._run_with_fallback(database_query, lambda: [])
  699. return {
  700. "source": source,
  701. "unit": unit,
  702. "faults": result,
  703. "notice": self._source_notice(source),
  704. }
  705. def list_alarms(self, current_time: str | None = None) -> dict[str, Any]:
  706. """列出在给定时刻命中的告警。
  707. 命中条件为 ``alarm_time_start <= 当前时间 <= alarm_time_end``,
  708. 即告警时间段覆盖“告警当前时间”。
  709. """
  710. moment = _parse_time(current_time) or datetime.now()
  711. def database_query():
  712. with get_connection() as connection:
  713. with connection.cursor() as cursor:
  714. cursor.execute(
  715. """
  716. SELECT id, device_code, device_part, device_point, alarm_type,
  717. alarm_des, alarm_time_start, alarm_time_end, status,
  718. alarm_level, scan_start_time, scan_end_time, alarm_info_json
  719. FROM compressor_alarm
  720. WHERE alarm_time_start <= %s AND alarm_time_end >= %s
  721. ORDER BY alarm_time_start DESC, id DESC
  722. """,
  723. (moment, moment),
  724. )
  725. rows = cursor.fetchall()
  726. return [_alarm_dict(row) for row in rows]
  727. result, source = self._run_with_fallback(database_query, lambda: [])
  728. return {
  729. "source": source,
  730. "currentTime": _time_string(moment),
  731. "alarms": result,
  732. "notice": self._source_notice(source),
  733. }
  734. def scan_alarm(self, payload: dict[str, Any]) -> dict[str, Any]:
  735. """执行告警算法并按原有协议写入 compressor_alarm。"""
  736. device_code = str(payload.get("device_code") or "").strip()
  737. alarm_type = str(payload.get("alarm_type") or "").strip()
  738. start = _parse_time(payload.get("forecast_time"))
  739. try:
  740. hours = float(payload.get("hours") or 0)
  741. except (TypeError, ValueError):
  742. raise ValueError("告警时长必须是数字") from None
  743. if not device_code:
  744. raise ValueError("压缩机不能为空")
  745. if not alarm_type:
  746. raise ValueError("算法不能为空")
  747. if start is None:
  748. raise ValueError("预警时间不能为空")
  749. if hours <= 0:
  750. raise ValueError("告警时长必须大于 0")
  751. if alarm_type == VALVE_FAULT_ALARM_TYPE:
  752. return self._scan_valve_fault(device_code, start, hours)
  753. if alarm_type == OIL_PRESSURE_ALARM_TYPE:
  754. # `hours` controls the alarm validity window, not the analysis window.
  755. # Oil pressure analysis always uses the previous 14 days below.
  756. analysis = self._analyze_oil_pressure_alarm(device_code, start)
  757. if analysis is None:
  758. return {
  759. "source": "database",
  760. "notice": None,
  761. "action": "no_alarm",
  762. "id": 0,
  763. }
  764. device_part = analysis["device_part"]
  765. device_point = OIL_PRESSURE_DEVICE_POINT
  766. alarm_des = analysis["basis"]
  767. status = 1
  768. alarm_level = analysis["stage"]
  769. scan_start = analysis["scan_start"]
  770. scan_end = analysis["scan_end"]
  771. alarm_info_json = json.dumps(analysis["info"], ensure_ascii=False, separators=(",", ":"))
  772. elif alarm_type == CRUCIFORM_FAULT_ALARM_TYPE:
  773. analysis = self._analyze_cruciform_vibration_alarm(device_code, start)
  774. if analysis is None:
  775. return {
  776. "source": "database",
  777. "notice": None,
  778. "action": "no_alarm",
  779. "id": 0,
  780. }
  781. device_part = analysis["device_part"]
  782. device_point = VIBRATION_DEVICE_POINT
  783. alarm_des = analysis["basis"]
  784. status = 1
  785. alarm_level = analysis["stage"]
  786. scan_start = analysis["scan_start"]
  787. scan_end = analysis["scan_end"]
  788. alarm_info_json = json.dumps(analysis["info"], ensure_ascii=False, separators=(",", ":"))
  789. else:
  790. # Keep the existing endpoint contract for algorithms not yet implemented.
  791. device_part = "测试组件"
  792. device_point = "测试组件"
  793. alarm_des = "测试数据"
  794. status = 0
  795. alarm_level = None
  796. scan_start = None
  797. scan_end = None
  798. alarm_info_json = None
  799. end = start + timedelta(hours=hours)
  800. def database_query():
  801. with get_connection() as connection:
  802. with connection.cursor() as cursor:
  803. cursor.execute(
  804. """
  805. SELECT id FROM compressor_alarm
  806. WHERE device_code = %s AND device_part = %s
  807. AND device_point = %s AND alarm_type = %s
  808. AND alarm_time_start <= %s AND alarm_time_end >= %s
  809. ORDER BY id ASC
  810. LIMIT 1
  811. """,
  812. (device_code, device_part, device_point, alarm_type, start, start),
  813. )
  814. row = cursor.fetchone()
  815. if row is not None:
  816. cursor.execute(
  817. """
  818. UPDATE compressor_alarm
  819. SET alarm_time_end = %s, alarm_des = %s, status = %s,
  820. alarm_level = %s, scan_start_time = %s, scan_end_time = %s,
  821. alarm_info_json = %s
  822. WHERE id = %s
  823. """,
  824. (end, alarm_des, status, alarm_level, scan_start, scan_end, alarm_info_json, row["id"]),
  825. )
  826. return {"action": "updated", "id": int(row["id"])}
  827. cursor.execute(
  828. """
  829. INSERT INTO compressor_alarm
  830. (device_code, device_part, device_point, alarm_type,
  831. alarm_des, alarm_time_start, alarm_time_end, status,
  832. alarm_level, scan_start_time, scan_end_time, alarm_info_json)
  833. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  834. """,
  835. (device_code, device_part, device_point, alarm_type, alarm_des, start, end, status,
  836. alarm_level, scan_start, scan_end, alarm_info_json),
  837. )
  838. return {"action": "inserted", "id": int(cursor.lastrowid)}
  839. result, source = self._run_with_fallback(
  840. database_query,
  841. lambda: {"action": "demo", "id": 0},
  842. )
  843. return {"source": source, "notice": self._source_notice(source), **result}
  844. @staticmethod
  845. def _valve_segment_values(values: list[float]) -> list[float]:
  846. if not values:
  847. return []
  848. raw_average = sum(values) / len(values)
  849. return [value for value in values if value >= raw_average * 0.5]
  850. @classmethod
  851. def _valve_segment_average(cls, values: list[float]) -> float | None:
  852. filtered = cls._valve_segment_values(values)
  853. return sum(filtered) / len(filtered) if filtered else None
  854. @staticmethod
  855. def _split_valve_segments(rows: list[dict[str, Any]]) -> list[list[dict[str, Any]]]:
  856. segments: list[list[dict[str, Any]]] = []
  857. current: list[dict[str, Any]] = []
  858. for row in rows:
  859. if current and row["sample_time"] - current[-1]["sample_time"] > timedelta(days=1):
  860. segments.append(current)
  861. current = []
  862. current.append(row)
  863. if current:
  864. segments.append(current)
  865. return segments
  866. @staticmethod
  867. def _previous_valve_segment(
  868. cursor: Any, part: str, point: str, before: datetime, current_first: datetime
  869. ) -> tuple[datetime, datetime] | None:
  870. """Find the nearest preceding running segment lasting at least 24 hours."""
  871. boundary = before
  872. boundary_id = 1 << 63
  873. newer = None
  874. previous_end = None
  875. previous_start = None
  876. while True:
  877. cursor.execute(
  878. """
  879. SELECT id, sample_time FROM wave_file
  880. WHERE device_part = %s AND device_point = %s
  881. AND measurement_type = %s AND rpm > 0
  882. AND (sample_time < %s OR (sample_time = %s AND id < %s))
  883. ORDER BY sample_time DESC, id DESC LIMIT 1000
  884. """,
  885. (part, point, "压力", boundary, boundary, boundary_id),
  886. )
  887. rows = cursor.fetchall()
  888. if not rows:
  889. break
  890. for row in rows:
  891. stamp = row["sample_time"]
  892. if newer is None:
  893. previous_end = stamp
  894. previous_start = stamp
  895. elif newer - stamp > timedelta(days=1):
  896. if previous_end - previous_start >= timedelta(hours=24):
  897. return previous_start, previous_end
  898. previous_end = stamp
  899. previous_start = stamp
  900. else:
  901. previous_start = stamp
  902. newer = stamp
  903. boundary = rows[-1]["sample_time"]
  904. boundary_id = int(rows[-1]["id"])
  905. if len(rows) < 1000:
  906. break
  907. if (
  908. previous_start is not None
  909. and previous_end is not None
  910. and previous_end - previous_start >= timedelta(hours=24)
  911. ):
  912. return previous_start, previous_end
  913. return None
  914. def _scan_valve_fault(self, device_code: str, forecast_time: datetime, hours: float) -> dict[str, Any]:
  915. unit_name = device_code if device_code.endswith("号机组") else f"{device_code.rstrip('#')}号机组"
  916. scan_start = forecast_time - timedelta(days=14)
  917. scan_end = forecast_time + timedelta(days=1)
  918. def database_query():
  919. with get_connection() as connection:
  920. with connection.cursor() as cursor:
  921. cursor.execute(
  922. """
  923. SELECT f.device_part, f.device_point, f.sample_time, pa.pressure
  924. FROM wave_file f
  925. LEFT JOIN statistic_pressure_angle pa
  926. ON f.id = pa.wave_file_id AND pa.angle = 230
  927. WHERE SUBSTRING(f.device_part, 1, 4) = %s
  928. AND f.rpm > 0
  929. AND f.measurement_type = %s
  930. AND f.sample_time >= %s AND f.sample_time < %s
  931. ORDER BY f.device_part, f.device_point, f.sample_time
  932. """,
  933. (unit_name, "压力", scan_start, scan_end),
  934. )
  935. current_rows = cursor.fetchall()
  936. grouped: dict[tuple[str, str], list[dict[str, Any]]] = {}
  937. for row in current_rows:
  938. grouped.setdefault((row["device_part"], row["device_point"]), []).append(row)
  939. results = []
  940. for (part, point), running_rows in grouped.items():
  941. segments = self._split_valve_segments(running_rows)
  942. if not segments:
  943. continue
  944. current = segments[-1]
  945. current_scan_start = current[0]["sample_time"]
  946. current_values = self._valve_segment_values(
  947. [float(row["pressure"]) for row in current if row["pressure"] is not None]
  948. )
  949. current_avg = sum(current_values) / len(current_values) if current_values else None
  950. if current_avg is None or current_avg == 0:
  951. continue
  952. previous_range = None
  953. previous_values: list[float] = []
  954. eligible_previous = [
  955. segment for segment in segments[:-1]
  956. if segment[-1]["sample_time"] - segment[0]["sample_time"] >= timedelta(hours=24)
  957. ]
  958. if eligible_previous:
  959. previous = eligible_previous[-1]
  960. previous_start = previous[0]["sample_time"]
  961. previous_end = previous[-1]["sample_time"]
  962. previous_values = self._valve_segment_values(
  963. [float(row["pressure"]) for row in previous if row["pressure"] is not None]
  964. )
  965. previous_range = (previous_start, previous_end)
  966. else:
  967. previous_range = self._previous_valve_segment(
  968. cursor, part, point, current_scan_start, current_scan_start
  969. )
  970. previous_start = previous_end = None
  971. previous_avg = None
  972. if previous_range is not None:
  973. previous_start, previous_end = previous_range
  974. if not previous_values:
  975. cursor.execute(
  976. """
  977. SELECT pa.pressure FROM wave_file f
  978. INNER JOIN statistic_pressure_angle pa
  979. ON f.id = pa.wave_file_id AND pa.angle = 230
  980. WHERE f.device_part = %s AND f.device_point = %s
  981. AND f.measurement_type = %s AND f.rpm > 0
  982. AND f.sample_time >= %s AND f.sample_time <= %s
  983. """,
  984. (part, point, "压力", previous_start, previous_end),
  985. )
  986. previous_values = self._valve_segment_values(
  987. [float(row["pressure"]) for row in cursor.fetchall() if row["pressure"] is not None]
  988. )
  989. previous_avg = sum(previous_values) / len(previous_values) if previous_values else None
  990. else:
  991. previous_values = []
  992. is_fault = previous_avg is not None and previous_avg > current_avg * 1.05
  993. if not is_fault:
  994. continue
  995. info = json.dumps(
  996. {
  997. "当前段平均值": current_avg,
  998. "上一段平均值": previous_avg,
  999. "当前段最高值": max(current_values) if current_values else None,
  1000. "当前段最低值": min(current_values) if current_values else None,
  1001. "上一段最高值": max(previous_values) if previous_values else None,
  1002. "上一段最低值": min(previous_values) if previous_values else None,
  1003. "上一段开始时间": _time_string(previous_start) if previous_start else None,
  1004. "上一段结束时间": _time_string(previous_end) if previous_end else None,
  1005. },
  1006. ensure_ascii=False,
  1007. )
  1008. description = (
  1009. f"230°压力上一段均值{previous_avg:.4f}高于当前段均值"
  1010. f"{current_avg:.4f}的105%"
  1011. )
  1012. status = 1
  1013. alarm_level = "1"
  1014. cursor.execute(
  1015. """
  1016. SELECT id FROM compressor_alarm
  1017. WHERE device_code = %s AND device_part = %s AND device_point = %s
  1018. AND alarm_type = %s AND alarm_time_start = %s
  1019. LIMIT 1
  1020. """,
  1021. (device_code, part, point, VALVE_FAULT_ALARM_TYPE, forecast_time),
  1022. )
  1023. existing = cursor.fetchone()
  1024. if existing:
  1025. cursor.execute(
  1026. """
  1027. UPDATE compressor_alarm SET alarm_time_end = %s, alarm_des = %s,
  1028. alarm_level = %s, alarm_info_json = %s, status = %s,
  1029. scan_start_time = %s, scan_end_time = %s
  1030. WHERE id = %s
  1031. """,
  1032. (forecast_time + timedelta(hours=hours), description, alarm_level, info,
  1033. status, current_scan_start, forecast_time, existing["id"]),
  1034. )
  1035. results.append(("updated", int(existing["id"])))
  1036. else:
  1037. cursor.execute(
  1038. """
  1039. INSERT INTO compressor_alarm
  1040. (device_code, device_part, device_point, alarm_type, alarm_des,
  1041. alarm_level, alarm_info_json, alarm_time_start, alarm_time_end,
  1042. scan_start_time, scan_end_time, status)
  1043. VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
  1044. """,
  1045. (device_code, part, point, VALVE_FAULT_ALARM_TYPE, description, alarm_level, info,
  1046. forecast_time, forecast_time + timedelta(hours=hours), current_scan_start, forecast_time, status),
  1047. )
  1048. results.append(("inserted", int(cursor.lastrowid)))
  1049. return {
  1050. "action": results[0][0] if results else "no_alarm",
  1051. "id": results[0][1] if results else 0,
  1052. "inserted": sum(action == "inserted" for action, _ in results),
  1053. "updated": sum(action == "updated" for action, _ in results),
  1054. }
  1055. result, source = self._run_with_fallback(database_query, lambda: {"action": "demo", "id": 0})
  1056. return {"source": source, "notice": self._source_notice(source), **result}
  1057. @staticmethod
  1058. def _alarm_limit(item: dict[str, Any], alarm_type: str) -> float | None:
  1059. for index in range(1, 5):
  1060. if str(item.get(f"AlarmType{index}") or "").strip() == alarm_type:
  1061. return _safe_float(item.get(f"AlarmLimit{index}"))
  1062. return None
  1063. def _analyze_oil_pressure_alarm(self, device_code: str, end: datetime) -> dict[str, Any] | None:
  1064. """Analyze oil-pressure stages in the 14 days before ``end``."""
  1065. unit = device_code.rstrip("#").strip()
  1066. if not device_code.endswith("#") or not device_code[:-1].isdigit():
  1067. raise ValueError("压缩机编号格式必须为数字+#")
  1068. start = end - timedelta(days=OIL_PRESSURE_SCAN_DAYS)
  1069. with get_connection() as connection:
  1070. with connection.cursor() as cursor:
  1071. cursor.execute(
  1072. """
  1073. SELECT ItemName, ItemDescription, AlarmType1, AlarmType2, AlarmType3, AlarmType4,
  1074. AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4
  1075. FROM site_point
  1076. WHERE ItemName = %s
  1077. """,
  1078. (f"YSJ{unit}_5",),
  1079. )
  1080. config = cursor.fetchone()
  1081. if config is None:
  1082. raise ValueError(f"未找到 {device_code} 的润滑油压力报警配置")
  1083. low = self._alarm_limit(config, "PVLow")
  1084. low_low = self._alarm_limit(config, "PVLowLow")
  1085. if low is None or low_low is None:
  1086. raise ValueError(f"{device_code} 缺少润滑油压力低压报警阈值")
  1087. point_descriptions = {}
  1088. for point in (5, 41):
  1089. cursor.execute(
  1090. "SELECT ItemName, ItemDescription FROM site_point WHERE ItemName = %s",
  1091. (f"YSJ{unit}_{point}",),
  1092. )
  1093. item = cursor.fetchone()
  1094. if item:
  1095. point_descriptions[f"YSJ_{point}"] = item.get("ItemDescription") or f"YSJ_{point}"
  1096. cursor.execute(
  1097. """
  1098. SELECT FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start,
  1099. COUNT(*) AS samples,
  1100. SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples,
  1101. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_avg,
  1102. GROUP_CONCAT(DISTINCT import_batch_id ORDER BY import_batch_id) AS batch_ids
  1103. FROM pks_long_sample
  1104. WHERE device_code = %s
  1105. AND sample_time >= %s AND sample_time < %s
  1106. GROUP BY FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600)
  1107. ORDER BY hour_start
  1108. """,
  1109. (device_code, start, end),
  1110. )
  1111. rows = []
  1112. for row in cursor.fetchall():
  1113. running = int(row["running_samples"] or 0)
  1114. value = _safe_float(row["oil_avg"])
  1115. if running / 720 >= 0.80 and value is not None:
  1116. rows.append({
  1117. "time": row["hour_start"], "value": value, "running": running,
  1118. "batch_ids": [int(value) for value in str(row["batch_ids"] or "").split(",") if value],
  1119. })
  1120. if not rows:
  1121. return None
  1122. baseline = float(np.median([row["value"] for row in rows]))
  1123. candidates = []
  1124. history = []
  1125. for row in rows:
  1126. recent = history[-5:] + [row]
  1127. trend_values = [item["value"] for item in recent]
  1128. trend = float(np.polyfit(range(len(trend_values)), trend_values, 1)[0]) if len(trend_values) >= 3 else None
  1129. decline = 1
  1130. newer = row
  1131. for older in reversed(history):
  1132. if newer["time"] - older["time"] != timedelta(hours=1):
  1133. break
  1134. if newer["value"] > older["value"]:
  1135. break
  1136. decline += 1
  1137. newer = older
  1138. current = row["value"]
  1139. deviation = (baseline - current) / baseline if baseline else 0.0
  1140. downward = decline >= 3 or (trend is not None and trend < -0.0005)
  1141. if current <= low_low or (deviation >= 0.20 and downward):
  1142. phase = "严重异常"
  1143. elif current <= low or (deviation >= 0.10 and downward):
  1144. phase = "异常"
  1145. elif deviation >= 0.05 and downward:
  1146. phase = "轻微"
  1147. else:
  1148. phase = "正常"
  1149. candidates.append({"row": recent[-1], "phase": phase, "trend": trend, "decline": decline, "deviation": deviation})
  1150. history.append(recent[-1])
  1151. abnormal = [item for item in candidates if item["phase"] != "正常"]
  1152. if not abnormal:
  1153. return None
  1154. stage_runs = []
  1155. current_run = []
  1156. for item in candidates:
  1157. if current_run and (
  1158. item["phase"] != current_run[-1]["phase"]
  1159. or item["row"]["time"] - current_run[-1]["row"]["time"] != timedelta(hours=1)
  1160. ):
  1161. if current_run[0]["phase"] != "正常":
  1162. stage_runs.append(current_run)
  1163. current_run = []
  1164. current_run.append(item)
  1165. if current_run and current_run[0]["phase"] != "正常":
  1166. stage_runs.append(current_run)
  1167. stage_rows = max(
  1168. stage_runs,
  1169. key=lambda run: (
  1170. {"轻微": 1, "异常": 2, "严重异常": 3}[run[0]["phase"]],
  1171. len(run),
  1172. ),
  1173. )
  1174. stage = stage_rows[0]["phase"]
  1175. first = stage_rows[0]
  1176. last = stage_rows[-1]
  1177. stage_values = [item["row"]["value"] for item in stage_rows]
  1178. basis = (
  1179. f"油压{first['row']['value']:.3f}~{last['row']['value']:.3f}MPa,"
  1180. f"最低{min(stage_values):.3f}MPa,连续{len(stage_rows)}个有效运行小时;"
  1181. f"周期基线油压{baseline:.3f}MPa,相对基线变化{last['deviation']:+.1%}"
  1182. )
  1183. batch_ids = sorted({batch for item in stage_rows for batch in item["row"]["batch_ids"]})
  1184. info = {
  1185. "机组": f"{unit}号机",
  1186. "批次": batch_ids[0] if len(batch_ids) == 1 else batch_ids,
  1187. "周期编号": f"{device_code}-P14",
  1188. "开始小时": first["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
  1189. "结束小时": last["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
  1190. "持续自然小时数": int((last["row"]["time"] - first["row"]["time"]).total_seconds() / 3600) + 1,
  1191. "有效运行小时数": len(stage_rows),
  1192. "周期基线油压": baseline,
  1193. "正常波动带": "",
  1194. "阶段开始小时油压": first["row"]["value"],
  1195. "阶段结束小时油压": last["row"]["value"],
  1196. "阶段最低小时油压": min(stage_values),
  1197. "阶段开始相对周期基线变化": first["deviation"],
  1198. "阶段结束相对周期基线变化": last["deviation"],
  1199. "阶段最大连续下降有效小时数": max(item["decline"] for item in stage_rows),
  1200. "阶段低报警小时数": sum(item["row"]["value"] <= low for item in stage_rows),
  1201. "阶段低低报警小时数": sum(item["row"]["value"] <= low_low for item in stage_rows),
  1202. "是否形成确认等级": "是" if len(stage_rows) >= 3 else "否",
  1203. }
  1204. return {
  1205. "stage": stage,
  1206. "basis": basis,
  1207. "device_part": "-".join(dict.fromkeys(point_descriptions.values())),
  1208. "scan_start": start,
  1209. "scan_end": end,
  1210. "info": info,
  1211. }
  1212. def _analyze_cruciform_vibration_alarm(self, device_code: str, end: datetime) -> dict[str, Any] | None:
  1213. """Build the strongest 14-day vibration stage for the crosshead alarm."""
  1214. unit = device_code.rstrip("#").strip()
  1215. if not device_code.endswith("#") or not device_code[:-1].isdigit():
  1216. raise ValueError("压缩机编号格式必须为数字+#")
  1217. start = end - timedelta(days=OIL_PRESSURE_SCAN_DAYS)
  1218. with get_connection() as connection:
  1219. with connection.cursor() as cursor:
  1220. configs = {}
  1221. descriptions = []
  1222. for point, side in ((10, "联轴器端"), (11, "链轮端"), (41, "运行状态")):
  1223. cursor.execute(
  1224. """
  1225. SELECT ItemDescription, AlarmType1, AlarmType2, AlarmType3, AlarmType4,
  1226. AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4
  1227. FROM site_point WHERE ItemName = %s
  1228. """,
  1229. (f"YSJ{unit}_{point}",),
  1230. )
  1231. item = cursor.fetchone()
  1232. if item is None:
  1233. raise ValueError(f"未找到 {device_code} 的 YSJ_{point} 点位配置")
  1234. descriptions.append(item.get("ItemDescription") or f"YSJ_{point}")
  1235. if point != 41:
  1236. high = self._alarm_limit(item, "PVHigh")
  1237. high_high = self._alarm_limit(item, "PVHighHigh")
  1238. if high is None or high_high is None:
  1239. raise ValueError(f"{device_code} YSJ_{point} 缺少振动报警阈值")
  1240. configs[side] = {"high": high, "high_high": high_high}
  1241. cursor.execute(
  1242. """
  1243. SELECT FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start,
  1244. SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples,
  1245. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_avg,
  1246. MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_max,
  1247. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_avg,
  1248. MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_max,
  1249. GROUP_CONCAT(DISTINCT import_batch_id ORDER BY import_batch_id) AS batch_ids
  1250. FROM pks_long_sample
  1251. WHERE device_code = %s
  1252. AND sample_time >= %s AND sample_time < %s
  1253. GROUP BY FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600)
  1254. ORDER BY hour_start
  1255. """,
  1256. (device_code, start, end),
  1257. )
  1258. rows = []
  1259. for row in cursor.fetchall():
  1260. running = int(row["running_samples"] or 0)
  1261. if running / 720 < 0.80:
  1262. continue
  1263. item = {
  1264. "time": row["hour_start"], "running": running,
  1265. "batch_ids": [int(value) for value in str(row["batch_ids"] or "").split(",") if value],
  1266. }
  1267. for key in ("coupling_avg", "coupling_max", "chain_avg", "chain_max"):
  1268. item[key] = _safe_float(row[key])
  1269. if item["coupling_max"] is not None or item["chain_max"] is not None:
  1270. rows.append(item)
  1271. if len(rows) < 3:
  1272. return None
  1273. signals = []
  1274. for index, row in enumerate(rows):
  1275. prior = rows[max(0, index - 336):index]
  1276. if len(prior) < 24:
  1277. continue
  1278. side_signals = []
  1279. for side, prefix in (("联轴器端", "coupling"), ("链轮端", "chain")):
  1280. current = row[f"{prefix}_max"]
  1281. values = [item[f"{prefix}_max"] for item in prior if item[f"{prefix}_max"] is not None]
  1282. averages = [item[f"{prefix}_avg"] for item in prior if item[f"{prefix}_avg"] is not None]
  1283. if current is None or len(values) < 24:
  1284. continue
  1285. baseline = float(np.median(values))
  1286. scale = max(float(np.median(np.abs(np.asarray(values) - baseline))) * 1.4826, abs(baseline) * 0.05, 1e-6)
  1287. peak_z = (current - baseline) / scale
  1288. average = row[f"{prefix}_avg"]
  1289. average_base = float(np.median(averages)) if averages else None
  1290. average_change = average / average_base - 1 if average is not None and average_base else 0.0
  1291. limit = configs[side]
  1292. severe = current >= limit["high_high"] or peak_z >= 8
  1293. abnormal = severe or current >= limit["high"] or peak_z >= 5 or average_change >= 0.10
  1294. mild = abnormal or peak_z >= 4 or average_change >= 0.05
  1295. if mild:
  1296. side_signals.append({
  1297. "side": side, "value": current, "baseline": baseline, "scale": scale,
  1298. "peak_z": peak_z, "average_change": average_change,
  1299. "level": "严重异常" if severe else "异常" if abnormal else "轻微",
  1300. })
  1301. if side_signals:
  1302. signals.append({"row": row, "signals": side_signals})
  1303. if not signals:
  1304. return None
  1305. runs = []
  1306. current_run = []
  1307. for signal in signals:
  1308. if current_run and signal["row"]["time"] - current_run[-1]["row"]["time"] > timedelta(hours=6):
  1309. runs.append(current_run)
  1310. current_run = []
  1311. current_run.append(signal)
  1312. if current_run:
  1313. runs.append(current_run)
  1314. run = max(
  1315. runs,
  1316. key=lambda values: (
  1317. max({"轻微": 1, "异常": 2, "严重异常": 3}[item["level"]] for signal in values for item in signal["signals"]),
  1318. len(values),
  1319. ),
  1320. )
  1321. all_signals = [item for signal in run for item in signal["signals"]]
  1322. strongest = max(all_signals, key=lambda item: ({"轻微": 1, "异常": 2, "严重异常": 3}[item["level"]], item["peak_z"]))
  1323. stage = strongest["level"]
  1324. first, last = run[0], run[-1]
  1325. batch_ids = sorted({batch for signal in run for batch in signal["row"]["batch_ids"]})
  1326. max_value = max(item["value"] for item in all_signals)
  1327. basis = (
  1328. f"{strongest['side']}振动峰值或均值偏离近期基线,阶段内{len(run)}个有效运行小时触发;"
  1329. f"最大振动值{max_value:.3f}mm/s,峰值高于近期基线{strongest['peak_z']:.2f}个稳健尺度"
  1330. f"(峰值/基线{max_value / max(strongest['baseline'], 1e-6):.2f}倍)"
  1331. )
  1332. info = {
  1333. "机组": f"{unit}号机", "批次": batch_ids[0] if len(batch_ids) == 1 else batch_ids,
  1334. "周期编号": f"{device_code}-V14", "开始小时": first["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
  1335. "结束小时": last["row"]["time"].strftime("%Y-%m-%d %H:%M:%S"),
  1336. "持续自然小时数": int((last["row"]["time"] - first["row"]["time"]).total_seconds() / 3600) + 1,
  1337. "有效运行小时数": len(run), "近期基线振动": strongest["baseline"],
  1338. "正常波动带": f"{strongest['baseline'] - 2 * strongest['scale']:.3f}~{strongest['baseline'] + 2 * strongest['scale']:.3f} mm/s",
  1339. "阶段开始小时振动": max(item["value"] for item in first["signals"]),
  1340. "阶段结束小时振动": max(item["value"] for item in last["signals"]),
  1341. "阶段最大振动值": max_value,
  1342. "阶段开始相对近期基线变化": max(item["value"] / item["baseline"] - 1 for item in first["signals"]),
  1343. "阶段结束相对近期基线变化": max(item["value"] / item["baseline"] - 1 for item in last["signals"]),
  1344. "阶段最大连续异常有效小时数": len(run),
  1345. "阶段高报警小时数": sum(any(item["value"] >= configs[item["side"]]["high"] for item in signal["signals"]) for signal in run),
  1346. "阶段高高报警小时数": sum(any(item["value"] >= configs[item["side"]]["high_high"] for item in signal["signals"]) for signal in run),
  1347. "是否形成确认等级": "是" if stage != "轻微" else "否",
  1348. "最强异常侧": strongest["side"], "最强侧近期峰值基线": strongest["baseline"],
  1349. "最强侧近期稳健尺度": strongest["scale"], "阶段最大峰值/基线比例": max_value / max(strongest["baseline"], 1e-6),
  1350. }
  1351. return {
  1352. "stage": stage, "basis": basis, "device_part": "-".join(dict.fromkeys(descriptions)),
  1353. "scan_start": start, "scan_end": end, "info": info,
  1354. }
  1355. def sync_pressure_angle(self, device_part: str, device_points: list[str] | None = None) -> dict[str, Any]:
  1356. """按机组与部位批量生成 statistic_pressure_angle 采样数据。
  1357. 对每个压力 wave_file,若 statistic_pressure_angle 已存在至少一条记录则跳过,
  1358. 否则读 wave_sample_one(第一个周期)按“每个整数角度取第一个值”生成 360 个角度点。
  1359. """
  1360. points = device_points or []
  1361. def database_query():
  1362. with get_connection() as connection:
  1363. with connection.cursor() as cursor:
  1364. clauses = ["measurement_type = %s", "device_part = %s"]
  1365. params: list[Any] = ["压力", device_part]
  1366. if points:
  1367. placeholders = ", ".join(["%s"] * len(points))
  1368. clauses.append(f"device_point IN ({placeholders})")
  1369. params.extend(points)
  1370. cursor.execute(
  1371. f"SELECT id FROM wave_file WHERE {' AND '.join(clauses)} ORDER BY id",
  1372. params,
  1373. )
  1374. file_ids = [int(row["id"]) for row in cursor.fetchall()]
  1375. if not file_ids:
  1376. return {"total": 0, "created": 0, "skipped": 0}
  1377. placeholders = ", ".join(["%s"] * len(file_ids))
  1378. cursor.execute(
  1379. f"SELECT DISTINCT wave_file_id FROM statistic_pressure_angle WHERE wave_file_id IN ({placeholders})",
  1380. file_ids,
  1381. )
  1382. existing = {int(row["wave_file_id"]) for row in cursor.fetchall()}
  1383. created = 0
  1384. skipped = 0
  1385. for file_id in file_ids:
  1386. if file_id in existing:
  1387. skipped += 1
  1388. continue
  1389. cursor.execute(
  1390. "SELECT sample_index, signal_value, second_value FROM wave_sample_one "
  1391. "WHERE wave_file_id = %s ORDER BY sample_index",
  1392. (file_id,),
  1393. )
  1394. rows = cursor.fetchall()
  1395. if not rows:
  1396. skipped += 1
  1397. continue
  1398. samples = np.asarray(
  1399. [
  1400. (float(row["sample_index"]), float(row["signal_value"]),
  1401. float(row["second_value"]) if row["second_value"] is not None else 0.0)
  1402. for row in rows
  1403. ],
  1404. dtype=float,
  1405. )
  1406. sampled = _sample_first_cycle_360(samples)
  1407. if not sampled:
  1408. skipped += 1
  1409. continue
  1410. cursor.executemany(
  1411. "INSERT INTO statistic_pressure_angle (wave_file_id, period, angle, pressure) "
  1412. "VALUES (%s, %s, %s, %s)",
  1413. [(file_id, 1, angle, pressure) for angle, pressure in sampled],
  1414. )
  1415. created += 1
  1416. return {"total": len(file_ids), "created": created, "skipped": skipped}
  1417. result, source = self._run_with_fallback(
  1418. database_query,
  1419. lambda: {"total": 0, "created": 0, "skipped": 0},
  1420. )
  1421. return {"source": source, "notice": self._source_notice(source), **result}
  1422. def query_pressure_angle(
  1423. self,
  1424. device_part: str,
  1425. device_points: list[str] | None,
  1426. min_time: str | None,
  1427. max_time: str | None,
  1428. angles: list[int],
  1429. ) -> dict[str, Any]:
  1430. """按时间范围和角度列表查询 statistic_pressure_angle 曲线数据。"""
  1431. if not angles:
  1432. raise ValueError("请至少输入一个角度")
  1433. angles = [int(angle) for angle in angles if 0 <= int(angle) <= 360]
  1434. def database_query():
  1435. with get_connection() as connection:
  1436. with connection.cursor() as cursor:
  1437. clauses = ["wf.measurement_type = %s", "wf.device_part = %s"]
  1438. params: list[Any] = ["压力", device_part]
  1439. if device_points:
  1440. placeholders = ", ".join(["%s"] * len(device_points))
  1441. clauses.append(f"wf.device_point IN ({placeholders})")
  1442. params.extend(device_points)
  1443. start = _parse_time(min_time) if min_time else None
  1444. end = _parse_time(max_time) if max_time else None
  1445. if start:
  1446. clauses.append("wf.sample_time >= %s")
  1447. params.append(start)
  1448. if end:
  1449. clauses.append("wf.sample_time <= %s")
  1450. params.append(end)
  1451. angle_placeholders = ", ".join(["%s"] * len(angles))
  1452. clauses.append(f"spa.angle IN ({angle_placeholders})")
  1453. params.extend(angles)
  1454. cursor.execute(
  1455. f"""
  1456. SELECT wf.sample_time, spa.angle, spa.pressure
  1457. FROM statistic_pressure_angle spa
  1458. INNER JOIN wave_file wf ON wf.id = spa.wave_file_id
  1459. WHERE {' AND '.join(clauses)}
  1460. ORDER BY wf.sample_time, spa.angle
  1461. """,
  1462. params,
  1463. )
  1464. rows = cursor.fetchall()
  1465. stop_days = self._query_stop_days(
  1466. cursor, device_part, device_points, start, end,
  1467. )
  1468. series: dict[int, list[dict[str, Any]]] = {}
  1469. for row in rows:
  1470. angle = int(round(float(row["angle"])))
  1471. series.setdefault(angle, []).append({
  1472. "time": _time_string(row["sample_time"]),
  1473. "pressure": float(row["pressure"]),
  1474. })
  1475. for day in sorted(stop_days):
  1476. day_time = f"{day} 00:00:00"
  1477. for angle in angles:
  1478. series.setdefault(angle, []).append({
  1479. "time": day_time,
  1480. "pressure": 0.0,
  1481. })
  1482. for angle in angles:
  1483. series.setdefault(angle, []).sort(key=lambda point: point["time"])
  1484. return {
  1485. "angles": angles,
  1486. "series": [
  1487. {"angle": angle, "points": series.get(angle, [])}
  1488. for angle in angles
  1489. ],
  1490. }
  1491. result, source = self._run_with_fallback(
  1492. database_query,
  1493. lambda: {"angles": angles, "series": [{"angle": angle, "points": []} for angle in angles]},
  1494. )
  1495. return {"source": source, "notice": self._source_notice(source), **result}
  1496. @staticmethod
  1497. def _query_stop_days(
  1498. cursor: Any,
  1499. device_part: str,
  1500. device_points: list[str] | None,
  1501. start: datetime | None,
  1502. end: datetime | None,
  1503. ) -> set[str]:
  1504. """返回停机日集合:该机组在时间范围内 rpm=0 且当天无 rpm>0 记录的日期。"""
  1505. clauses = ["device_part = %s", "rpm = 0"]
  1506. params: list[Any] = [device_part]
  1507. if device_points:
  1508. placeholders = ", ".join(["%s"] * len(device_points))
  1509. clauses.append(f"device_point IN ({placeholders})")
  1510. params.extend(device_points)
  1511. if start:
  1512. clauses.append("sample_time >= %s")
  1513. params.append(start)
  1514. if end:
  1515. clauses.append("sample_time <= %s")
  1516. params.append(end)
  1517. cursor.execute(
  1518. f"""
  1519. SELECT DISTINCT DATE(sample_time) AS d
  1520. FROM wave_file
  1521. WHERE {' AND '.join(clauses)}
  1522. """,
  1523. params,
  1524. )
  1525. stop_days = {str(row["d"]) for row in cursor.fetchall()}
  1526. run_clauses = ["device_part = %s", "rpm > 0"]
  1527. run_params: list[Any] = [device_part]
  1528. if device_points:
  1529. placeholders = ", ".join(["%s"] * len(device_points))
  1530. run_clauses.append(f"device_point IN ({placeholders})")
  1531. run_params.extend(device_points)
  1532. if start:
  1533. run_clauses.append("sample_time >= %s")
  1534. run_params.append(start)
  1535. if end:
  1536. run_clauses.append("sample_time <= %s")
  1537. run_params.append(end)
  1538. cursor.execute(
  1539. f"""
  1540. SELECT DISTINCT DATE(sample_time) AS d
  1541. FROM wave_file
  1542. WHERE {' AND '.join(run_clauses)}
  1543. """,
  1544. run_params,
  1545. )
  1546. run_days = {str(row["d"]) for row in cursor.fetchall()}
  1547. return stop_days - run_days
  1548. @staticmethod
  1549. def _group_time_points(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
  1550. grouped: OrderedDict[Any, dict[str, Any]] = OrderedDict()
  1551. for row in rows:
  1552. # Channels from one acquisition batch can be stamped a few seconds
  1553. # apart (e.g. 10:30:02 vs 10:30:05), yet the acquisition cadence is
  1554. # 15 minutes. Align on the minute so one batch is one time point.
  1555. sample_time = row["sample_time"]
  1556. key = str(sample_time)[:16]
  1557. point = grouped.setdefault(
  1558. key,
  1559. {
  1560. "sampleTime": _time_string(sample_time),
  1561. "files": {},
  1562. },
  1563. )
  1564. device_point = row["device_point"]
  1565. point["files"].setdefault(device_point, {
  1566. "id": int(row["id"]),
  1567. "devicePoint": device_point,
  1568. "measurementType": row["measurement_type"],
  1569. "sampleCount": int(row["sample_count"] or 0),
  1570. "sampleFrequencyHz": int(row["sample_frequency_hz"] or 0),
  1571. "rpm": float(row["rpm"] or 0),
  1572. "status": int(row.get("tspluse_status") or 0),
  1573. })
  1574. return [
  1575. {"index": index, **point}
  1576. for index, point in enumerate(grouped.values())
  1577. ]
  1578. def _pks_time_points(
  1579. self,
  1580. unit: str,
  1581. columns: list[tuple[str, str]],
  1582. start: datetime | None,
  1583. end: datetime | None,
  1584. include_stopped: bool,
  1585. ) -> list[dict[str, Any]]:
  1586. """PKS-only 时间点:按 15 分钟槽从 pks_long_sample 生成骨架。
  1587. columns 为 [(item_name, pks 列名), ...],siteValues 直挂所选各列
  1588. 原值。停机/运转只由该机组批次内的 YSJ_41 决定(>0=运转、=0=停机),
  1589. 未勾选停机时仅保留运转锚点;勾选停机则停机锚点一并显示。
  1590. """
  1591. batch = UNIT_BATCH[unit]
  1592. clauses = ["import_batch_id = %s"]
  1593. params: list[Any] = [batch]
  1594. if start:
  1595. clauses.append("sample_time >= %s")
  1596. params.append(start)
  1597. if end:
  1598. clauses.append("sample_time <= %s")
  1599. params.append(end)
  1600. with get_connection() as connection:
  1601. with connection.cursor() as cursor:
  1602. cursor.execute(
  1603. f"""
  1604. SELECT MIN(sample_time) AS anchor
  1605. FROM pks_long_sample
  1606. WHERE {' AND '.join(clauses)}
  1607. GROUP BY FLOOR(UNIX_TIMESTAMP(sample_time) / 900)
  1608. ORDER BY anchor
  1609. """,
  1610. params,
  1611. )
  1612. anchors = [row["anchor"] for row in cursor.fetchall()]
  1613. if not anchors:
  1614. return []
  1615. if len(anchors) > PKS_STRIP_MAX_ANCHORS:
  1616. step = (len(anchors) + PKS_STRIP_MAX_ANCHORS - 1) // PKS_STRIP_MAX_ANCHORS
  1617. sampled = anchors[::step]
  1618. if sampled[-1] != anchors[-1]:
  1619. sampled.append(anchors[-1])
  1620. anchors = sampled
  1621. chunk_size = 1000
  1622. select_expr = ", ".join(f"`{column}`" for _item, column in columns)
  1623. values_by_time: dict[datetime, dict[str, float | None]] = {anchor: {} for anchor in anchors}
  1624. running_by_time: dict[datetime, bool] = {}
  1625. for index in range(0, len(anchors), chunk_size):
  1626. chunk = anchors[index:index + chunk_size]
  1627. placeholders = ", ".join(["%s"] * len(chunk))
  1628. with get_connection() as connection:
  1629. with connection.cursor() as cursor:
  1630. cursor.execute(
  1631. f"SELECT sample_time, {select_expr}, `YSJ_41` FROM pks_long_sample "
  1632. f"WHERE import_batch_id = %s AND sample_time IN ({placeholders})",
  1633. (batch, *chunk),
  1634. )
  1635. rows = cursor.fetchall()
  1636. for row in rows:
  1637. anchor_time = row["sample_time"]
  1638. values = values_by_time[anchor_time]
  1639. for item_name, column in columns:
  1640. value = row[column]
  1641. values[item_name] = float(value) if value is not None else None
  1642. running_value = row["YSJ_41"]
  1643. running_by_time[anchor_time] = bool(
  1644. running_value is not None and running_value > 0
  1645. )
  1646. points: list[dict[str, Any]] = []
  1647. for anchor in anchors:
  1648. running = running_by_time.get(anchor, False)
  1649. if not include_stopped and not running:
  1650. continue
  1651. points.append(
  1652. {
  1653. "sampleTime": _time_string(anchor),
  1654. "files": {},
  1655. "siteValues": values_by_time[anchor],
  1656. "machineRunning": running,
  1657. },
  1658. )
  1659. for index, point in enumerate(points):
  1660. point["index"] = index
  1661. return points
  1662. def _pks_only_wave_window(
  1663. self,
  1664. device_part: str,
  1665. points: list[dict[str, Any]],
  1666. site_points: list[str],
  1667. unit: str,
  1668. ) -> dict[str, Any]:
  1669. """PKS-only 波形窗口:不取 wave_file,仅构建所选 PKS 点位的微曲线。"""
  1670. def database_query():
  1671. return {
  1672. "points": points,
  1673. "siteSeries": self._build_site_series(site_points, unit, points),
  1674. }
  1675. def demo_query():
  1676. return {"points": points, "siteSeries": {"points": [], "series": []}}
  1677. result, source = self._run_with_fallback(database_query, demo_query)
  1678. return {
  1679. "devicePart": device_part,
  1680. "devicePoints": [],
  1681. "primaryPoint": None,
  1682. "points": points,
  1683. "xMin": 0,
  1684. "xMax": len(points),
  1685. "series": [],
  1686. "secondSeries": {
  1687. "name": "周期数据",
  1688. "color": "#f56c6c",
  1689. "sourceMeasurementType": None,
  1690. "sourceDevicePoint": None,
  1691. "data": [],
  1692. "finiteCount": 0,
  1693. "nonZeroCount": 0,
  1694. "min": None,
  1695. "max": None,
  1696. },
  1697. "angleSeries": {"color": "#d59b2b", "data": []},
  1698. "volumeSeries": {"color": "#4d9e6f", "data": [], "info": None},
  1699. "cycles": [],
  1700. "triggerXs": [],
  1701. "files": [],
  1702. "diagnostics": [],
  1703. "siteSeries": result["siteSeries"],
  1704. "source": source,
  1705. "notice": self._source_notice(source),
  1706. }
  1707. def alarm_compressor_options(self) -> dict[str, Any]:
  1708. """Return running compressor groups using device_part's machine prefix."""
  1709. def database_query():
  1710. with get_connection() as connection:
  1711. with connection.cursor() as cursor:
  1712. cursor.execute(
  1713. """
  1714. SELECT SUBSTRING(device_part, 1, 4) AS device_name,
  1715. COUNT(1) AS count_no
  1716. FROM wave_file
  1717. WHERE rpm > 0 AND device_part <> ''
  1718. GROUP BY SUBSTRING(device_part, 1, 4)
  1719. ORDER BY device_name
  1720. """,
  1721. )
  1722. return [
  1723. {
  1724. "deviceName": row["device_name"],
  1725. "deviceCode": f"{str(row['device_name'])[0]}#",
  1726. "count": int(row["count_no"]),
  1727. }
  1728. for row in cursor.fetchall()
  1729. if row["device_name"]
  1730. ]
  1731. result, source = self._run_with_fallback(database_query, lambda: [])
  1732. return {"source": source, "notice": self._source_notice(source), "items": result}
  1733. def wave_window(
  1734. self,
  1735. device_part: str,
  1736. device_points: list[str],
  1737. points: list[dict[str, Any]],
  1738. max_points: int,
  1739. no_sampling: bool = False,
  1740. first_cycle_only: bool = False,
  1741. site_points: list[str] | None = None,
  1742. ) -> dict[str, Any]:
  1743. selected_points = _validate_device_points(device_points, allow_empty=True)
  1744. if not points:
  1745. raise ValueError("至少选择一个时间点")
  1746. if len(points) > 200:
  1747. raise ValueError("单次最多预览 200 个时间点,请缩小时间窗口")
  1748. max_points = min(max(int(max_points), 256), 200000)
  1749. unit = _unit_number(device_part)
  1750. site_list = list(site_points or [])
  1751. site_columns: list[tuple[str, str]] = []
  1752. if unit:
  1753. site_columns = [
  1754. (item, column)
  1755. for item in site_list
  1756. if (column := _site_point_column(item, unit)) is not None
  1757. ]
  1758. pks_only = bool(unit and site_columns and not selected_points)
  1759. if not selected_points and not pks_only:
  1760. raise ValueError("请至少选择一个测试点位")
  1761. if pks_only:
  1762. return self._pks_only_wave_window(device_part, points, site_list, unit)
  1763. def database_query():
  1764. if first_cycle_only:
  1765. return self._build_first_cycle_window(
  1766. device_part,
  1767. selected_points,
  1768. points,
  1769. self._load_db_wave,
  1770. )
  1771. return self._build_wave_window(
  1772. device_part,
  1773. selected_points,
  1774. points,
  1775. max_points,
  1776. self._load_db_wave,
  1777. no_sampling,
  1778. )
  1779. def demo_query():
  1780. if first_cycle_only:
  1781. return self._build_first_cycle_window(
  1782. device_part,
  1783. selected_points,
  1784. points,
  1785. self._load_demo_wave,
  1786. )
  1787. return self._build_wave_window(
  1788. device_part,
  1789. selected_points,
  1790. points,
  1791. max_points,
  1792. self._load_demo_wave,
  1793. no_sampling,
  1794. )
  1795. result, source = self._run_with_fallback(database_query, demo_query)
  1796. # 全场点位(PKS)系列:每周期(PKS 5s)一小段、跨时间点连续的 15 分钟窗口曲线。
  1797. if first_cycle_only and site_points:
  1798. try:
  1799. unit = _unit_number(device_part)
  1800. result["siteSeries"] = (
  1801. self._build_site_series(site_points, unit, result["points"])
  1802. if unit is not None
  1803. else {"points": [], "series": []}
  1804. )
  1805. except Exception:
  1806. result["siteSeries"] = {"points": [], "series": []}
  1807. else:
  1808. result["siteSeries"] = {"points": [], "series": []}
  1809. result["source"] = source
  1810. result["notice"] = self._source_notice(source)
  1811. return result
  1812. def period_detail(self, wave_file_id: int, period_number: int) -> dict[str, Any]:
  1813. if wave_file_id <= 0 or period_number <= 0:
  1814. raise ValueError("wave_file_id 和周期编号必须为正整数")
  1815. if period_number != 1:
  1816. raise ValueError("当前数据仅保留第一个周期")
  1817. def database_query():
  1818. metadata, samples = self._load_db_wave(wave_file_id)
  1819. return self._build_period_detail(metadata, samples, period_number)
  1820. def demo_query():
  1821. metadata, samples = self._load_demo_wave(
  1822. wave_file_id,
  1823. "压力",
  1824. DEMO_POINT_NAME,
  1825. DEMO_START,
  1826. )
  1827. return self._build_period_detail(metadata, samples, period_number)
  1828. result, source = self._run_with_fallback(database_query, demo_query)
  1829. result["source"] = source
  1830. result["notice"] = self._source_notice(source)
  1831. return result
  1832. def annotation_config(self) -> dict[str, Any]:
  1833. return {
  1834. "source": self.source,
  1835. "annotationWidth": settings.annotation_width,
  1836. "notice": self._source_notice(self.source),
  1837. }
  1838. def list_annotations(self, wave_file_ids: list[int]) -> dict[str, Any]:
  1839. ids = sorted({int(value) for value in wave_file_ids if value})
  1840. if not ids:
  1841. return {
  1842. "source": self.source,
  1843. "annotations": [],
  1844. "notice": self._source_notice(self.source),
  1845. }
  1846. def database_query():
  1847. placeholders = ", ".join(["%s"] * len(ids))
  1848. with get_connection() as connection:
  1849. with connection.cursor() as cursor:
  1850. cursor.execute(
  1851. f"""
  1852. SELECT id, wave_file_id, label, period_start, period_end,
  1853. sample_index_start, sample_index_end
  1854. FROM wave_annotation
  1855. WHERE wave_file_id IN ({placeholders})
  1856. ORDER BY id ASC
  1857. """,
  1858. ids,
  1859. )
  1860. rows = cursor.fetchall()
  1861. return [_annotation_dict(row) for row in rows]
  1862. def demo_query():
  1863. return [
  1864. _annotation_dict(annotation)
  1865. for annotation in self._demo_annotations.values()
  1866. if annotation["wave_file_id"] in ids
  1867. ]
  1868. annotations, source = self._run_with_fallback(database_query, demo_query)
  1869. return {
  1870. "source": source,
  1871. "annotations": annotations,
  1872. "notice": self._source_notice(source),
  1873. }
  1874. def create_annotation(self, payload: dict[str, Any]) -> dict[str, Any]:
  1875. self._validate_annotation(payload)
  1876. wave_file_id = int(payload["wave_file_id"])
  1877. label = payload["label"]
  1878. period_start = int(payload["period_start"])
  1879. period_end = int(payload["period_end"])
  1880. sample_index_start = int(payload["sample_index_start"])
  1881. sample_index_end = int(payload["sample_index_end"])
  1882. def database_query():
  1883. with get_connection() as connection:
  1884. with connection.cursor() as cursor:
  1885. cursor.execute(
  1886. """
  1887. INSERT INTO wave_annotation
  1888. (wave_file_id, label, period_start, period_end,
  1889. sample_index_start, sample_index_end)
  1890. VALUES (%s, %s, %s, %s, %s, %s)
  1891. """,
  1892. (
  1893. wave_file_id,
  1894. label,
  1895. period_start,
  1896. period_end,
  1897. sample_index_start,
  1898. sample_index_end,
  1899. ),
  1900. )
  1901. annotation_id = cursor.lastrowid
  1902. cursor.execute(
  1903. """
  1904. SELECT id, wave_file_id, label, period_start, period_end,
  1905. sample_index_start, sample_index_end
  1906. FROM wave_annotation
  1907. WHERE id = %s
  1908. """,
  1909. (annotation_id,),
  1910. )
  1911. return _annotation_dict(cursor.fetchone())
  1912. def demo_query():
  1913. annotation_id = self._demo_annotation_seq
  1914. self._demo_annotation_seq += 1
  1915. annotation = {
  1916. "id": annotation_id,
  1917. "wave_file_id": wave_file_id,
  1918. "label": label,
  1919. "period_start": period_start,
  1920. "period_end": period_end,
  1921. "sample_index_start": sample_index_start,
  1922. "sample_index_end": sample_index_end,
  1923. }
  1924. self._demo_annotations[annotation_id] = annotation
  1925. return _annotation_dict(annotation)
  1926. result, source = self._run_with_fallback(database_query, demo_query)
  1927. result["source"] = source
  1928. result["notice"] = self._source_notice(source)
  1929. return result
  1930. def delete_annotation(self, annotation_id: int) -> dict[str, Any]:
  1931. if annotation_id <= 0:
  1932. raise ValueError("标注 id 必须为正整数")
  1933. def database_query():
  1934. with get_connection() as connection:
  1935. with connection.cursor() as cursor:
  1936. cursor.execute(
  1937. "DELETE FROM wave_annotation WHERE id = %s",
  1938. (annotation_id,),
  1939. )
  1940. return int(cursor.rowcount)
  1941. def demo_query():
  1942. if annotation_id not in self._demo_annotations:
  1943. return 0
  1944. del self._demo_annotations[annotation_id]
  1945. return 1
  1946. deleted, source = self._run_with_fallback(database_query, demo_query)
  1947. if not deleted:
  1948. raise ValueError(f"标注 id={annotation_id} 不存在")
  1949. return {
  1950. "deleted": annotation_id,
  1951. "source": source,
  1952. "notice": self._source_notice(source),
  1953. }
  1954. @staticmethod
  1955. def _validate_annotation(payload: dict[str, Any]) -> None:
  1956. label = payload.get("label")
  1957. if label not in ANNOTATION_LABELS:
  1958. raise ValueError("样本类型只能是 正常 或 异常")
  1959. for field in ("wave_file_id", "period_start", "period_end", "sample_index_start", "sample_index_end"):
  1960. if payload.get(field) is None:
  1961. raise ValueError(f"{field} 不能为空")
  1962. if int(payload["wave_file_id"]) <= 0:
  1963. raise ValueError("wave_file_id 必须为正整数")
  1964. if int(payload["period_start"]) <= 0 or int(payload["period_end"]) <= 0:
  1965. raise ValueError("周期编号必须为正整数")
  1966. if int(payload["period_start"]) > int(payload["period_end"]):
  1967. raise ValueError("起始周期不能大于结束周期")
  1968. if int(payload["sample_index_start"]) < 0 or int(payload["sample_index_end"]) < 0:
  1969. raise ValueError("采样点索引不能为负")
  1970. if int(payload["sample_index_start"]) > int(payload["sample_index_end"]):
  1971. raise ValueError("起始采样点不能大于结束采样点")
  1972. def health(self) -> dict[str, Any]:
  1973. return {
  1974. "status": "ok",
  1975. "source": self.source,
  1976. "databaseError": self._last_db_error or None,
  1977. }
  1978. def _source_notice(self, source: str) -> str | None:
  1979. if source == "demo":
  1980. if self._last_db_error:
  1981. return f"当前为演示数据:数据库暂不可用({self._last_db_error})"
  1982. return "当前为演示数据:可设置 DEMO_MODE=never 强制使用数据库"
  1983. return "已连接 MySQL 数据库"
  1984. def _build_wave_window(
  1985. self,
  1986. device_part: str,
  1987. device_points: list[str],
  1988. points: list[dict[str, Any]],
  1989. max_points: int,
  1990. loader: Callable[..., tuple[dict[str, Any], np.ndarray]],
  1991. no_sampling: bool = False,
  1992. ) -> dict[str, Any]:
  1993. series_data: dict[str, list[dict[str, Any]]] = {point: [] for point in device_points}
  1994. angle_data: list[dict[str, Any]] = []
  1995. volume_data: list[dict[str, Any]] = []
  1996. volume_info: dict[str, Any] | None = None
  1997. cycles: list[dict[str, Any]] = []
  1998. triggers: list[float] = []
  1999. files: list[dict[str, Any]] = []
  2000. diagnostics: list[dict[str, Any]] = []
  2001. second_series_data: list[dict[str, Any]] = []
  2002. second_finite_count = 0
  2003. second_non_zero_count = 0
  2004. second_min: float | None = None
  2005. second_max: float | None = None
  2006. primary = _primary_device_point(device_points)
  2007. primary_type = DEVICE_POINT_TO_TYPE[primary]
  2008. pressure_points = [point for point in device_points if DEVICE_POINT_TO_TYPE[point] == "压力"]
  2009. load_cache: dict[int, tuple[dict[str, Any], np.ndarray]] = {}
  2010. single_cycle = loader.__name__ != "_load_demo_wave"
  2011. for slot, point in enumerate(points):
  2012. target_per_file = max(256, int(np.ceil(max_points / max(len(points), 1))))
  2013. slot_files: dict[str, dict[str, Any]] = {}
  2014. for device_point in device_points:
  2015. file_info = point.get("files", {}).get(device_point)
  2016. if file_info:
  2017. slot_files[device_point] = file_info
  2018. # The per-slot "period source" supplies the 周期数据/体积/角度 series
  2019. # and the background cycle bands. Primary point is preferred; when a
  2020. # selected point's timestamp differs (e.g. 10:30:02 vs 10:30:05), a
  2021. # slot may only contain another device point, which is used instead
  2022. # so the period data and cycle bands do not disappear there.
  2023. source_point = primary if primary in slot_files else (next(iter(slot_files)) if slot_files else None)
  2024. source_samples: np.ndarray | None = None
  2025. source_detected: list[Any] = []
  2026. source_angle_full: np.ndarray | None = None
  2027. source_volume: np.ndarray | None = None
  2028. source_required: set[int] | None = None
  2029. source_id: int | None = None
  2030. source_type = ""
  2031. if source_point is not None:
  2032. source_info = slot_files[source_point]
  2033. source_id = int(source_info["id"] if isinstance(source_info, dict) else source_info)
  2034. source_type = DEVICE_POINT_TO_TYPE[source_point]
  2035. try:
  2036. source_meta, source_samples = self._load_for_window(
  2037. loader,
  2038. source_id,
  2039. source_type,
  2040. device_part + source_point,
  2041. point["sampleTime"],
  2042. load_cache,
  2043. )
  2044. except ValueError:
  2045. source_samples = None
  2046. if source_samples is not None and len(source_samples):
  2047. source_detected, source_diagnostic = _detect_source_cycles(source_samples, single_cycle)
  2048. source_angle_vector = build_angle_vector(len(source_samples), source_detected)
  2049. source_angle_full = np.full(len(source_samples), np.nan, dtype=float)
  2050. for detected_cycle in source_detected:
  2051. source_angle_full[detected_cycle.start_offset:detected_cycle.end_offset] = detected_cycle.angle
  2052. source_volume, current_volume_info = self._build_volume_vector(
  2053. len(source_samples),
  2054. source_detected,
  2055. device_part + primary,
  2056. )
  2057. if current_volume_info is not None:
  2058. volume_info = current_volume_info
  2059. source_indices = source_samples[:, 0].astype(np.int64)
  2060. source_span = max(len(source_samples), 1)
  2061. finite_second = source_samples[:, 2][np.isfinite(source_samples[:, 2])]
  2062. if len(finite_second):
  2063. second_finite_count += int(len(finite_second))
  2064. second_non_zero_count += int(np.count_nonzero(finite_second != 0))
  2065. current_min = float(np.min(finite_second))
  2066. current_max = float(np.max(finite_second))
  2067. second_min = current_min if second_min is None else min(second_min, current_min)
  2068. second_max = current_max if second_max is None else max(second_max, current_max)
  2069. source_required = {0, len(source_samples) - 1}
  2070. for cycle in source_detected:
  2071. source_required.update(
  2072. {
  2073. cycle.start_offset,
  2074. max(cycle.end_offset - 1, cycle.start_offset),
  2075. *cycle.trigger_offsets,
  2076. },
  2077. )
  2078. start_x = slot + cycle.start_offset / source_span
  2079. end_x = slot + cycle.end_offset / source_span
  2080. cycles.append(
  2081. {
  2082. "id": f"{source_id}:{cycle.number}",
  2083. "waveFileId": source_id,
  2084. "periodNo": cycle.number,
  2085. "pointIndex": slot,
  2086. "sampleTime": point["sampleTime"],
  2087. "startX": start_x,
  2088. "endX": end_x,
  2089. "startSampleIndex": int(source_indices[cycle.start_offset]),
  2090. "endSampleIndex": int(
  2091. source_indices[max(cycle.end_offset - 1, cycle.start_offset)],
  2092. ),
  2093. "sourceType": source_type,
  2094. "devicePoint": source_point,
  2095. "background": True,
  2096. },
  2097. )
  2098. for trigger_offset in sorted(source_required):
  2099. if any(
  2100. trigger_offset == run_offset
  2101. for cycle in source_detected
  2102. for run_offset in cycle.trigger_offsets
  2103. ):
  2104. triggers.append(slot + trigger_offset / source_span)
  2105. diagnostics.append(
  2106. {
  2107. "waveFileId": source_id,
  2108. "measurementType": source_type,
  2109. "devicePoint": source_point,
  2110. "sampleTime": point["sampleTime"],
  2111. **source_diagnostic,
  2112. },
  2113. )
  2114. source_base_index = int(source_indices[0])
  2115. second_chosen = downsample_indices(source_samples[:, 2], target_per_file, source_required)
  2116. for offset in second_chosen:
  2117. second_value = _safe_float(source_samples[offset, 2])
  2118. if second_value is None:
  2119. continue
  2120. x = slot + (int(source_indices[offset]) - source_base_index) / source_span
  2121. second_series_data.append(
  2122. {
  2123. "value": [x, second_value],
  2124. "x": x,
  2125. "rawValue": second_value,
  2126. "sampleIndex": int(source_indices[offset]),
  2127. "waveFileId": source_id,
  2128. "sampleTime": point["sampleTime"],
  2129. },
  2130. )
  2131. angle_chosen = downsample_indices(source_angle_vector, target_per_file, source_required)
  2132. for offset in angle_chosen:
  2133. value = _safe_float(source_angle_vector[offset])
  2134. if value is None:
  2135. continue
  2136. x = slot + (int(source_indices[offset]) - source_base_index) / source_span
  2137. angle_data.append(
  2138. {
  2139. "value": [x, value],
  2140. "x": x,
  2141. "angle": value,
  2142. "sampleIndex": int(source_indices[offset]),
  2143. "waveFileId": source_id,
  2144. "sampleTime": point["sampleTime"],
  2145. },
  2146. )
  2147. volume_chosen = downsample_indices(source_volume, target_per_file, source_required)
  2148. for offset in volume_chosen:
  2149. value = _safe_float(source_volume[offset])
  2150. if value is None:
  2151. continue
  2152. x = slot + (int(source_indices[offset]) - source_base_index) / source_span
  2153. volume_data.append(
  2154. {
  2155. "value": [x, value],
  2156. "x": x,
  2157. "volume": value,
  2158. "sampleIndex": int(source_indices[offset]),
  2159. "waveFileId": source_id,
  2160. "sampleTime": point["sampleTime"],
  2161. },
  2162. )
  2163. for device_point in device_points:
  2164. file_info = slot_files.get(device_point)
  2165. if not file_info:
  2166. continue
  2167. file_id = int(file_info["id"] if isinstance(file_info, dict) else file_info)
  2168. measurement_type = DEVICE_POINT_TO_TYPE[device_point]
  2169. try:
  2170. metadata, samples = self._load_for_window(
  2171. loader,
  2172. file_id,
  2173. measurement_type,
  2174. device_part + device_point,
  2175. point["sampleTime"],
  2176. load_cache,
  2177. )
  2178. except ValueError:
  2179. continue
  2180. sample_count = len(samples)
  2181. if not sample_count:
  2182. continue
  2183. sample_indices_for_file = samples[:, 0].astype(np.int64)
  2184. file_base_index = int(sample_indices_for_file[0])
  2185. file_span = max(sample_count, 1)
  2186. file_required = {0, sample_count - 1}
  2187. own_volume: np.ndarray | None = None
  2188. own_angle: np.ndarray | None = None
  2189. if device_point in pressure_points:
  2190. if source_point == device_point and source_samples is not None:
  2191. own_volume = source_volume
  2192. own_angle = source_angle_full
  2193. if source_required is not None:
  2194. file_required.update(source_required)
  2195. else:
  2196. detected_own, _ = _detect_source_cycles(samples, single_cycle)
  2197. for cycle in detected_own:
  2198. file_required.update(
  2199. {
  2200. cycle.start_offset,
  2201. max(cycle.end_offset - 1, cycle.start_offset),
  2202. *cycle.trigger_offsets,
  2203. },
  2204. )
  2205. cycles.append(
  2206. {
  2207. "id": f"{file_id}:{cycle.number}",
  2208. "waveFileId": file_id,
  2209. "periodNo": cycle.number,
  2210. "pointIndex": slot,
  2211. "sampleTime": point["sampleTime"],
  2212. "startX": slot + cycle.start_offset / file_span,
  2213. "endX": slot + cycle.end_offset / file_span,
  2214. "startSampleIndex": int(sample_indices_for_file[cycle.start_offset]),
  2215. "endSampleIndex": int(
  2216. sample_indices_for_file[max(cycle.end_offset - 1, cycle.start_offset)],
  2217. ),
  2218. "sourceType": measurement_type,
  2219. "devicePoint": device_point,
  2220. "background": False,
  2221. },
  2222. )
  2223. if detected_own:
  2224. angle_own = build_angle_vector(len(samples), detected_own)
  2225. full_own = np.full(len(samples), np.nan, dtype=float)
  2226. for detected_cycle in detected_own:
  2227. full_own[detected_cycle.start_offset:detected_cycle.end_offset] = detected_cycle.angle
  2228. own_volume, _ = self._build_volume_vector(
  2229. len(samples),
  2230. detected_own,
  2231. device_part + device_point,
  2232. )
  2233. own_angle = full_own
  2234. if no_sampling:
  2235. target_per_file_own = sample_count
  2236. else:
  2237. target_per_file_own = target_per_file
  2238. chosen = downsample_indices(samples[:, 1], target_per_file_own, file_required)
  2239. for offset in chosen:
  2240. x = slot + (int(sample_indices_for_file[offset]) - file_base_index) / file_span
  2241. raw_value = float(samples[offset, 1])
  2242. series_data[device_point].append(
  2243. {
  2244. "value": [x, raw_value],
  2245. "x": x,
  2246. "rawValue": raw_value,
  2247. "sampleIndex": int(sample_indices_for_file[offset]),
  2248. "waveFileId": file_id,
  2249. "sampleTime": point["sampleTime"],
  2250. "secondValue": _safe_float(samples[offset, 2]),
  2251. "volume": _safe_float(own_volume[offset]) if own_volume is not None else None,
  2252. "angle360": _safe_float(own_angle[offset]) if own_angle is not None else None,
  2253. },
  2254. )
  2255. files.append(
  2256. {
  2257. "id": file_id,
  2258. "pointIndex": slot,
  2259. "sampleTime": point["sampleTime"],
  2260. "devicePoint": device_point,
  2261. "measurementType": measurement_type,
  2262. "sampleCount": int(metadata.get("sample_count") or sample_count),
  2263. "sampleFrequencyHz": int(metadata.get("sample_frequency_hz") or 0),
  2264. "pointName": str(metadata.get("point_name") or device_part + device_point),
  2265. "rpm": float(metadata.get("rpm") or 0),
  2266. "status": int(metadata.get("tspluse_status") or 0),
  2267. "fileName": str(metadata.get("file_name") or ""),
  2268. },
  2269. )
  2270. load_cache.clear()
  2271. for device_point in series_data:
  2272. series_data[device_point].sort(key=lambda item: item["x"])
  2273. second_series_data.sort(key=lambda item: item["x"])
  2274. angle_data.sort(key=lambda item: item["x"])
  2275. volume_data.sort(key=lambda item: item["x"])
  2276. extents = {device_point: _series_extent(series_data[device_point]) for device_point in device_points}
  2277. return {
  2278. "devicePart": device_part,
  2279. "devicePoints": device_points,
  2280. "primaryPoint": primary,
  2281. "points": points,
  2282. "xMin": 0,
  2283. "xMax": len(points),
  2284. "series": [
  2285. {
  2286. "devicePoint": device_point,
  2287. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  2288. "color": MEASUREMENT_COLORS[DEVICE_POINT_TO_TYPE[device_point]],
  2289. "data": series_data[device_point],
  2290. "min": extents[device_point][0],
  2291. "max": extents[device_point][1],
  2292. }
  2293. for device_point in device_points
  2294. ],
  2295. "secondSeries": {
  2296. "name": "周期数据",
  2297. "color": "#f56c6c",
  2298. "sourceMeasurementType": primary_type,
  2299. "sourceDevicePoint": primary,
  2300. "data": second_series_data,
  2301. "finiteCount": second_finite_count,
  2302. "nonZeroCount": second_non_zero_count,
  2303. "min": second_min,
  2304. "max": second_max,
  2305. },
  2306. "angleSeries": {
  2307. "color": "#d59b2b",
  2308. "data": angle_data,
  2309. },
  2310. "volumeSeries": {
  2311. "color": "#4d9e6f",
  2312. "data": volume_data,
  2313. "info": volume_info,
  2314. },
  2315. "cycles": cycles,
  2316. "triggerXs": sorted(set(triggers)),
  2317. "files": files,
  2318. "diagnostics": diagnostics,
  2319. }
  2320. def _build_first_cycle_window(
  2321. self,
  2322. device_part: str,
  2323. device_points: list[str],
  2324. points: list[dict[str, Any]],
  2325. loader: Callable[..., tuple[dict[str, Any], np.ndarray]],
  2326. ) -> dict[str, Any]:
  2327. """Build a slot-layout window where each file shows only its first cycle.
  2328. The primary point is selected in pressure-cap, pressure-shaft, then
  2329. selected-point order. It is the point used for the first-cycle slice;
  2330. files without recorded bounds still appear in the file list.
  2331. The slice for every primary file is fetched in a single JOIN query (range
  2332. scan on the ``(wave_file_id, sample_index)`` primary key). Other selected
  2333. device points contribute no curve but still appear in the file list. The
  2334. curve is continuous: all cycle samples are kept and each file spans its
  2335. own x slot.
  2336. """
  2337. primary = _primary_device_point(device_points)
  2338. primary_type = DEVICE_POINT_TO_TYPE[primary]
  2339. series_data: dict[str, list[dict[str, Any]]] = {point: [] for point in device_points}
  2340. second_series_data: list[dict[str, Any]] = []
  2341. angle_data: list[dict[str, Any]] = []
  2342. volume_data: list[dict[str, Any]] = []
  2343. volume_info: dict[str, Any] | None = None
  2344. cycles: list[dict[str, Any]] = []
  2345. triggers: list[float] = []
  2346. files: list[dict[str, Any]] = []
  2347. diagnostics: list[dict[str, Any]] = []
  2348. second_finite_count = 0
  2349. second_non_zero_count = 0
  2350. second_min: float | None = None
  2351. second_max: float | None = None
  2352. slot_sources: list[tuple[int, dict[str, Any], str, int]] = []
  2353. all_file_ids: set[int] = set()
  2354. for slot, point in enumerate(points):
  2355. for device_point in device_points:
  2356. file_info = point.get("files", {}).get(device_point)
  2357. if not file_info:
  2358. continue
  2359. file_id = int(file_info["id"] if isinstance(file_info, dict) else file_info)
  2360. slot_sources.append((slot, point, device_point, file_id))
  2361. all_file_ids.add(file_id)
  2362. if not slot_sources:
  2363. return self._assemble_first_cycle_window(
  2364. device_part, device_points, primary, points, series_data, second_series_data,
  2365. angle_data, volume_data, volume_info,
  2366. second_finite_count, second_non_zero_count, second_min, second_max,
  2367. cycles, triggers, files, diagnostics, 0,
  2368. )
  2369. is_demo = loader.__name__ == "_load_demo_wave"
  2370. metas: dict[int, dict[str, Any]] = {}
  2371. # file_id -> (cycle_start, cycle_end, padded_slice[offset, signal, second])
  2372. slices: dict[int, tuple[int, int, np.ndarray]] = {}
  2373. if is_demo:
  2374. for _slot, point, device_point, source_id in slot_sources:
  2375. metadata, samples = self._load_for_window(
  2376. loader, source_id, DEVICE_POINT_TO_TYPE[device_point],
  2377. device_part + device_point, point["sampleTime"], {},
  2378. )
  2379. metas[source_id] = metadata
  2380. detected, _ = detect_cycles(samples)
  2381. if detected:
  2382. cycle = detected[0]
  2383. start_si = int(samples[cycle.start_offset, 0])
  2384. end_si = int(samples[cycle.end_offset, 0])
  2385. pad_end = min(cycle.end_offset + FIRST_CYCLE_PAD, len(samples))
  2386. slices[source_id] = (start_si, end_si, samples[cycle.start_offset:pad_end])
  2387. else:
  2388. source_ids = [source_id for _, _, _, source_id in slot_sources]
  2389. with get_connection() as connection:
  2390. with connection.cursor() as cursor:
  2391. placeholders = ", ".join(["%s"] * len(all_file_ids))
  2392. cursor.execute(
  2393. f"""
  2394. SELECT id, point_name, measurement_type, sample_time,
  2395. sample_count, sample_frequency_hz, rpm, tspluse_status,
  2396. file_name, cycle_start, cycle_end
  2397. FROM wave_file
  2398. WHERE id IN ({placeholders})
  2399. """,
  2400. tuple(all_file_ids),
  2401. )
  2402. for row in cursor.fetchall():
  2403. metas[int(row["id"])] = row
  2404. source_placeholders = ", ".join(["%s"] * len(source_ids))
  2405. cursor.execute(
  2406. f"""
  2407. SELECT ws.wave_file_id,
  2408. ws.sample_index,
  2409. CAST(ws.signal_value AS FLOAT) AS sig,
  2410. CAST(ws.second_value AS FLOAT) AS sec
  2411. FROM wave_sample_one ws
  2412. WHERE ws.wave_file_id IN ({source_placeholders})
  2413. ORDER BY ws.wave_file_id, ws.sample_index
  2414. """,
  2415. tuple(source_ids),
  2416. )
  2417. raw: dict[int, list[tuple[float, float, float]]] = {}
  2418. for row in cursor.fetchall():
  2419. raw.setdefault(int(row["wave_file_id"]), []).append(
  2420. (
  2421. float(row["sample_index"]),
  2422. float(row["sig"]),
  2423. float(row["sec"]) if row["sec"] is not None else float("nan"),
  2424. ),
  2425. )
  2426. for file_id, rows in raw.items():
  2427. meta = metas.get(file_id)
  2428. if meta is None or not rows:
  2429. continue
  2430. # wave_sample_one 只保留第一个周期,返回的数组本身就是该周期。
  2431. slices[file_id] = (
  2432. int(rows[0][0]),
  2433. int(rows[-1][0]) + 1,
  2434. np.asarray(rows, dtype=float),
  2435. )
  2436. # Per-slot period source: primary preferred, else the first selected
  2437. # point with a file at that slot (timestamps may differ by seconds).
  2438. source_by_slot: dict[int, str] = {}
  2439. for _slot, _point, device_point, _file_id in slot_sources:
  2440. if _slot not in source_by_slot:
  2441. source_by_slot[_slot] = device_point
  2442. for _slot, _point, device_point, _file_id in slot_sources:
  2443. if device_point == primary:
  2444. source_by_slot[_slot] = primary
  2445. for slot, point, device_point, source_id in slot_sources:
  2446. sample_time = point["sampleTime"]
  2447. is_slot_source = device_point == source_by_slot.get(slot, device_point)
  2448. cycle_bounds = slices.get(source_id)
  2449. cycle_start = 0
  2450. cycle_end = 0
  2451. slice_arr: np.ndarray | None = None
  2452. if cycle_bounds is not None:
  2453. cycle_start, cycle_end, slice_arr = cycle_bounds
  2454. cycle_len = len(slice_arr) if slice_arr is not None else 0
  2455. has_cycle = slice_arr is not None and 1 < cycle_len
  2456. if has_cycle:
  2457. angle: np.ndarray | None = None
  2458. volume: np.ndarray | None = None
  2459. point_volume_info: dict[str, Any] | None = None
  2460. stored_cycle = _stored_first_cycle(slice_arr)
  2461. detected = [stored_cycle] if stored_cycle is not None else []
  2462. if detected:
  2463. cycle = detected[0]
  2464. angle = cycle.angle
  2465. volume, point_volume_info = self._build_volume_vector(
  2466. len(slice_arr),
  2467. [cycle],
  2468. device_part + primary,
  2469. )
  2470. if is_slot_source and point_volume_info is not None:
  2471. volume_info = point_volume_info
  2472. for offset in range(cycle_len):
  2473. sample_index = int(slice_arr[offset, 0])
  2474. raw_value = float(slice_arr[offset, 1])
  2475. second = _safe_float(slice_arr[offset, 2])
  2476. angle360 = _safe_float(angle[offset]) if angle is not None and offset < len(angle) else None
  2477. volume_value = _safe_float(volume[offset]) if volume is not None and offset < len(volume) else None
  2478. x = slot + (sample_index - cycle_start) / cycle_len
  2479. series_data[device_point].append(
  2480. {
  2481. "value": [x, raw_value],
  2482. "x": x,
  2483. "rawValue": raw_value,
  2484. "sampleIndex": sample_index,
  2485. "waveFileId": source_id,
  2486. "sampleTime": sample_time,
  2487. "secondValue": second,
  2488. "volume": volume_value,
  2489. "angle360": angle360,
  2490. },
  2491. )
  2492. if is_slot_source:
  2493. if volume_value is not None:
  2494. volume_data.append(
  2495. {
  2496. "value": [x, volume_value],
  2497. "x": x,
  2498. "volume": volume_value,
  2499. "sampleIndex": sample_index,
  2500. "waveFileId": source_id,
  2501. "sampleTime": sample_time,
  2502. },
  2503. )
  2504. if angle360 is not None:
  2505. display_angle = angle360 if angle360 <= 180.0 else 360.0 - angle360
  2506. angle_data.append(
  2507. {
  2508. "value": [x, display_angle],
  2509. "x": x,
  2510. "angle": display_angle,
  2511. "sampleIndex": sample_index,
  2512. "waveFileId": source_id,
  2513. "sampleTime": sample_time,
  2514. },
  2515. )
  2516. if second is not None:
  2517. second_finite_count += 1
  2518. if second != 0:
  2519. second_non_zero_count += 1
  2520. second_min = second if second_min is None else min(second_min, second)
  2521. second_max = second if second_max is None else max(second_max, second)
  2522. second_series_data.append(
  2523. {
  2524. "value": [x, second],
  2525. "x": x,
  2526. "rawValue": second,
  2527. "sampleIndex": sample_index,
  2528. "waveFileId": source_id,
  2529. "sampleTime": sample_time,
  2530. },
  2531. )
  2532. if offset > 0:
  2533. prev = _safe_float(slice_arr[offset - 1, 2])
  2534. if second is not None and (prev is None or prev < 30) and second >= 30:
  2535. triggers.append(x)
  2536. cycles.append(
  2537. {
  2538. "id": f"{source_id}:1",
  2539. "waveFileId": source_id,
  2540. "periodNo": 1,
  2541. "pointIndex": slot,
  2542. "sampleTime": sample_time,
  2543. "startX": float(slot),
  2544. "endX": float(slot + 1),
  2545. "startSampleIndex": cycle_start,
  2546. "endSampleIndex": cycle_end - 1,
  2547. "sourceType": DEVICE_POINT_TO_TYPE[device_point],
  2548. "devicePoint": device_point,
  2549. "background": is_slot_source,
  2550. },
  2551. )
  2552. diagnostics.append(
  2553. {
  2554. "waveFileId": source_id,
  2555. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  2556. "devicePoint": device_point,
  2557. "sampleTime": sample_time,
  2558. "cycleStart": cycle_start,
  2559. "cycleEnd": cycle_end,
  2560. "cycleSampleCount": cycle_len,
  2561. },
  2562. )
  2563. for slot, point in enumerate(points):
  2564. for device_point in device_points:
  2565. file_info = point.get("files", {}).get(device_point)
  2566. if not file_info:
  2567. continue
  2568. file_id = int(file_info["id"] if isinstance(file_info, dict) else file_info)
  2569. metadata = metas.get(file_id)
  2570. cycle_start = metadata.get("cycle_start") if metadata else None
  2571. cycle_end = metadata.get("cycle_end") if metadata else None
  2572. files.append(
  2573. {
  2574. "id": file_id,
  2575. "pointIndex": slot,
  2576. "sampleTime": point["sampleTime"],
  2577. "devicePoint": device_point,
  2578. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  2579. "sampleCount": int(
  2580. (metadata.get("sample_count") if metadata else file_info.get("sampleCount") or 0) or 0,
  2581. ),
  2582. "sampleFrequencyHz": int(
  2583. (metadata.get("sample_frequency_hz") if metadata else file_info.get("sampleFrequencyHz") or 0) or 0,
  2584. ),
  2585. "pointName": str(
  2586. (metadata.get("point_name") if metadata else device_part + device_point)
  2587. or device_part + device_point,
  2588. ),
  2589. "rpm": float((metadata.get("rpm") if metadata else file_info.get("rpm") or 0) or 0),
  2590. "status": int(
  2591. (metadata.get("tspluse_status") if metadata else file_info.get("status") or 0) or 0,
  2592. ),
  2593. "fileName": str((metadata.get("file_name") if metadata else "") or ""),
  2594. "cycleStart": int(cycle_start) if cycle_start is not None else None,
  2595. "cycleEnd": int(cycle_end) if cycle_end is not None else None,
  2596. },
  2597. )
  2598. return self._assemble_first_cycle_window(
  2599. device_part, device_points, primary, points, series_data, second_series_data,
  2600. angle_data, volume_data, volume_info,
  2601. second_finite_count, second_non_zero_count, second_min, second_max,
  2602. cycles, triggers, files, diagnostics,
  2603. max(len(slot_sources) - len(cycles), 0),
  2604. )
  2605. @staticmethod
  2606. def _assemble_first_cycle_window(
  2607. device_part: str,
  2608. device_points: list[str],
  2609. primary: str,
  2610. points: list[dict[str, Any]],
  2611. series_data: dict[str, list[dict[str, Any]]],
  2612. second_series_data: list[dict[str, Any]],
  2613. angle_data: list[dict[str, Any]],
  2614. volume_data: list[dict[str, Any]],
  2615. volume_info: dict[str, Any] | None,
  2616. second_finite_count: int,
  2617. second_non_zero_count: int,
  2618. second_min: float | None,
  2619. second_max: float | None,
  2620. cycles: list[dict[str, Any]],
  2621. triggers: list[float],
  2622. files: list[dict[str, Any]],
  2623. diagnostics: list[dict[str, Any]],
  2624. missing_cycles: int = 0,
  2625. ) -> dict[str, Any]:
  2626. for device_point in series_data:
  2627. series_data[device_point].sort(key=lambda item: item["x"])
  2628. second_series_data.sort(key=lambda item: item["x"])
  2629. angle_data.sort(key=lambda item: item["x"])
  2630. volume_data.sort(key=lambda item: item["x"])
  2631. extents = {device_point: _series_extent(series_data[device_point]) for device_point in device_points}
  2632. return {
  2633. "devicePart": device_part,
  2634. "devicePoints": device_points,
  2635. "primaryPoint": primary,
  2636. "points": points,
  2637. "xMin": 0,
  2638. "xMax": len(points),
  2639. "series": [
  2640. {
  2641. "devicePoint": device_point,
  2642. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  2643. "color": MEASUREMENT_COLORS[DEVICE_POINT_TO_TYPE[device_point]],
  2644. "data": series_data[device_point],
  2645. "min": extents[device_point][0],
  2646. "max": extents[device_point][1],
  2647. }
  2648. for device_point in device_points
  2649. ],
  2650. "secondSeries": {
  2651. "name": "周期数据",
  2652. "color": "#f56c6c",
  2653. "sourceMeasurementType": DEVICE_POINT_TO_TYPE[primary],
  2654. "sourceDevicePoint": primary,
  2655. "data": second_series_data,
  2656. "finiteCount": second_finite_count,
  2657. "nonZeroCount": second_non_zero_count,
  2658. "min": second_min,
  2659. "max": second_max,
  2660. },
  2661. "angleSeries": {
  2662. "color": "#d59b2b",
  2663. "data": angle_data,
  2664. },
  2665. "volumeSeries": {
  2666. "color": "#4d9e6f",
  2667. "data": volume_data,
  2668. "info": volume_info,
  2669. },
  2670. "cycles": cycles,
  2671. "triggerXs": sorted(set(triggers)),
  2672. "files": files,
  2673. "diagnostics": diagnostics,
  2674. "firstCycleMode": True,
  2675. "firstCycleNotice": (
  2676. f"窗口内 {missing_cycles} 个文件没有可用的单周期采样数据(wave_sample_one),"
  2677. "多为停机(rpm=0)文件,无法绘制首周期曲线。"
  2678. if missing_cycles > 0
  2679. else None
  2680. ),
  2681. }
  2682. @staticmethod
  2683. def _build_volume_vector(
  2684. sample_count: int,
  2685. cycles: list[DetectedCycle],
  2686. point_name: str,
  2687. ) -> tuple[np.ndarray, dict[str, Any] | None]:
  2688. volume = np.full(sample_count, np.nan, dtype=float)
  2689. cylinder_name = next(
  2690. (name for name in CYLINDER_BORE_MM if name in point_name),
  2691. None,
  2692. )
  2693. if cylinder_name is None or not cycles:
  2694. return volume, None
  2695. bore_mm = CYLINDER_BORE_MM[cylinder_name]
  2696. clearance = CLEARANCE_VOLUME_L_BY_BORE[bore_mm]
  2697. crank_radius = PISTON_STROKE_MM / 2.0
  2698. piston_area = np.pi * (bore_mm / 2.0) ** 2
  2699. for cycle in cycles:
  2700. angle_rad = np.deg2rad(cycle.angle)
  2701. travel = (
  2702. crank_radius * (1.0 - np.cos(angle_rad))
  2703. + CONNECTING_ROD_LENGTH_MM
  2704. - np.sqrt(
  2705. CONNECTING_ROD_LENGTH_MM**2
  2706. - (crank_radius * np.sin(angle_rad)) ** 2,
  2707. )
  2708. )
  2709. volume[cycle.start_offset : cycle.end_offset] = (
  2710. clearance + piston_area * travel / 1_000_000.0
  2711. )
  2712. finite = volume[np.isfinite(volume)]
  2713. return volume, {
  2714. "cylinder": cylinder_name,
  2715. "boreMm": bore_mm,
  2716. "clearanceVolumeL": clearance,
  2717. "minVolumeL": float(np.min(finite)) if len(finite) else None,
  2718. "maxVolumeL": float(np.max(finite)) if len(finite) else None,
  2719. }
  2720. @staticmethod
  2721. def _build_period_detail(
  2722. metadata: dict[str, Any],
  2723. samples: np.ndarray,
  2724. period_number: int,
  2725. ) -> dict[str, Any]:
  2726. stored_cycle = _stored_first_cycle(samples) if metadata.get("single_cycle") else None
  2727. if stored_cycle is not None:
  2728. detected = [stored_cycle]
  2729. diagnostics = {"completeCycleCount": 1, "storedFirstCycle": 1}
  2730. else:
  2731. detected, diagnostics = detect_cycles(samples)
  2732. cycle = next((item for item in detected if item.number == period_number), None)
  2733. if cycle is None:
  2734. raise ValueError(f"没有找到周期 {period_number}")
  2735. angles360 = np.linspace(0.0, 359.0, 360)
  2736. pressure = np.interp(angles360, cycle.angle, cycle.signal)
  2737. display_angles = np.where(angles360 <= 180.0, angles360, 360.0 - angles360)
  2738. point_name = str(metadata.get("point_name") or "")
  2739. cylinder_name = next(
  2740. (name for name in CYLINDER_BORE_MM if name in point_name),
  2741. None,
  2742. )
  2743. volume: np.ndarray | None = None
  2744. volume_info: dict[str, Any] | None = None
  2745. if cylinder_name:
  2746. bore_mm = CYLINDER_BORE_MM[cylinder_name]
  2747. clearance_volume = CLEARANCE_VOLUME_L_BY_BORE[bore_mm]
  2748. angle_rad = np.deg2rad(angles360)
  2749. crank_radius = PISTON_STROKE_MM / 2.0
  2750. piston_travel = (
  2751. crank_radius * (1.0 - np.cos(angle_rad))
  2752. + CONNECTING_ROD_LENGTH_MM
  2753. - np.sqrt(
  2754. CONNECTING_ROD_LENGTH_MM**2
  2755. - (crank_radius * np.sin(angle_rad)) ** 2,
  2756. )
  2757. )
  2758. piston_area = np.pi * (bore_mm / 2.0) ** 2
  2759. volume = clearance_volume + piston_area * piston_travel / 1_000_000.0
  2760. volume_info = {
  2761. "cylinder": cylinder_name,
  2762. "boreMm": bore_mm,
  2763. "clearanceVolumeL": clearance_volume,
  2764. "minVolumeL": float(np.min(volume)),
  2765. "maxVolumeL": float(np.max(volume)),
  2766. }
  2767. start_index = int(samples[cycle.start_offset, 0])
  2768. end_offset = min(cycle.end_offset, len(samples) - 1)
  2769. end_index = int(samples[max(cycle.end_offset - 1, cycle.start_offset), 0])
  2770. return {
  2771. "waveFile": {
  2772. "id": int(metadata["id"]),
  2773. "pointName": point_name,
  2774. "measurementType": metadata.get("measurement_type"),
  2775. "sampleTime": _time_string(metadata.get("sample_time")),
  2776. "sampleFrequencyHz": int(metadata.get("sample_frequency_hz") or 0),
  2777. "sampleCount": int(metadata.get("sample_count") or len(samples)),
  2778. },
  2779. "period": {
  2780. "periodNo": cycle.number,
  2781. "startSampleIndex": start_index,
  2782. "endSampleIndex": end_index,
  2783. "sampleCount": int(cycle.end_offset - cycle.start_offset),
  2784. "triggerSampleIndices": [
  2785. int(samples[offset, 0])
  2786. for offset in cycle.trigger_offsets
  2787. if 0 <= offset < len(samples)
  2788. ],
  2789. },
  2790. "angles": display_angles.tolist(),
  2791. "angles360": angles360.tolist(),
  2792. "pressure": pressure.tolist(),
  2793. "volume": volume.tolist() if volume is not None else None,
  2794. "volumeInfo": volume_info,
  2795. "phases": [
  2796. {
  2797. "name": name,
  2798. "color": color,
  2799. "start": start,
  2800. "end": end,
  2801. }
  2802. for name, color, start, end in PHASES
  2803. ],
  2804. "diagnostics": diagnostics,
  2805. }
  2806. @staticmethod
  2807. def _load_for_window(
  2808. loader: Callable[..., tuple[dict[str, Any], np.ndarray]],
  2809. file_id: int,
  2810. measurement_type: str,
  2811. point_name: str,
  2812. sample_time: str,
  2813. load_cache: dict[int, tuple[dict[str, Any], np.ndarray]],
  2814. ) -> tuple[dict[str, Any], np.ndarray]:
  2815. if file_id not in load_cache:
  2816. if loader.__name__ == "_load_demo_wave":
  2817. load_cache[file_id] = loader(file_id, measurement_type, point_name, sample_time)
  2818. else:
  2819. load_cache[file_id] = loader(file_id)
  2820. return load_cache[file_id]
  2821. def _load_db_wave(self, file_id: int) -> tuple[dict[str, Any], np.ndarray]:
  2822. with get_connection() as connection:
  2823. with connection.cursor() as cursor:
  2824. cursor.execute(
  2825. """
  2826. SELECT id, point_name, measurement_type, sample_frequency_hz,
  2827. sample_count, sample_time, rpm, file_name, tspluse_status,
  2828. cycle_start, cycle_end
  2829. FROM wave_file
  2830. WHERE id = %s
  2831. """,
  2832. (file_id,),
  2833. )
  2834. metadata = cursor.fetchone()
  2835. if metadata is None:
  2836. raise ValueError(f"wave_file.id={file_id} 不存在")
  2837. # wave_sample_one 只保留第一个周期,直接读取该文件全部样本。
  2838. cursor.execute(
  2839. """
  2840. SELECT sample_index, signal_value, second_value
  2841. FROM wave_sample_one
  2842. WHERE wave_file_id = %s
  2843. ORDER BY sample_index ASC
  2844. """,
  2845. (file_id,),
  2846. )
  2847. rows = cursor.fetchall()
  2848. if not rows:
  2849. raise ValueError(f"wave_file.id={file_id} 没有采样数据")
  2850. metadata["single_cycle"] = True
  2851. samples = np.asarray(
  2852. [
  2853. (
  2854. float(row["sample_index"]),
  2855. float(row["signal_value"]),
  2856. float(row["second_value"]) if row["second_value"] is not None else np.nan,
  2857. )
  2858. for row in rows
  2859. ],
  2860. dtype=float,
  2861. )
  2862. return metadata, samples
  2863. @staticmethod
  2864. @lru_cache(maxsize=24)
  2865. def _demo_samples(file_id: int, measurement_type: str) -> np.ndarray:
  2866. count = DEMO_SAMPLE_COUNT
  2867. index = np.arange(count, dtype=float)
  2868. revolution = DEMO_REVOLUTION_SAMPLES
  2869. phase = (index % revolution) / revolution * 2 * np.pi
  2870. second = np.zeros(count, dtype=float)
  2871. for revolution_start in range(0, count, revolution):
  2872. for pulse in range(PULSES_PER_REVOLUTION):
  2873. pulse_start = revolution_start + int(round(pulse * revolution / PULSES_PER_REVOLUTION))
  2874. width = 22 if pulse == 0 else 8
  2875. pulse_end = min(count, pulse_start + width)
  2876. second[pulse_start:pulse_end] = 40.0
  2877. variation = (file_id % 17) / 17.0
  2878. if measurement_type == "压力":
  2879. signal = (
  2880. 4.2
  2881. + 1.8 * np.sin(phase - 0.4)
  2882. + 0.55 * np.sin(2 * phase + variation)
  2883. + 0.22 * np.sin(7 * phase)
  2884. )
  2885. signal += 0.2 * np.maximum(np.sin(phase - 0.2), 0) ** 5
  2886. elif measurement_type == "位移":
  2887. signal = 0.5 + 0.18 * np.cos(phase) + 0.035 * np.sin(3 * phase + variation)
  2888. else:
  2889. signal = 0.15 * np.sin(phase * 2 + variation) + 0.04 * np.sin(11 * phase)
  2890. signal += 0.018 * np.cos(index / 37.0)
  2891. return np.column_stack((index, signal, second))
  2892. def _load_demo_wave(
  2893. self,
  2894. file_id: int,
  2895. measurement_type: str,
  2896. point_name: str,
  2897. sample_time: str,
  2898. ) -> tuple[dict[str, Any], np.ndarray]:
  2899. samples = self._demo_samples(file_id, measurement_type)
  2900. metadata = {
  2901. "id": file_id,
  2902. "point_name": point_name,
  2903. "measurement_type": measurement_type,
  2904. "sample_frequency_hz": 25600,
  2905. "sample_count": len(samples),
  2906. "sample_time": sample_time,
  2907. "rpm": 998.0,
  2908. "file_name": f"demo-{file_id}.dat",
  2909. }
  2910. return metadata, samples
  2911. @staticmethod
  2912. def _demo_options() -> list[dict[str, Any]]:
  2913. end = DEMO_START + timedelta(minutes=5 * (DEMO_POINT_COUNT - 1))
  2914. return [
  2915. {
  2916. "devicePart": DEMO_DEVICE_PART,
  2917. "devicePoint": device_point,
  2918. "measurementType": DEVICE_POINT_TO_TYPE[device_point],
  2919. "minTime": _time_string(DEMO_START),
  2920. "maxTime": _time_string(end),
  2921. "fileCount": DEMO_POINT_COUNT,
  2922. }
  2923. for device_point in DEVICE_POINTS
  2924. ]
  2925. @staticmethod
  2926. def _demo_time_points(
  2927. device_part: str,
  2928. device_points: list[str],
  2929. start: datetime | None,
  2930. end: datetime | None,
  2931. ) -> list[dict[str, Any]]:
  2932. points = []
  2933. for index in range(DEMO_POINT_COUNT):
  2934. timestamp = DEMO_START + timedelta(minutes=5 * index)
  2935. if start and timestamp < start:
  2936. continue
  2937. if end and timestamp > end:
  2938. continue
  2939. files = {}
  2940. for device_point in device_points:
  2941. measurement_type = DEVICE_POINT_TO_TYPE[device_point]
  2942. files[device_point] = {
  2943. "id": DEMO_ID_BY_TYPE[measurement_type] + index,
  2944. "devicePoint": device_point,
  2945. "measurementType": measurement_type,
  2946. "sampleCount": DEMO_SAMPLE_COUNT,
  2947. "sampleFrequencyHz": 25600,
  2948. "rpm": 998.0,
  2949. }
  2950. points.append(
  2951. {
  2952. "index": len(points),
  2953. "sampleTime": _time_string(timestamp),
  2954. "files": files,
  2955. },
  2956. )
  2957. return points
  2958. data_service = DataService()