scan_pks_fault_evolution.py 80 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605
  1. #!/usr/bin/env python3
  2. """Build long-term oil-pressure and short-term vibration evolution records.
  3. This is a new, read-only scanner. It deliberately ignores inlet and exhaust
  4. pressure. The scanner reads the alarm configuration from site_point for each
  5. unit by AlarmType rather than assuming that AlarmLimit1 has the same meaning
  6. on every machine.
  7. """
  8. from __future__ import annotations
  9. import argparse
  10. import csv
  11. import math
  12. import os
  13. import sys
  14. from collections import defaultdict
  15. from datetime import datetime, timedelta
  16. from pathlib import Path
  17. from statistics import median
  18. import pymysql
  19. sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
  20. from app.config import settings # noqa: E402
  21. BATCH_TO_UNIT = {30: "7号机", 31: "8号机", 32: "9号机"}
  22. EXPECTED_PER_HOUR = 720
  23. MIN_HOUR_COVERAGE = 0.80
  24. MIN_DAILY_HOURS = 4
  25. MIN_HISTORY_HOURS = 24
  26. BASELINE_DAYS = 3
  27. RESET_GAP_HOURS = 12
  28. RESET_RISE_ABS = 0.04
  29. RESET_RISE_REL = 0.10
  30. HOUR_ALARM_SAMPLE_RATIO = 0.20
  31. HOURLY_BASELINE_HOURS = 24
  32. HOURLY_TREND_HOURS = 6
  33. OIL_CONFIRM_HOURS = 3
  34. OIL_DOWNGRADE_HOURS = 2
  35. OIL_RECOVER_HOURS = 3
  36. VIBRATION_CONFIRM_HOURS = 3
  37. VIBRATION_MERGE_GAP_HOURS = 6
  38. VIBRATION_RESTART_GAP_HOURS = 12
  39. SITE_POINT_COLUMNS = (
  40. "ItemName, ItemDescription, Units, "
  41. "AlarmType1, AlarmType2, AlarmType3, AlarmType4, "
  42. "AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4"
  43. )
  44. HOURLY_SQL = """
  45. SELECT
  46. FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start,
  47. COUNT(*) AS samples,
  48. SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples,
  49. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_avg,
  50. MIN(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_min,
  51. MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_max,
  52. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_avg,
  53. MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_max,
  54. AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_avg,
  55. MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_max,
  56. SUM(CASE WHEN YSJ_41 > 0 AND YSJ_5 <= %s THEN 1 ELSE 0 END) AS oil_low_count,
  57. SUM(CASE WHEN YSJ_41 > 0 AND YSJ_5 <= %s THEN 1 ELSE 0 END) AS oil_low_low_count,
  58. SUM(CASE WHEN YSJ_41 > 0 AND YSJ_10 >= %s THEN 1 ELSE 0 END) AS coupling_high_count,
  59. SUM(CASE WHEN YSJ_41 > 0 AND YSJ_10 >= %s THEN 1 ELSE 0 END) AS coupling_high_high_count,
  60. SUM(CASE WHEN YSJ_41 > 0 AND YSJ_11 >= %s THEN 1 ELSE 0 END) AS chain_high_count,
  61. SUM(CASE WHEN YSJ_41 > 0 AND YSJ_11 >= %s THEN 1 ELSE 0 END) AS chain_high_high_count
  62. FROM pks_long_sample
  63. WHERE import_batch_id = %s
  64. AND sample_time >= %s
  65. AND sample_time < %s
  66. GROUP BY FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600)
  67. ORDER BY hour_start
  68. """
  69. NUMERIC_HOURLY_FIELDS = (
  70. "oil_avg", "oil_min", "oil_max", "coupling_avg", "coupling_max",
  71. "chain_avg", "chain_max", "oil_low_count", "oil_low_low_count",
  72. "coupling_high_count", "coupling_high_high_count", "chain_high_count",
  73. "chain_high_high_count",
  74. )
  75. def connect():
  76. return pymysql.connect(
  77. host=settings.db_host,
  78. port=settings.db_port,
  79. user=settings.db_user,
  80. password=settings.db_password,
  81. database=settings.db_name,
  82. charset="utf8mb4",
  83. cursorclass=pymysql.cursors.DictCursor,
  84. connect_timeout=settings.db_connect_timeout,
  85. read_timeout=600,
  86. write_timeout=60,
  87. autocommit=True,
  88. )
  89. def parse_dt(value: str) -> datetime:
  90. return datetime.strptime(value, "%Y-%m-%d %H:%M:%S")
  91. def chunks(start: datetime, end: datetime, days: int):
  92. cursor = start
  93. while cursor < end:
  94. nxt = min(cursor + timedelta(days=days), end)
  95. yield cursor, nxt
  96. cursor = nxt
  97. def number(value):
  98. if value is None:
  99. return None
  100. try:
  101. value = float(value)
  102. except (TypeError, ValueError):
  103. return None
  104. return value if math.isfinite(value) else None
  105. def alarm_entry(item: dict, alarm_type: str) -> tuple[float | None, str]:
  106. for index in range(1, 5):
  107. configured_type = item.get(f"AlarmType{index}")
  108. if configured_type is not None and str(configured_type).strip() == alarm_type:
  109. return number(item.get(f"AlarmLimit{index}")), f"AlarmType{index}/AlarmLimit{index}"
  110. return None, ""
  111. def load_alarm_config() -> dict[int, dict]:
  112. """Read alarm limits by AlarmType, because AlarmLimit positions differ."""
  113. result = {}
  114. connection = connect()
  115. try:
  116. with connection.cursor() as cursor:
  117. names = [f"YSJ{unit}_{point}" for unit in (7, 8, 9) for point in (5, 10, 11)]
  118. placeholders = ",".join(["%s"] * len(names))
  119. cursor.execute(
  120. f"SELECT {SITE_POINT_COLUMNS} FROM site_point WHERE ItemName IN ({placeholders})",
  121. names,
  122. )
  123. items = {row["ItemName"]: row for row in cursor.fetchall()}
  124. finally:
  125. connection.close()
  126. for batch, unit in BATCH_TO_UNIT.items():
  127. unit_number = unit[0]
  128. oil = items.get(f"YSJ{unit_number}_5")
  129. coupling = items.get(f"YSJ{unit_number}_10")
  130. chain = items.get(f"YSJ{unit_number}_11")
  131. if not oil or not coupling or not chain:
  132. raise RuntimeError(f"site_point 缺少 {unit} 的 YSJ_5/10/11 配置")
  133. oil_entries = {
  134. "low": alarm_entry(oil, "PVLow"),
  135. "low_low": alarm_entry(oil, "PVLowLow"),
  136. "high": alarm_entry(oil, "PVHigh"),
  137. "high_high": alarm_entry(oil, "PVHighHigh"),
  138. }
  139. coupling_entries = {
  140. "high": alarm_entry(coupling, "PVHigh"),
  141. "high_high": alarm_entry(coupling, "PVHighHigh"),
  142. }
  143. chain_entries = {
  144. "high": alarm_entry(chain, "PVHigh"),
  145. "high_high": alarm_entry(chain, "PVHighHigh"),
  146. }
  147. oil_limits = {key: value for key, (value, _source) in oil_entries.items()}
  148. coupling_limits = {key: value for key, (value, _source) in coupling_entries.items()}
  149. chain_limits = {key: value for key, (value, _source) in chain_entries.items()}
  150. if any(value is None for value in oil_limits.values()):
  151. raise RuntimeError(f"{unit} YSJ_5 缺少完整低压报警阈值")
  152. if any(value is None for value in (*coupling_limits.values(), *chain_limits.values())):
  153. raise RuntimeError(f"{unit} YSJ_10/11 缺少完整振动报警阈值")
  154. result[batch] = {
  155. "unit": unit,
  156. "batch": batch,
  157. "oil": oil_limits,
  158. "coupling": coupling_limits,
  159. "chain": chain_limits,
  160. "alarm_sources": {
  161. "oil": {key: source for key, (_value, source) in oil_entries.items()},
  162. "coupling": {key: source for key, (_value, source) in coupling_entries.items()},
  163. "chain": {key: source for key, (_value, source) in chain_entries.items()},
  164. },
  165. "metadata": {
  166. "oil_description": oil.get("ItemDescription") or "",
  167. "oil_units": oil.get("Units") or "",
  168. "coupling_description": coupling.get("ItemDescription") or "",
  169. "coupling_units": coupling.get("Units") or "",
  170. "chain_description": chain.get("ItemDescription") or "",
  171. "chain_units": chain.get("Units") or "",
  172. },
  173. }
  174. return result
  175. def query_params(batch: int, start: datetime, end: datetime, config: dict) -> tuple:
  176. return (
  177. config["oil"]["low"], config["oil"]["low_low"],
  178. config["coupling"]["high"], config["coupling"]["high_high"],
  179. config["chain"]["high"], config["chain"]["high_high"],
  180. batch, start, end,
  181. )
  182. def fetch_hourly(batch: int, start: datetime, end: datetime, chunk_days: int, config: dict):
  183. rows = []
  184. connection = connect()
  185. try:
  186. with connection.cursor() as cursor:
  187. for index, (chunk_start, chunk_end) in enumerate(chunks(start, end, chunk_days), 1):
  188. params = query_params(batch, chunk_start, chunk_end, config)
  189. cursor.execute("EXPLAIN " + HOURLY_SQL, params)
  190. plan = cursor.fetchone()
  191. if not plan or plan.get("key") not in ("PRIMARY",):
  192. raise RuntimeError(
  193. f"安全检查失败:batch={batch} chunk={chunk_start} "
  194. f"EXPLAIN 未使用 PRIMARY: {plan}"
  195. )
  196. cursor.execute(HOURLY_SQL, params)
  197. chunk_count = 0
  198. for raw in cursor.fetchall():
  199. row = {
  200. "批次": batch,
  201. "机组": config["unit"],
  202. "小时": raw["hour_start"].strftime("%Y-%m-%d %H:%M:%S"),
  203. "样本数": int(raw["samples"] or 0),
  204. "运行样本数": int(raw["running_samples"] or 0),
  205. }
  206. row["运行覆盖率"] = row["运行样本数"] / EXPECTED_PER_HOUR
  207. for field in NUMERIC_HOURLY_FIELDS:
  208. row[field] = number(raw[field])
  209. row["_dt"] = raw["hour_start"]
  210. rows.append(row)
  211. chunk_count += 1
  212. print(
  213. f"batch={batch} chunk={index} {chunk_start}..{chunk_end}: "
  214. f"{chunk_count:,} hourly rows"
  215. )
  216. finally:
  217. connection.close()
  218. return rows
  219. def median_value(values):
  220. values = [value for value in values if value is not None]
  221. return median(values) if values else None
  222. def max_consecutive(values):
  223. longest = 0
  224. current = 0
  225. for value in values:
  226. if value:
  227. current += 1
  228. longest = max(longest, current)
  229. else:
  230. current = 0
  231. return longest
  232. def robust_baseline(values, minimum: int = MIN_HISTORY_HOURS):
  233. values = [value for value in values if value is not None]
  234. if len(values) < minimum:
  235. return None, None
  236. center = median(values)
  237. mad = median(abs(value - center) for value in values)
  238. # Keep a meaningful relative scale when a sensor is very stable.
  239. scale = max(1.4826 * mad, abs(center) * 0.05, 1e-6)
  240. return center, scale
  241. def linear_slope(values):
  242. pairs = [(index, value) for index, value in enumerate(values) if value is not None]
  243. if len(pairs) < 3:
  244. return None
  245. xbar = sum(x for x, _ in pairs) / len(pairs)
  246. ybar = sum(y for _, y in pairs) / len(pairs)
  247. denominator = sum((x - xbar) ** 2 for x, _ in pairs)
  248. if denominator == 0:
  249. return None
  250. return sum((x - xbar) * (y - ybar) for x, y in pairs) / denominator
  251. def consecutive_decline_hours(history: list[dict], current: dict) -> int:
  252. """Count contiguous valid hours that have not risen from the prior hour."""
  253. count = 1
  254. newer = current
  255. for older in reversed(history):
  256. if newer["_dt"] - older["_dt"] != timedelta(hours=1):
  257. break
  258. if newer["oil_avg"] is None or older["oil_avg"] is None or newer["oil_avg"] > older["oil_avg"]:
  259. break
  260. count += 1
  261. newer = older
  262. return count
  263. def consecutive_flag_hours(history: list[dict], current: dict, field: str) -> int:
  264. count = 1 if current.get(field) else 0
  265. if not count:
  266. return 0
  267. newer = current
  268. for older in reversed(history):
  269. if newer["_dt"] - older["_dt"] != timedelta(hours=1) or not older.get(field):
  270. break
  271. count += 1
  272. newer = older
  273. return count
  274. def build_daily(rows: list[dict], configs: dict[int, dict]):
  275. grouped = defaultdict(list)
  276. for row in rows:
  277. if row["运行覆盖率"] >= MIN_HOUR_COVERAGE and row["oil_avg"] is not None:
  278. grouped[(row["批次"], row["_dt"].date())].append(row)
  279. daily = []
  280. for (batch, day), values in sorted(grouped.items()):
  281. if len(values) < MIN_DAILY_HOURS:
  282. continue
  283. config = configs[batch]
  284. running_samples = sum(row["运行样本数"] for row in values)
  285. low_count = sum(row["oil_low_count"] or 0 for row in values)
  286. low_low_count = sum(row["oil_low_low_count"] or 0 for row in values)
  287. oil_values = [row["oil_avg"] for row in values]
  288. low_hours = [
  289. (row["oil_low_count"] or 0) / max(row["运行样本数"], 1) >= HOUR_ALARM_SAMPLE_RATIO
  290. for row in values
  291. ]
  292. low_low_hours = [
  293. (row["oil_low_low_count"] or 0) / max(row["运行样本数"], 1) >= HOUR_ALARM_SAMPLE_RATIO
  294. for row in values
  295. ]
  296. daily.append({
  297. "批次": batch,
  298. "机组": config["unit"],
  299. "日期": day.isoformat(),
  300. "_date": day,
  301. "有效运行小时数": len(values),
  302. "日运行覆盖率": running_samples / (EXPECTED_PER_HOUR * 24),
  303. "日油压中位数": median(oil_values),
  304. "日油压最小值": min(row["oil_min"] for row in values if row["oil_min"] is not None),
  305. "日油压最大值": max(row["oil_max"] for row in values if row["oil_max"] is not None),
  306. # Raw sample ratios are retained for diagnosis. Alarm hours use
  307. # a 20% running-sample gate so isolated 5-second drops do not
  308. # become a full-day severe pressure phase.
  309. "低报警小时数": sum(low_hours),
  310. "低低报警小时数": sum(low_low_hours),
  311. "低报警最长连续小时数": max_consecutive(low_hours),
  312. "低低报警最长连续小时数": max_consecutive(low_low_hours),
  313. "低报警样本数": low_count,
  314. "低低报警样本数": low_low_count,
  315. "低报警运行样本比例": low_count / max(running_samples, 1),
  316. "低低报警运行样本比例": low_low_count / max(running_samples, 1),
  317. "低报警阈值": config["oil"]["low"],
  318. "低低报警阈值": config["oil"]["low_low"],
  319. })
  320. return daily
  321. def assign_cycles(daily: list[dict]):
  322. by_batch = defaultdict(list)
  323. for row in daily:
  324. by_batch[row["批次"]].append(row)
  325. for batch, series in by_batch.items():
  326. series.sort(key=lambda row: row["_date"])
  327. cycle_number = 0
  328. previous = None
  329. for row in series:
  330. reset_reason = ""
  331. if previous is None:
  332. cycle_number += 1
  333. reset_reason = "扫描范围内首个有效运行日"
  334. else:
  335. gap_days = (row["_date"] - previous["_date"]).days
  336. rise = row["日油压中位数"] - previous["日油压中位数"]
  337. rise_rel = rise / max(abs(previous["日油压中位数"]), 1e-6)
  338. if gap_days > RESET_GAP_DAYS:
  339. cycle_number += 1
  340. reset_reason = f"有效运行间隔{gap_days}天,疑似重新开机周期"
  341. elif rise >= RESET_RISE_ABS and rise_rel >= RESET_RISE_REL:
  342. cycle_number += 1
  343. reset_reason = f"日油压回升{rise:.3f}MPa,疑似恢复后重新运行"
  344. if cycle_number == 0:
  345. cycle_number = 1
  346. row["周期编号"] = f"{row['机组']}-P{cycle_number:02d}"
  347. row["周期起点依据"] = reset_reason
  348. previous = row
  349. def classify_oil_daily(daily: list[dict]):
  350. by_cycle = defaultdict(list)
  351. for row in daily:
  352. by_cycle[row["周期编号"]].append(row)
  353. for series in by_cycle.values():
  354. series.sort(key=lambda row: row["_date"])
  355. baseline_days = series[:BASELINE_DAYS]
  356. cycle_baseline = median(row["日油压中位数"] for row in baseline_days)
  357. low_alarm = series[0]["低报警阈值"]
  358. low_low_alarm = series[0]["低低报警阈值"]
  359. for index, row in enumerate(series):
  360. row["周期基线油压"] = cycle_baseline
  361. relative = row["日油压中位数"] / cycle_baseline - 1 if cycle_baseline else None
  362. recent = series[max(0, index - 4): index + 1]
  363. recent_values = [item["日油压中位数"] for item in recent]
  364. slope = linear_slope(recent_values)
  365. decline_days = 1
  366. cursor = index
  367. while cursor > 0 and series[cursor]["日油压中位数"] <= series[cursor - 1]["日油压中位数"]:
  368. decline_days += 1
  369. cursor -= 1
  370. smoothed = median_value([item["日油压中位数"] for item in series[max(0, index - 2): index + 1]])
  371. row["相对基线变化"] = relative
  372. row["近5运行日油压斜率"] = slope
  373. row["连续下降运行日数"] = decline_days
  374. row["平滑油压"] = smoothed
  375. row["是否基线期"] = "是" if index < BASELINE_DAYS else "否"
  376. downward = (slope is not None and slope < -0.002) or decline_days >= 3
  377. recovery = index > 0 and row["日油压中位数"] > series[index - 1]["日油压中位数"] * 1.03
  378. # Long-term phase is based on the daily level and trend. Raw
  379. # alarm samples are retained separately so isolated 5-second
  380. # drops do not promote an otherwise normal day to severe.
  381. severe_by_absolute = row["日油压中位数"] <= low_low_alarm and (
  382. decline_days >= 2 or row["低低报警小时数"] >= 2
  383. )
  384. severe_by_evolution = (
  385. relative is not None
  386. and downward
  387. and relative <= -0.20
  388. and (decline_days >= 2 or row["日油压中位数"] <= low_alarm)
  389. )
  390. abnormal_by_absolute = row["日油压中位数"] <= low_alarm and downward
  391. abnormal_by_evolution = relative is not None and relative <= -0.10 and downward
  392. if index < BASELINE_DAYS:
  393. phase = "正常"
  394. elif severe_by_absolute or severe_by_evolution:
  395. phase = "严重异常"
  396. elif abnormal_by_absolute or abnormal_by_evolution:
  397. phase = "异常"
  398. elif relative is not None and relative <= -0.05 and downward:
  399. phase = "轻微"
  400. else:
  401. phase = "正常"
  402. row["压力阶段"] = phase
  403. short_alarm = (
  404. row["低报警样本数"] > 0
  405. and row["日油压中位数"] > low_alarm
  406. ) or (
  407. row["低低报警样本数"] > 0
  408. and row["日油压中位数"] > low_low_alarm
  409. )
  410. row["短时低压标记"] = "是" if short_alarm else "否"
  411. row["短时低压小时数"] = row["低报警小时数"]
  412. row["短时低低报警小时数"] = row["低低报警小时数"]
  413. reasons = []
  414. if severe_by_absolute:
  415. reasons.append("日油压达到低低报警区间且具有持续性")
  416. if severe_by_evolution:
  417. reasons.append("相对周期基线下降至少20%且仍在下降")
  418. if abnormal_by_absolute:
  419. reasons.append("日油压进入低报警区间且仍在下降")
  420. if abnormal_by_evolution:
  421. reasons.append("相对周期基线下降至少10%且仍在下降")
  422. if phase == "轻微":
  423. reasons.append("相对周期基线下降至少5%且仍在下降")
  424. if short_alarm:
  425. reasons.append("存在短时低压样本,但日油压未进入对应长期阶段")
  426. row["阶段判定依据"] = ";".join(reasons) or "周期基线范围内"
  427. if row["短时低压标记"] == "是" and phase == "正常":
  428. shape = "短时低压波动型"
  429. elif recovery:
  430. shape = "回升型"
  431. elif downward and decline_days >= 3:
  432. shape = "渐进下降型"
  433. else:
  434. shape = "稳定型"
  435. row["变化形态"] = shape
  436. row["是否低报警"] = "是" if row["低报警小时数"] > 0 else "否"
  437. row["是否低低报警"] = "是" if row["低低报警小时数"] > 0 else "否"
  438. def build_oil_stage_segments(daily: list[dict]):
  439. """Merge daily phase labels into reviewable pressure-evolution stages."""
  440. by_cycle = defaultdict(list)
  441. for row in daily:
  442. by_cycle[row["周期编号"]].append(row)
  443. segments = []
  444. for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_date"]):
  445. series.sort(key=lambda row: row["_date"])
  446. current = []
  447. for row in series:
  448. if not current or (
  449. row["压力阶段"] == current[-1]["压力阶段"]
  450. and (row["_date"] - current[-1]["_date"]).days <= RESET_GAP_DAYS
  451. ):
  452. current.append(row)
  453. else:
  454. segments.append(_build_oil_stage_segment(cycle_id, current, len(segments) + 1))
  455. current = [row]
  456. if current:
  457. segments.append(_build_oil_stage_segment(cycle_id, current, len(segments) + 1))
  458. # Stage numbers should restart within each pressure cycle.
  459. by_cycle_segments = defaultdict(list)
  460. for segment in segments:
  461. by_cycle_segments[segment["周期编号"]].append(segment)
  462. for cycle_segments in by_cycle_segments.values():
  463. for index, segment in enumerate(cycle_segments, 1):
  464. segment["阶段序号"] = index
  465. return sorted(segments, key=lambda row: (row["批次"], row["开始日期"], row["阶段序号"]))
  466. def _build_oil_stage_segment(cycle_id: str, series: list[dict], sequence: int):
  467. first = series[0]
  468. last = series[-1]
  469. phase = first["压力阶段"]
  470. return {
  471. "机组": first["机组"],
  472. "批次": first["批次"],
  473. "周期编号": cycle_id,
  474. "阶段序号": sequence,
  475. "压力阶段": phase,
  476. "开始日期": first["日期"],
  477. "结束日期": last["日期"],
  478. "持续自然日数": (last["_date"] - first["_date"]).days + 1,
  479. "有效运行日数": len(series),
  480. "阶段运行小时数": sum(row["有效运行小时数"] for row in series),
  481. "阶段低报警小时数": sum(row["低报警小时数"] for row in series),
  482. "阶段低低报警小时数": sum(row["低低报警小时数"] for row in series),
  483. "阶段短时低压小时数": sum(row["短时低压小时数"] for row in series),
  484. "周期基线油压": first["周期基线油压"],
  485. "阶段开始日油压": first["日油压中位数"],
  486. "阶段结束日油压": last["日油压中位数"],
  487. "阶段最低日油压": min(row["日油压中位数"] for row in series),
  488. "阶段开始相对基线变化": first["相对基线变化"],
  489. "阶段结束相对基线变化": last["相对基线变化"],
  490. "阶段最大连续下降运行日数": max(row["连续下降运行日数"] for row in series),
  491. "阶段变化形态": "渐进下降型" if any(row["变化形态"] == "渐进下降型" for row in series) else "稳定/短时波动型",
  492. "阶段触发依据": "; ".join(
  493. reason for reason in (first.get("周期起点依据", ""), last.get("变化形态", "")) if reason
  494. ),
  495. }
  496. def classify_oil_hourly(rows: list[dict], daily: list[dict], configs: dict[int, dict]):
  497. """Assign a causal pressure level to each valid running hour.
  498. The cycle baseline is the pressure baseline already established for the
  499. corresponding running cycle. Recent history is limited to prior valid
  500. hours, so this output can be used to implement an online hourly updater.
  501. """
  502. cycle_by_date = {(row["批次"], row["_date"]): row["周期编号"] for row in daily}
  503. daily_baseline = {row["周期编号"]: row["周期基线油压"] for row in daily}
  504. grouped = defaultdict(list)
  505. for row in rows:
  506. cycle_id = cycle_by_date.get((row["批次"], row["_dt"].date()))
  507. if row["运行覆盖率"] < MIN_HOUR_COVERAGE or row["oil_avg"] is None or cycle_id is None:
  508. continue
  509. item = dict(row)
  510. item["周期编号"] = cycle_id
  511. grouped[(row["批次"], cycle_id)].append(item)
  512. result = []
  513. classified_keys = set()
  514. for (batch, cycle_id), series in grouped.items():
  515. series.sort(key=lambda row: row["_dt"])
  516. config = configs[batch]
  517. cycle_base = daily_baseline[cycle_id]
  518. history = []
  519. for row in series:
  520. recent = history[-HOURLY_BASELINE_HOURS:]
  521. recent_base = median_value([item["oil_avg"] for item in recent])
  522. trend_values = [item["oil_avg"] for item in history[-(HOURLY_TREND_HOURS - 1):]] + [row["oil_avg"]]
  523. trend = linear_slope(trend_values)
  524. decline_hours = consecutive_decline_hours(history, row)
  525. cycle_relative = row["oil_avg"] / cycle_base - 1 if cycle_base else None
  526. recent_relative = row["oil_avg"] / recent_base - 1 if recent_base else None
  527. running = max(row["运行样本数"], 1)
  528. low_ratio = (row["oil_low_count"] or 0) / running
  529. low_low_ratio = (row["oil_low_low_count"] or 0) / running
  530. low_hour = low_ratio >= HOUR_ALARM_SAMPLE_RATIO
  531. low_low_hour = low_low_ratio >= HOUR_ALARM_SAMPLE_RATIO
  532. low_low_streak = consecutive_flag_hours(history, {**row, "_low_low_hour": low_low_hour}, "_low_low_hour")
  533. downward = (
  534. decline_hours >= 3
  535. or (trend is not None and trend < -0.0005)
  536. or (recent_relative is not None and recent_relative <= -0.03 and trend is not None and trend < 0)
  537. )
  538. low_alarm = config["oil"]["low"]
  539. low_low_alarm = config["oil"]["low_low"]
  540. severe = (
  541. low_low_streak >= 3
  542. ) or (
  543. cycle_relative is not None and cycle_relative <= -0.20
  544. and (downward or row["oil_avg"] <= low_alarm)
  545. )
  546. abnormal = not severe and (
  547. (row["oil_avg"] <= low_alarm and (low_hour or decline_hours >= 2))
  548. or (cycle_relative is not None and cycle_relative <= -0.10 and downward)
  549. )
  550. mild = not severe and not abnormal and (
  551. (cycle_relative is not None and cycle_relative <= -0.05 and downward)
  552. or (recent_relative is not None and recent_relative <= -0.03 and downward)
  553. )
  554. if severe:
  555. level = "严重异常"
  556. reason = "相对周期基线下降至少20%并持续下降"
  557. if row["oil_avg"] <= low_low_alarm:
  558. reason = "小时油压达到低低报警区间"
  559. elif abnormal:
  560. level = "异常"
  561. reason = "相对周期基线下降至少10%并持续下降"
  562. if row["oil_avg"] <= low_alarm:
  563. reason = "小时油压进入低报警区间"
  564. elif mild:
  565. level = "轻微"
  566. reason = "相对周期基线下降至少5%并持续下降"
  567. else:
  568. level = "正常"
  569. reason = "周期基线范围内"
  570. short_low = (low_hour or low_low_hour) and level == "正常"
  571. if short_low:
  572. reason += ";存在短时低压样本,未形成长期下降等级"
  573. result.append({
  574. "机组": config["unit"],
  575. "批次": batch,
  576. "周期编号": cycle_id,
  577. "小时": row["小时"],
  578. "运行样本数": row["运行样本数"],
  579. "运行覆盖率": row["运行覆盖率"],
  580. "小时油压均值": row["oil_avg"],
  581. "小时油压最小值": row["oil_min"],
  582. "小时油压最大值": row["oil_max"],
  583. "周期基线油压": cycle_base,
  584. "近期24有效小时基线油压": recent_base,
  585. "相对周期基线变化": cycle_relative,
  586. "相对近期基线变化": recent_relative,
  587. "近6有效小时油压斜率": trend,
  588. "连续下降有效小时数": decline_hours,
  589. "历史有效运行小时数": len(history),
  590. "低报警样本比例": low_ratio,
  591. "低低报警样本比例": low_low_ratio,
  592. "低报警阈值": low_alarm,
  593. "低低报警阈值": low_low_alarm,
  594. "是否低报警小时": "是" if low_hour else "否",
  595. "是否低低报警小时": "是" if low_low_hour else "否",
  596. "压力等级": level,
  597. "短时低压标记": "是" if short_low else "否",
  598. "是否纳入等级判定": "是",
  599. "小时判定依据": reason,
  600. "_dt": row["_dt"],
  601. })
  602. classified_keys.add((batch, row["_dt"]))
  603. history.append(row)
  604. # Keep partial-running hours visible, but do not let them affect the
  605. # long-term level or its historical baseline.
  606. for row in rows:
  607. key = (row["批次"], row["_dt"])
  608. if row["运行样本数"] <= 0 or key in classified_keys:
  609. continue
  610. result.append({
  611. "机组": row["机组"],
  612. "批次": row["批次"],
  613. "周期编号": cycle_by_date.get((row["批次"], row["_dt"].date()), ""),
  614. "小时": row["小时"],
  615. "运行样本数": row["运行样本数"],
  616. "运行覆盖率": row["运行覆盖率"],
  617. "小时油压均值": row["oil_avg"],
  618. "小时油压最小值": row["oil_min"],
  619. "小时油压最大值": row["oil_max"],
  620. "周期基线油压": "",
  621. "近期24有效小时基线油压": "",
  622. "相对周期基线变化": "",
  623. "相对近期基线变化": "",
  624. "近6有效小时油压斜率": "",
  625. "连续下降有效小时数": "",
  626. "历史有效运行小时数": "",
  627. "低报警样本比例": (row["oil_low_count"] or 0) / max(row["运行样本数"], 1),
  628. "低低报警样本比例": (row["oil_low_low_count"] or 0) / max(row["运行样本数"], 1),
  629. "低报警阈值": configs[row["批次"]]["oil"]["low"],
  630. "低低报警阈值": configs[row["批次"]]["oil"]["low_low"],
  631. "是否低报警小时": "",
  632. "是否低低报警小时": "",
  633. "低低报警连续小时数": "",
  634. "压力等级": "数据不足",
  635. "短时低压标记": "",
  636. "正常波动标记": "",
  637. "是否纳入等级判定": "否",
  638. "小时判定依据": "运行覆盖率不足,未参与长期小时等级判定",
  639. "_dt": row["_dt"],
  640. })
  641. return sorted(result, key=lambda row: (row["批次"], row["_dt"]))
  642. def build_oil_hour_stage_segments(hourly_levels: list[dict]):
  643. """Merge confirmed levels and short unconfirmed triggers into segments."""
  644. grouped = defaultdict(list)
  645. for row in hourly_levels:
  646. if row["压力等级"] in ("轻微", "异常", "严重异常"):
  647. level = row["压力等级"]
  648. elif row["压力等级"] == "轻微" and row.get("是否形成确认等级") == "否":
  649. level = "轻微"
  650. else:
  651. continue
  652. grouped[(row["批次"], row["周期编号"], level)].append(row)
  653. segments = []
  654. for (batch, cycle_id, level), series in grouped.items():
  655. series.sort(key=lambda row: row["_dt"])
  656. current = []
  657. for row in series:
  658. if not current or (
  659. row["_dt"] - current[-1]["_dt"] <= timedelta(hours=VIBRATION_MERGE_GAP_HOURS)
  660. and row["压力等级"] == current[-1]["压力等级"]
  661. ):
  662. current.append(row)
  663. else:
  664. segments.append(_build_oil_hour_stage_segment(current, level))
  665. current = [row]
  666. if current:
  667. segments.append(_build_oil_hour_stage_segment(current, level))
  668. return sorted(segments, key=lambda row: (row["批次"], row["开始小时"]))
  669. def assign_hourly_cycles(rows: list[dict]):
  670. """Create oil-pressure cycles directly from valid running hours."""
  671. grouped = defaultdict(list)
  672. for row in rows:
  673. if row["运行覆盖率"] >= MIN_HOUR_COVERAGE and row["oil_avg"] is not None:
  674. grouped[row["批次"]].append(row)
  675. for batch, series in grouped.items():
  676. series.sort(key=lambda row: row["_dt"])
  677. cycle_number = 0
  678. previous = None
  679. for row in series:
  680. gap_hours = None if previous is None else (row["_dt"] - previous["_dt"]).total_seconds() / 3600
  681. # A one-hour pressure recovery is not a restart. It can be a
  682. # sensor dropout or part of the same fault episode. Start a new
  683. # cycle only after a long gap in valid running hours.
  684. if previous is None or (gap_hours is not None and gap_hours > RESET_GAP_HOURS):
  685. cycle_number += 1
  686. row["_cycle_id"] = f"{row['机组']}-P{cycle_number:02d}"
  687. previous = row
  688. return grouped
  689. def classify_oil_hourly_state_machine(rows: list[dict], configs: dict[int, dict]):
  690. """Classify oil pressure with upgrade/downgrade hysteresis per cycle."""
  691. cycles = assign_hourly_cycles(rows)
  692. result = []
  693. for batch, batch_series in cycles.items():
  694. config = configs[batch]
  695. by_cycle = defaultdict(list)
  696. for row in batch_series:
  697. by_cycle[row["_cycle_id"]].append(row)
  698. for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_dt"]):
  699. series.sort(key=lambda item: item["_dt"])
  700. cycle_base = median_value([row["oil_avg"] for row in series[:HOURLY_BASELINE_HOURS]])
  701. stable_values = [row["oil_avg"] for row in series[:HOURLY_BASELINE_HOURS]]
  702. stable_center = median_value(stable_values)
  703. stable_mad = median(abs(value - stable_center) for value in stable_values) if stable_center is not None else 0
  704. stable_scale = max(1.4826 * stable_mad, abs(stable_center or 0) * 0.01, 1e-6)
  705. history = []
  706. state = "正常"
  707. pending = None
  708. pending_count = 0
  709. downgrade_count = 0
  710. for row in series:
  711. recent = history[-HOURLY_BASELINE_HOURS:]
  712. recent_base = median_value([item["oil_avg"] for item in recent])
  713. trend_values = [item["oil_avg"] for item in history[-(HOURLY_TREND_HOURS - 1):]] + [row["oil_avg"]]
  714. trend = linear_slope(trend_values)
  715. decline_hours = consecutive_decline_hours(history, row)
  716. cycle_relative = row["oil_avg"] / cycle_base - 1 if cycle_base else None
  717. recent_relative = row["oil_avg"] / recent_base - 1 if recent_base else None
  718. running = max(row["运行样本数"], 1)
  719. low_ratio = (row["oil_low_count"] or 0) / running
  720. low_low_ratio = (row["oil_low_low_count"] or 0) / running
  721. low_hour = low_ratio >= HOUR_ALARM_SAMPLE_RATIO
  722. low_low_hour = low_low_ratio >= HOUR_ALARM_SAMPLE_RATIO
  723. low_low_streak = consecutive_flag_hours(
  724. history, {**row, "_low_low_hour": low_low_hour}, "_low_low_hour"
  725. )
  726. downward = decline_hours >= 3 or (trend is not None and trend < -0.0005)
  727. immediate = "严重预警" if low_low_hour else "低压预警" if low_hour else ""
  728. low_alarm = config["oil"]["low"]
  729. low_low_alarm = config["oil"]["low_low"]
  730. deviation_z = ((stable_center - row["oil_avg"]) / stable_scale) if stable_center is not None else 0.0
  731. if (
  732. (row["oil_avg"] <= low_low_alarm and low_low_streak >= 2)
  733. or (deviation_z >= 8 and decline_hours >= OIL_CONFIRM_HOURS)
  734. ):
  735. raw = "严重异常"
  736. elif (
  737. (row["oil_avg"] <= low_alarm and decline_hours >= OIL_CONFIRM_HOURS)
  738. or (deviation_z >= 4 and decline_hours >= OIL_CONFIRM_HOURS)
  739. ):
  740. raw = "异常"
  741. elif deviation_z >= 2 or low_hour:
  742. raw = "轻微"
  743. else:
  744. raw = "正常"
  745. rank = {"正常": 0, "轻微": 1, "异常": 2, "严重异常": 3}
  746. if rank[raw] > rank[state]:
  747. if pending == raw:
  748. pending_count += 1
  749. else:
  750. pending, pending_count = raw, 1
  751. required = 1 if raw == "轻微" else OIL_CONFIRM_HOURS
  752. if pending_count >= required:
  753. state = raw
  754. pending, pending_count = None, 0
  755. downgrade_count = 0
  756. elif rank[raw] < rank[state]:
  757. if pending == raw:
  758. pending_count += 1
  759. else:
  760. pending, pending_count = raw, 1
  761. required = OIL_DOWNGRADE_HOURS if raw != "正常" else OIL_RECOVER_HOURS
  762. if pending_count >= required:
  763. state = raw
  764. pending, pending_count = None, 0
  765. downgrade_count += 1
  766. elif raw == "正常" and state != "正常":
  767. downgrade_count += 1
  768. pending, pending_count = None, 0
  769. if downgrade_count >= 6:
  770. state = "正常"
  771. downgrade_count = 0
  772. else:
  773. pending, pending_count = None, 0
  774. downgrade_count = 0
  775. level_label = {"正常": "正常", "轻微": "轻微", "异常": "异常", "严重异常": "严重"}[state]
  776. reason = (
  777. f"当前油压{row['oil_avg']:.3f}MPa,近期基线{recent_base:.3f}MPa,"
  778. f"偏离稳定基线波动尺度{deviation_z:.1f}倍,连续下降{decline_hours}小时"
  779. if recent_base is not None else "近期基线不足"
  780. )
  781. if row["oil_avg"] <= low_low_alarm:
  782. reason += f",低于低低报警阈值{low_low_alarm:.3f}MPa"
  783. elif row["oil_avg"] <= low_alarm:
  784. reason += f",低于低报警阈值{low_alarm:.3f}MPa"
  785. if state != "正常":
  786. reason += f",{level_label}状态已连续确认"
  787. if low_low_streak >= 3:
  788. reason += f";低低报警连续{low_low_streak}小时"
  789. if immediate:
  790. reason += f";{immediate}"
  791. unconfirmed = state == "正常" and raw != "正常"
  792. if unconfirmed:
  793. reason += ";连续时间不足3小时,记录为轻微"
  794. output_level = "轻微" if unconfirmed else state
  795. result.append({
  796. "机组": config["unit"], "批次": batch, "周期编号": cycle_id,
  797. "小时": row["小时"], "运行样本数": row["运行样本数"], "运行覆盖率": row["运行覆盖率"],
  798. "小时油压均值": row["oil_avg"], "小时油压最小值": row["oil_min"], "小时油压最大值": row["oil_max"],
  799. "周期基线油压": cycle_base, "近期24有效小时基线油压": recent_base,
  800. "相对周期基线变化": cycle_relative, "相对近期基线变化": recent_relative,
  801. "近6有效小时油压斜率": trend, "连续下降有效小时数": decline_hours,
  802. "历史有效运行小时数": len(history), "低报警样本比例": low_ratio, "低低报警样本比例": low_low_ratio,
  803. "低报警阈值": config["oil"]["low"], "低低报警阈值": config["oil"]["low_low"],
  804. "是否低报警小时": "是" if low_hour else "否", "是否低低报警小时": "是" if low_low_hour else "否",
  805. "低低报警连续小时数": low_low_streak,
  806. "压力等级": output_level, "原始小时等级": raw, "立即报警": immediate,
  807. "正常波动标记": "否",
  808. "是否形成确认等级": "否" if unconfirmed else "是",
  809. "是否纳入等级判定": "是", "小时判定依据": reason, "_dt": row["_dt"],
  810. })
  811. row["_low_hour"] = low_hour
  812. row["_low_low_hour"] = low_low_hour
  813. history.append(row)
  814. return sorted(result, key=lambda row: (row["批次"], row["_dt"]))
  815. def classify_oil_hourly_segments(rows: list[dict], configs: dict[int, dict]):
  816. """Classify pressure by amplitude zone first, then confirm by duration.
  817. The configured low/low-low limits define the engineering zones. A short
  818. excursion is retained as mild; only a run of three or more consecutive
  819. valid hours in the low zone is confirmed abnormal, and the low-low zone is
  820. confirmed severe. No fixed 5/10/20 percent thresholds are used.
  821. """
  822. cycles = assign_hourly_cycles(rows)
  823. output = []
  824. for batch, batch_series in cycles.items():
  825. config = configs[batch]
  826. by_cycle = defaultdict(list)
  827. for row in batch_series:
  828. by_cycle[row["_cycle_id"]].append(row)
  829. for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_dt"]):
  830. series.sort(key=lambda item: item["_dt"])
  831. baseline_values = [row["oil_avg"] for row in series[:HOURLY_BASELINE_HOURS]]
  832. cycle_base = median_value(baseline_values)
  833. low = config["oil"]["low"]
  834. low_low = config["oil"]["low_low"]
  835. candidates = []
  836. history = []
  837. for row in series:
  838. current = row["oil_avg"]
  839. trend = linear_slope([item["oil_avg"] for item in history[-5:]] + [current])
  840. decline = consecutive_decline_hours(history, row)
  841. deviation = (cycle_base - current) / cycle_base if cycle_base else 0.0
  842. low_ratio = (row["oil_low_count"] or 0) / max(row["运行样本数"], 1)
  843. low_low_ratio = (row["oil_low_low_count"] or 0) / max(row["运行样本数"], 1)
  844. low_zone = current <= low
  845. low_low_zone = current <= low_low
  846. downward = (
  847. decline >= 3
  848. or (trend is not None and trend < -0.0005)
  849. )
  850. # Relative decline is a data-derived evolution signal. The
  851. # engineering alarm limits remain higher priority, while the
  852. # 5/10/20% bands describe progression before the alarm limit.
  853. severe_candidate = low_low_zone or (deviation >= 0.20 and downward)
  854. abnormal_candidate = low_zone or (deviation >= 0.10 and downward)
  855. mild_candidate = deviation >= 0.05 and downward
  856. if severe_candidate:
  857. raw = "严重异常"
  858. elif abnormal_candidate:
  859. raw = "异常"
  860. elif mild_candidate:
  861. raw = "轻微"
  862. else:
  863. raw = "正常"
  864. candidates.append({
  865. "row": row, "raw": raw, "trend": trend, "decline": decline,
  866. "deviation": deviation, "low_ratio": low_ratio, "low_low_ratio": low_low_ratio,
  867. "low_zone": low_zone, "low_low_zone": low_low_zone,
  868. })
  869. history.append(row)
  870. index = 0
  871. while index < len(candidates):
  872. raw = candidates[index]["raw"]
  873. if raw == "正常":
  874. final = "正常"
  875. end = index
  876. else:
  877. end = index
  878. while end + 1 < len(candidates) and candidates[end + 1]["raw"] == raw:
  879. end += 1
  880. run_length = end - index + 1
  881. # Alarm limits have highest priority for the hourly
  882. # candidate and immediate alarm fields. Long-term stage
  883. # labels still require persistence: a one/two-hour
  884. # low-low excursion is a short severe warning, not a
  885. # confirmed long-term severe stage.
  886. if run_length >= OIL_CONFIRM_HOURS:
  887. final = raw
  888. elif raw == "严重异常":
  889. final = "正常"
  890. else:
  891. final = "轻微"
  892. for item in candidates[index:end + 1]:
  893. row = item["row"]
  894. level_reason = {
  895. "轻微": "相对周期基线下降至少5%并持续下降,或普通异常持续不足3小时",
  896. "异常": "相对周期基线下降至少10%并持续下降,或低报警区间连续达到3小时",
  897. "严重异常": "相对周期基线下降至少20%并持续下降,或低低报警区间连续达到3小时",
  898. "正常": "未形成长期确认等级" if raw != "正常" else "未达到压力异常条件",
  899. }[final]
  900. if item["low_low_zone"]:
  901. level_reason += f";低低报警样本比例{item['low_low_ratio']:.1%}"
  902. elif item["low_zone"]:
  903. level_reason += f";低报警样本比例{item['low_ratio']:.1%}"
  904. output.append({
  905. "机组": config["unit"], "批次": batch, "周期编号": cycle_id,
  906. "小时": row["小时"], "运行样本数": row["运行样本数"], "运行覆盖率": row["运行覆盖率"],
  907. "小时油压均值": row["oil_avg"], "小时油压最小值": row["oil_min"], "小时油压最大值": row["oil_max"],
  908. "周期基线油压": cycle_base, "近期24有效小时基线油压": "",
  909. "相对周期基线变化": (row["oil_avg"] / cycle_base - 1) if cycle_base else None,
  910. "相对近期基线变化": "", "近6有效小时油压斜率": item["trend"],
  911. "连续下降有效小时数": item["decline"], "历史有效运行小时数": index,
  912. "低报警样本比例": item["low_ratio"], "低低报警样本比例": item["low_low_ratio"],
  913. "低报警阈值": low, "低低报警阈值": low_low,
  914. "是否低报警小时": "是" if item["low_zone"] else "否",
  915. "是否低低报警小时": "是" if item["low_low_zone"] else "否",
  916. "低低报警连续小时数": 0,
  917. "压力等级": final, "原始小时等级": raw,
  918. "立即报警": "严重预警" if item["low_low_zone"] else "低压预警" if item["low_zone"] else "",
  919. "正常波动标记": "是" if raw != "正常" and final == "轻微" else "否",
  920. "是否形成确认等级": "是" if final in ("异常", "严重异常") else "否",
  921. "小时判定依据": f"油压{row['oil_avg']:.3f}MPa,周期基线{cycle_base:.3f}MPa,"
  922. f"相对周期基线变化{((row['oil_avg'] / cycle_base - 1) if cycle_base else 0):+.1%},"
  923. f"{level_reason}",
  924. "_dt": row["_dt"],
  925. })
  926. index = end + 1
  927. return sorted(output, key=lambda row: (row["批次"], row["_dt"]))
  928. def _build_oil_hour_stage_segment(series: list[dict], level: str):
  929. first = series[0]
  930. last = series[-1]
  931. first_reason = series[0]["小时判定依据"]
  932. return {
  933. "机组": first["机组"],
  934. "批次": first["批次"],
  935. "周期编号": first["周期编号"],
  936. "压力阶段": level,
  937. "开始小时": first["小时"],
  938. "结束小时": last["小时"],
  939. "持续自然小时数": int((last["_dt"] - first["_dt"]).total_seconds() / 3600) + 1,
  940. "有效运行小时数": len(series),
  941. "周期基线油压": first["周期基线油压"],
  942. "正常波动带": first.get("正常波动带", ""),
  943. "正常波动带": first.get("正常波动带", ""),
  944. "阶段开始小时油压": first["小时油压均值"],
  945. "阶段结束小时油压": last["小时油压均值"],
  946. "阶段最低小时油压": min(row["小时油压均值"] for row in series),
  947. "阶段开始相对周期基线变化": first["相对周期基线变化"],
  948. "阶段结束相对周期基线变化": last["相对周期基线变化"],
  949. "阶段最大连续下降有效小时数": max(row["连续下降有效小时数"] for row in series),
  950. "阶段低报警小时数": sum(row["是否低报警小时"] == "是" for row in series),
  951. "阶段低低报警小时数": sum(row["是否低低报警小时"] == "是" for row in series),
  952. "阶段触发依据": (
  953. f"油压{first['小时油压均值']:.3f}~{last['小时油压均值']:.3f}MPa,"
  954. f"最低{min(row['小时油压均值'] for row in series):.3f}MPa,"
  955. f"连续{len(series)}个有效运行小时;{first_reason}"
  956. ),
  957. "是否形成确认等级": "是" if level != "轻微" or any(row.get("是否形成确认等级") == "是" for row in series) else "否",
  958. }
  959. def summarize_hourly_oil_cycles(hourly_levels: list[dict], hourly_segments: list[dict]):
  960. by_cycle = defaultdict(list)
  961. for row in hourly_levels:
  962. if row["压力等级"] != "数据不足":
  963. by_cycle[row["周期编号"]].append(row)
  964. segments_by_cycle = defaultdict(list)
  965. for segment in hourly_segments:
  966. segments_by_cycle[segment["周期编号"]].append(segment)
  967. result = []
  968. for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_dt"]):
  969. series.sort(key=lambda row: row["_dt"])
  970. segments = segments_by_cycle[cycle_id]
  971. levels = ["正常波动", "轻微", "异常", "严重异常"]
  972. fields = {
  973. "机组": series[0]["机组"], "批次": series[0]["批次"], "周期编号": cycle_id,
  974. "周期开始小时": series[0]["小时"], "周期结束小时": series[-1]["小时"],
  975. "周期有效运行小时数": len(series), "周期基线油压": series[0]["周期基线油压"],
  976. "正常波动带": series[0].get("正常波动带", ""),
  977. "正常波动带": series[0].get("正常波动带", ""),
  978. "阶段演化顺序": "→".join(segment["压力阶段"] for segment in segments),
  979. "周期最低小时油压": min(row["小时油压均值"] for row in series),
  980. "是否达到轻微": "是" if any(row["压力等级"] == "轻微" for row in series) else "否",
  981. "是否达到异常": "是" if any(row["压力等级"] == "异常" for row in series) else "否",
  982. "是否达到严重异常": "是" if any(row["压力等级"] == "严重异常" for row in series) else "否",
  983. "周期短时低压小时数": sum(row["立即报警"] != "" and row["压力等级"] == "正常" for row in series),
  984. "正常波动段数": sum(segment["压力阶段"] == "正常波动" for segment in segments),
  985. "正常波动小时数": sum(segment["有效运行小时数"] for segment in segments if segment["压力阶段"] == "正常波动"),
  986. }
  987. for level in levels:
  988. selected = [segment for segment in segments if segment["压力阶段"] == level]
  989. fields[f"{level}阶段开始小时"] = selected[0]["开始小时"] if selected else ""
  990. fields[f"{level}阶段结束小时"] = selected[-1]["结束小时"] if selected else ""
  991. fields[f"{level}阶段运行小时数"] = sum(segment["有效运行小时数"] for segment in selected)
  992. result.append(fields)
  993. return result
  994. def summarize_oil_cycles(daily: list[dict], segments: list[dict], hourly_segments: list[dict]):
  995. by_cycle = defaultdict(list)
  996. for row in daily:
  997. by_cycle[row["周期编号"]].append(row)
  998. segments_by_cycle = defaultdict(list)
  999. for segment in segments:
  1000. segments_by_cycle[segment["周期编号"]].append(segment)
  1001. hourly_segments_by_cycle = defaultdict(list)
  1002. for segment in hourly_segments:
  1003. hourly_segments_by_cycle[segment["周期编号"]].append(segment)
  1004. result = []
  1005. batches = sorted({row["批次"] for row in daily})
  1006. for batch in batches:
  1007. batch_cycles = sorted(
  1008. ((cycle_id, series) for cycle_id, series in by_cycle.items() if series[0]["批次"] == batch),
  1009. key=lambda item: item[1][0]["_date"],
  1010. )
  1011. previous_summary = None
  1012. for cycle_id, series in batch_cycles:
  1013. series.sort(key=lambda row: row["_date"])
  1014. cycle_segments = segments_by_cycle[cycle_id]
  1015. cycle_hourly_segments = hourly_segments_by_cycle[cycle_id]
  1016. phase_segments = {
  1017. phase: [segment for segment in cycle_segments if segment["压力阶段"] == phase]
  1018. for phase in ("正常", "轻微", "异常", "严重异常")
  1019. }
  1020. start = series[0]["_date"]
  1021. end = series[-1]["_date"]
  1022. previous_gap = ""
  1023. recurrence = "否"
  1024. if previous_summary is not None:
  1025. previous_gap = (start - previous_summary["_end"]).days
  1026. if (
  1027. previous_summary["达到异常阶段"] == "是"
  1028. and any(row["压力阶段"] in ("异常", "严重异常") for row in series)
  1029. ):
  1030. recurrence = "疑似恢复后再次出现"
  1031. abnormal = any(row["压力阶段"] in ("异常", "严重异常") for row in series)
  1032. severe = any(row["压力阶段"] == "严重异常" for row in series)
  1033. order = "→".join(segment["压力阶段"] for segment in cycle_segments)
  1034. def phase_value(phase: str, field: str):
  1035. values = phase_segments[phase]
  1036. return values[0][field] if values else ""
  1037. def phase_hours(phase: str):
  1038. return sum(segment["阶段运行小时数"] for segment in phase_segments[phase])
  1039. def hourly_phase_value(phase: str, field: str):
  1040. values = [segment for segment in cycle_hourly_segments if segment["压力阶段"] == phase]
  1041. if not values:
  1042. return ""
  1043. values.sort(key=lambda segment: segment["开始小时"])
  1044. return values[0][field] if field == "开始小时" else values[-1][field]
  1045. def hourly_phase_hours(phase: str):
  1046. return sum(
  1047. segment["有效运行小时数"]
  1048. for segment in cycle_hourly_segments
  1049. if segment["压力阶段"] == phase
  1050. )
  1051. summary = {
  1052. "机组": series[0]["机组"],
  1053. "批次": series[0]["批次"],
  1054. "周期编号": cycle_id,
  1055. "周期开始日期": series[0]["日期"],
  1056. "周期结束日期": series[-1]["日期"],
  1057. "周期有效运行天数": len(series),
  1058. "周期有效运行小时数": sum(row["有效运行小时数"] for row in series),
  1059. "周期基线油压": series[0]["周期基线油压"],
  1060. "正常阶段开始": phase_value("正常", "开始日期"),
  1061. "正常阶段结束": phase_value("正常", "结束日期"),
  1062. "正常阶段运行小时数": phase_hours("正常"),
  1063. "轻微阶段开始": phase_value("轻微", "开始日期"),
  1064. "轻微阶段结束": phase_value("轻微", "结束日期"),
  1065. "轻微阶段运行小时数": phase_hours("轻微"),
  1066. "异常阶段开始": phase_value("异常", "开始日期"),
  1067. "异常阶段结束": phase_value("异常", "结束日期"),
  1068. "异常阶段运行小时数": phase_hours("异常"),
  1069. "严重异常阶段开始": phase_value("严重异常", "开始日期"),
  1070. "严重异常阶段结束": phase_value("严重异常", "结束日期"),
  1071. "严重异常阶段运行小时数": phase_hours("严重异常"),
  1072. "周期最低日油压": min(row["日油压中位数"] for row in series),
  1073. "周期最低压力阶段": max(
  1074. series,
  1075. key=lambda row: {"正常": 0, "轻微": 1, "异常": 2, "严重异常": 3}
  1076. .get(row["压力阶段"], 0),
  1077. )["压力阶段"],
  1078. "达到异常阶段": "是" if abnormal else "否",
  1079. "达到严重异常阶段": "是" if severe else "否",
  1080. "阶段演化顺序": order,
  1081. "周期变化形态": "渐进下降型" if any(row["变化形态"] == "渐进下降型" for row in series) else "稳定/短时波动型",
  1082. "周期起点依据": series[0]["周期起点依据"],
  1083. "与上一周期结束间隔天数": previous_gap,
  1084. "是否疑似恢复后复发": recurrence,
  1085. "周期短时低压小时数": sum(row["短时低压小时数"] for row in series),
  1086. "小时分级轻微开始": hourly_phase_value("轻微", "开始小时"),
  1087. "小时分级轻微结束": hourly_phase_value("轻微", "结束小时"),
  1088. "小时分级轻微运行小时数": hourly_phase_hours("轻微"),
  1089. "小时分级异常开始": hourly_phase_value("异常", "开始小时"),
  1090. "小时分级异常结束": hourly_phase_value("异常", "结束小时"),
  1091. "小时分级异常运行小时数": hourly_phase_hours("异常"),
  1092. "小时分级严重异常开始": hourly_phase_value("严重异常", "开始小时"),
  1093. "小时分级严重异常结束": hourly_phase_value("严重异常", "结束小时"),
  1094. "小时分级严重异常运行小时数": hourly_phase_hours("严重异常"),
  1095. "_end": end,
  1096. }
  1097. result.append(summary)
  1098. previous_summary = summary
  1099. for row in result:
  1100. row.pop("_end", None)
  1101. return result
  1102. def vibration_hour_signals(rows: list[dict], configs: dict[int, dict]):
  1103. signals = []
  1104. grouped = defaultdict(list)
  1105. for row in rows:
  1106. if row["运行覆盖率"] >= MIN_HOUR_COVERAGE:
  1107. grouped[row["批次"]].append(row)
  1108. for batch, series in grouped.items():
  1109. series.sort(key=lambda row: row["_dt"])
  1110. config = configs[batch]
  1111. valid = []
  1112. last_valid_dt = None
  1113. for row in series:
  1114. if row["运行覆盖率"] < MIN_HOUR_COVERAGE:
  1115. continue
  1116. restart = last_valid_dt is None or row["_dt"] - last_valid_dt > timedelta(hours=VIBRATION_RESTART_GAP_HOURS)
  1117. if restart:
  1118. # The first complete running hour after a restart is a
  1119. # transition hour. Keep older valid hours for the long-term
  1120. # robust reference, but do not evaluate this transition hour.
  1121. last_valid_dt = row["_dt"]
  1122. continue
  1123. prior = valid[-336:]
  1124. side_signals = []
  1125. for side in ("coupling", "chain"):
  1126. avg_center, avg_scale = robust_baseline([item[f"{side}_avg"] for item in prior])
  1127. max_center, max_scale = robust_baseline([item[f"{side}_max"] for item in prior])
  1128. current_avg = row[f"{side}_avg"]
  1129. current_max = row[f"{side}_max"]
  1130. if current_max is None:
  1131. continue
  1132. avg_z = ((current_avg - avg_center) / avg_scale) if avg_center is not None and current_avg is not None else 0.0
  1133. peak_z = (current_max - max_center) / max_scale if max_center is not None and max_scale else 0.0
  1134. running = max(row["运行样本数"], 1)
  1135. high_ratio = (row[f"{side}_high_count"] or 0) / running
  1136. high_high_ratio = (row[f"{side}_high_high_count"] or 0) / running
  1137. limits = config[side]
  1138. local_prior = valid[:6] if len(valid) >= 6 else []
  1139. local_avg = median_value([item[f"{side}_avg"] for item in local_prior])
  1140. local_max = median_value([item[f"{side}_max"] for item in local_prior])
  1141. local_rise = False
  1142. severe = (
  1143. high_high_ratio >= 0.01
  1144. or current_max >= limits["high_high"]
  1145. or peak_z >= 8
  1146. )
  1147. abnormal = (
  1148. high_ratio >= 0.01
  1149. or current_max >= limits["high"]
  1150. or peak_z >= 5
  1151. or avg_z >= 4
  1152. or local_rise
  1153. )
  1154. mild = peak_z >= 4 or avg_z >= 3 or local_rise
  1155. if severe or abnormal or mild:
  1156. side_signals.append({
  1157. "side": side,
  1158. "avg_z": avg_z,
  1159. "peak_z": peak_z,
  1160. "high_ratio": high_ratio,
  1161. "high_high_ratio": high_high_ratio,
  1162. "severe": severe,
  1163. "abnormal": abnormal,
  1164. "value": current_max,
  1165. "avg": current_avg,
  1166. "prior_max_center": max_center,
  1167. "prior_max_scale": max_scale,
  1168. "prior_avg_center": avg_center,
  1169. })
  1170. if side_signals:
  1171. strongest = max(side_signals, key=lambda item: (item["severe"], item["peak_z"], item["avg_z"]))
  1172. signal_level = 3 if any(item["severe"] for item in side_signals) else 2 if any(item["abnormal"] for item in side_signals) else 1
  1173. reasons = []
  1174. for item in side_signals:
  1175. label = "联轴器端" if item["side"] == "coupling" else "链轮端"
  1176. reasons.append(
  1177. f"{label}均值偏离{item['avg_z']:.1f}倍、峰值偏离{item['peak_z']:.1f}倍、"
  1178. f"高报警比例{item['high_ratio']:.1%}"
  1179. )
  1180. signals.append({
  1181. "批次": batch,
  1182. "机组": config["unit"],
  1183. "_dt": row["_dt"],
  1184. "_level": signal_level,
  1185. "_severe": signal_level == 3,
  1186. "_side_signals": side_signals,
  1187. "_strongest": strongest,
  1188. "_reason": ";".join(reasons),
  1189. })
  1190. valid.append(row)
  1191. last_valid_dt = row["_dt"]
  1192. return signals
  1193. def merge_vibration_signals(signals: list[dict], configs: dict[int, dict]):
  1194. grouped = defaultdict(list)
  1195. for signal in signals:
  1196. grouped[signal["批次"]].append(signal)
  1197. events = []
  1198. for batch, series in grouped.items():
  1199. series.sort(key=lambda item: item["_dt"])
  1200. current = []
  1201. for signal in series:
  1202. if not current or signal["_dt"] - current[-1]["_dt"] <= timedelta(hours=VIBRATION_MERGE_GAP_HOURS):
  1203. current.append(signal)
  1204. else:
  1205. events.append(build_vibration_event(current, configs[batch]))
  1206. current = [signal]
  1207. if current:
  1208. events.append(build_vibration_event(current, configs[batch]))
  1209. for index, event in enumerate(events, 1):
  1210. event["周期编号"] = f'{event["机组"]}-V{index:02d}'
  1211. filtered = []
  1212. for event in events:
  1213. # A single post-start relative-rise signal is not sufficient evidence;
  1214. # absolute severe signals remain eligible on their own.
  1215. if event["_严重信号小时数"] > 0 or event["_触发小时数"] >= 2:
  1216. filtered.append(event)
  1217. return sorted(filtered, key=lambda row: (row["批次"], row["开始小时"]))
  1218. def vibration_score(duration_hours: int, valid_hours: int, avg_change: float, max_value: float, peak_z: float, high_hours: int, high_high_hours: int, config: dict, side: str):
  1219. duration_score = min(40.0, max(0.0, (duration_hours - 1) / 12 * 40))
  1220. avg_score = min(20.0, max(0.0, (avg_change - 0.10) / 0.40 * 20)) if avg_change >= 0.10 else 0.0
  1221. high = config[side]["high"]
  1222. high_high = config[side]["high_high"]
  1223. if max_value < high:
  1224. max_score = min(10.0, max(0.0, (max_value / max(high, 1e-6) - 0.8) / 0.2 * 10))
  1225. elif max_value < high_high:
  1226. max_score = 10.0 + (max_value - high) / max(high_high - high, 1e-6) * 10.0
  1227. else:
  1228. max_score = 20.0
  1229. peak_score = min(20.0, max(0.0, (peak_z - 4.0) / 8.0 * 20.0)) if peak_z >= 4.0 else 0.0
  1230. total = duration_score + avg_score + max_score + peak_score
  1231. hard_severe = max_value >= high_high or high_high_hours > 0 or (peak_z >= 8.0 and valid_hours >= 2)
  1232. if hard_severe or total >= 60.0:
  1233. level = "严重异常"
  1234. elif total >= 30.0:
  1235. level = "异常"
  1236. else:
  1237. level = "轻微"
  1238. return total, level, duration_score, avg_score, max_score, peak_score
  1239. def build_vibration_event(signals: list[dict], config: dict):
  1240. start = signals[0]["_dt"]
  1241. end = signals[-1]["_dt"]
  1242. all_side = [item for signal in signals for item in signal["_side_signals"]]
  1243. strongest = max(all_side, key=lambda item: (item["severe"], item["peak_z"], item["value"]))
  1244. main_side = "联轴器端" if strongest["side"] == "coupling" else "链轮端"
  1245. max_value = max(item["value"] for item in all_side if item["value"] is not None)
  1246. max_peak_z = max(item["peak_z"] for item in all_side)
  1247. high_hours = sum(
  1248. 1 for signal in signals
  1249. if any(item["high_ratio"] > 0 for item in signal["_side_signals"])
  1250. )
  1251. high_high_hours = sum(
  1252. 1 for signal in signals
  1253. if any(item["high_high_ratio"] > 0 for item in signal["_side_signals"])
  1254. )
  1255. severe_hours = sum(1 for signal in signals if signal["_severe"])
  1256. duration_hours = int((end - start).total_seconds() / 3600) + 1
  1257. running_days = len({signal["_dt"].date() for signal in signals})
  1258. side_avg_changes = []
  1259. for side in ("coupling", "chain"):
  1260. values = [item["avg"] for item in all_side if item["side"] == side and item["avg"] is not None]
  1261. prior_values = [item["prior_avg_center"] for item in all_side if item["side"] == side and item["prior_avg_center"] is not None]
  1262. if values and prior_values and median(prior_values) != 0:
  1263. side_avg_changes.append((median(values) / median(prior_values) - 1, side))
  1264. avg_change, _ = max(side_avg_changes, key=lambda item: item[0], default=(0.0, strongest["side"]))
  1265. peak_ratio = max_value / max(strongest["prior_max_center"], 1e-6)
  1266. if max_value >= config[strongest["side"]]["high_high"] or severe_hours > 0 and (high_high_hours > 0 or max_peak_z >= 8):
  1267. level = "严重异常"
  1268. elif max_value >= config[strongest["side"]]["high"] and high_hours >= 2:
  1269. level = "严重异常"
  1270. elif max_value >= config[strongest["side"]]["high"] or max_peak_z >= 5 or avg_change >= 0.10:
  1271. level = "异常"
  1272. else:
  1273. level = "轻微"
  1274. if peak_ratio >= 1.5 and avg_change < 0.25:
  1275. shape = "间歇冲击型"
  1276. elif avg_change >= 0.10 and len(signals) >= 3:
  1277. shape = "持续抬升型"
  1278. else:
  1279. shape = "混合型"
  1280. peak_signal = max(signals, key=lambda signal: signal["_strongest"]["peak_z"])
  1281. first_signal = signals[0]
  1282. last_signal = signals[-1]
  1283. first_value = first_signal["_strongest"]["value"]
  1284. last_value = last_signal["_strongest"]["value"]
  1285. baseline = strongest["prior_max_center"]
  1286. scale = max(strongest.get("prior_max_scale", 0.0), 0.0)
  1287. normal_band = f"{baseline - 2 * scale:.3f}~{baseline + 2 * scale:.3f} mm/s" if baseline is not None else ""
  1288. consecutive = 0
  1289. longest_consecutive = 0
  1290. for signal in signals:
  1291. if signal["_level"] >= 2:
  1292. consecutive += 1
  1293. longest_consecutive = max(longest_consecutive, consecutive)
  1294. else:
  1295. consecutive = 0
  1296. score, level, duration_score, avg_score, max_score, peak_score = vibration_score(
  1297. duration_hours, len(signals), avg_change, max_value, max_peak_z,
  1298. high_hours, high_high_hours, config, strongest["side"]
  1299. )
  1300. confirmed = "是" if level != "轻微" else "否"
  1301. if shape == "持续抬升型":
  1302. basis = f"{main_side}振动均值持续抬升,阶段内{longest_consecutive or len(signals)}个有效运行小时触发"
  1303. elif shape == "间歇冲击型":
  1304. basis = f"{main_side}重复出现振动冲击,阶段内{len(signals)}个有效运行小时触发"
  1305. else:
  1306. basis = f"{main_side}振动峰值或均值偏离近期基线,阶段内{len(signals)}个有效运行小时触发"
  1307. if high_high_hours:
  1308. basis += f";高高报警小时{high_high_hours}个"
  1309. elif high_hours:
  1310. basis += f";高报警小时{high_hours}个"
  1311. basis += f";最大振动值{max_value:.3f}mm/s,峰值高于近期基线{max_peak_z:.2f}个稳健尺度(峰值/基线{peak_ratio:.2f}倍);综合评分{score:.1f}分"
  1312. return {
  1313. "机组": config["unit"],
  1314. "批次": config["batch"],
  1315. "周期编号": "",
  1316. "振动阶段": level,
  1317. "开始小时": start.strftime("%Y-%m-%d %H:%M:%S"),
  1318. "结束小时": end.strftime("%Y-%m-%d %H:%M:%S"),
  1319. "持续自然小时数": duration_hours,
  1320. "有效运行小时数": len(signals),
  1321. "近期基线振动": baseline,
  1322. "正常波动带": normal_band,
  1323. "阶段开始小时振动": first_value,
  1324. "阶段结束小时振动": last_value,
  1325. "阶段最大振动值": max_value,
  1326. "阶段开始相对近期基线变化": (first_value / baseline - 1) if baseline else None,
  1327. "阶段结束相对近期基线变化": (last_value / baseline - 1) if baseline else None,
  1328. "阶段最大连续异常有效小时数": longest_consecutive,
  1329. "阶段综合评分": score,
  1330. "持续时间得分": duration_score,
  1331. "平均振动得分": avg_score,
  1332. "最大振动值得分": max_score,
  1333. "峰值偏离得分": peak_score,
  1334. "最强异常侧": main_side,
  1335. "最强侧近期峰值基线": baseline,
  1336. "最强侧近期稳健尺度": scale,
  1337. "阶段最大峰值/基线比例": peak_ratio,
  1338. "阶段高报警小时数": high_hours,
  1339. "阶段高高报警小时数": high_high_hours,
  1340. "_严重信号小时数": severe_hours,
  1341. "_触发小时数": len(signals),
  1342. "是否形成确认等级": confirmed,
  1343. "阶段触发依据": basis,
  1344. }
  1345. def write_csv(path: Path, rows: list[dict], fields: list[str], headers: dict[str, str]):
  1346. path.parent.mkdir(parents=True, exist_ok=True)
  1347. with path.open("w", encoding="utf-8-sig", newline="") as handle:
  1348. writer = csv.DictWriter(handle, fieldnames=fields, extrasaction="ignore")
  1349. writer.writerow({field: headers.get(field, field) for field in fields})
  1350. writer.writerows(rows)
  1351. def format_number(value):
  1352. if value is None or value == "":
  1353. return ""
  1354. if isinstance(value, float):
  1355. return round(value, 6)
  1356. return value
  1357. def clean_rows(rows):
  1358. cleaned = []
  1359. for row in rows:
  1360. item = {}
  1361. for key, value in row.items():
  1362. if key.startswith("_"):
  1363. continue
  1364. item[key] = format_number(value)
  1365. cleaned.append(item)
  1366. return cleaned
  1367. def write_thresholds(path: Path, configs: dict[int, dict]):
  1368. rows = []
  1369. for batch in sorted(configs):
  1370. config = configs[batch]
  1371. rows.extend([
  1372. {
  1373. "机组": config["unit"], "批次": batch, "点位": "YSJ_5",
  1374. "点位描述": config["metadata"]["oil_description"], "单位": config["metadata"]["oil_units"],
  1375. "低报警阈值": config["oil"]["low"], "低低报警阈值": config["oil"]["low_low"],
  1376. "高报警阈值": config["oil"]["high"], "高高报警阈值": config["oil"]["high_high"],
  1377. "低报警配置": config["alarm_sources"]["oil"]["low"],
  1378. "低低报警配置": config["alarm_sources"]["oil"]["low_low"],
  1379. "高报警配置": config["alarm_sources"]["oil"]["high"],
  1380. "高高报警配置": config["alarm_sources"]["oil"]["high_high"],
  1381. },
  1382. {
  1383. "机组": config["unit"], "批次": batch, "点位": "YSJ_10",
  1384. "点位描述": config["metadata"]["coupling_description"], "单位": config["metadata"]["coupling_units"],
  1385. "低报警阈值": "", "低低报警阈值": "",
  1386. "高报警阈值": config["coupling"]["high"], "高高报警阈值": config["coupling"]["high_high"],
  1387. "低报警配置": "", "低低报警配置": "",
  1388. "高报警配置": config["alarm_sources"]["coupling"]["high"],
  1389. "高高报警配置": config["alarm_sources"]["coupling"]["high_high"],
  1390. },
  1391. {
  1392. "机组": config["unit"], "批次": batch, "点位": "YSJ_11",
  1393. "点位描述": config["metadata"]["chain_description"], "单位": config["metadata"]["chain_units"],
  1394. "低报警阈值": "", "低低报警阈值": "",
  1395. "高报警阈值": config["chain"]["high"], "高高报警阈值": config["chain"]["high_high"],
  1396. "低报警配置": "", "低低报警配置": "",
  1397. "高报警配置": config["alarm_sources"]["chain"]["high"],
  1398. "高高报警配置": config["alarm_sources"]["chain"]["high_high"],
  1399. },
  1400. ])
  1401. fields = [
  1402. "机组", "批次", "点位", "点位描述", "单位", "低报警阈值", "低低报警阈值",
  1403. "高报警阈值", "高高报警阈值", "低报警配置", "低低报警配置", "高报警配置", "高高报警配置",
  1404. ]
  1405. write_csv(path, rows, fields, {field: field for field in fields})
  1406. def main():
  1407. parser = argparse.ArgumentParser(description="只读生成 PKS 润滑油压力长期演化和振动短期异常段")
  1408. parser.add_argument("--batch", nargs="+", type=int, choices=sorted(BATCH_TO_UNIT), default=list(BATCH_TO_UNIT))
  1409. parser.add_argument("--start", default="2025-04-01 00:00:00")
  1410. parser.add_argument("--end", default="2026-05-01 00:00:00")
  1411. parser.add_argument("--chunk-days", type=int, default=14)
  1412. parser.add_argument("--output-dir", default=str(Path(__file__).resolve().parents[1] / "cache" / "pks_scanner" / "evolution"))
  1413. args = parser.parse_args()
  1414. start, end = parse_dt(args.start), parse_dt(args.end)
  1415. if end <= start or args.chunk_days < 1:
  1416. parser.error("时间范围或 chunk-days 无效")
  1417. print("只读模式:不修改数据库;只分析 YSJ_5、YSJ_10、YSJ_11、YSJ_41")
  1418. configs = load_alarm_config()
  1419. print("已按 AlarmType 读取每台机器的 site_point 阈值")
  1420. all_rows = []
  1421. for batch in args.batch:
  1422. all_rows.extend(fetch_hourly(batch, start, end, args.chunk_days, configs[batch]))
  1423. all_rows.sort(key=lambda row: (row["批次"], row["_dt"]))
  1424. oil_hourly_levels = classify_oil_hourly_segments(all_rows, configs)
  1425. oil_hour_stage_segments = build_oil_hour_stage_segments(oil_hourly_levels)
  1426. oil_cycles = summarize_hourly_oil_cycles(oil_hourly_levels, oil_hour_stage_segments)
  1427. signals = vibration_hour_signals(all_rows, configs)
  1428. vibration_events = merge_vibration_signals(signals, configs)
  1429. out = Path(args.output_dir)
  1430. write_thresholds(out / "报警阈值配置.csv", configs)
  1431. hourly_fields = [
  1432. "批次", "机组", "小时", "样本数", "运行样本数", "运行覆盖率",
  1433. "oil_avg", "oil_min", "oil_max", "coupling_avg", "coupling_max",
  1434. "chain_avg", "chain_max", "oil_low_count", "oil_low_low_count",
  1435. "coupling_high_count", "coupling_high_high_count", "chain_high_count",
  1436. "chain_high_high_count",
  1437. ]
  1438. hourly_headers = {
  1439. "oil_avg": "润滑油压力均值", "oil_min": "润滑油压力最小值", "oil_max": "润滑油压力最大值",
  1440. "coupling_avg": "联轴器端振动均值", "coupling_max": "联轴器端振动最大值",
  1441. "chain_avg": "链轮端振动均值", "chain_max": "链轮端振动最大值",
  1442. "oil_low_count": "润滑油低报警样本数", "oil_low_low_count": "润滑油低低报警样本数",
  1443. "coupling_high_count": "联轴器端高报警样本数", "coupling_high_high_count": "联轴器端高高报警样本数",
  1444. "chain_high_count": "链轮端高报警样本数", "chain_high_high_count": "链轮端高高报警样本数",
  1445. "批次": "批次", "机组": "机组", "小时": "小时", "样本数": "样本数",
  1446. "运行样本数": "运行样本数", "运行覆盖率": "运行覆盖率",
  1447. }
  1448. hourly_level_fields = [
  1449. "机组", "批次", "周期编号", "小时", "运行样本数", "运行覆盖率", "小时油压均值",
  1450. "小时油压最小值", "小时油压最大值", "周期基线油压", "近期24有效小时基线油压",
  1451. "相对周期基线变化", "相对近期基线变化", "近6有效小时油压斜率", "连续下降有效小时数",
  1452. "历史有效运行小时数", "低报警样本比例", "低低报警样本比例", "低报警阈值", "低低报警阈值",
  1453. "是否低报警小时", "是否低低报警小时", "压力等级", "原始小时等级", "立即报警",
  1454. "正常波动标记", "低低报警连续小时数", "是否纳入等级判定", "小时判定依据",
  1455. ]
  1456. write_csv(out / "润滑油压力小时等级.csv", clean_rows(oil_hourly_levels), hourly_level_fields, {field: field for field in hourly_level_fields})
  1457. hour_stage_fields = [
  1458. "机组", "批次", "周期编号", "压力阶段", "开始小时", "结束小时", "持续自然小时数",
  1459. "有效运行小时数", "周期基线油压", "正常波动带", "阶段开始小时油压", "阶段结束小时油压", "阶段最低小时油压",
  1460. "阶段开始相对周期基线变化", "阶段结束相对周期基线变化", "阶段最大连续下降有效小时数",
  1461. "阶段低报警小时数", "阶段低低报警小时数", "是否形成确认等级", "阶段触发依据",
  1462. ]
  1463. write_csv(out / "润滑油压力小时阶段段.csv", clean_rows(oil_hour_stage_segments), hour_stage_fields, {field: field for field in hour_stage_fields})
  1464. cycle_fields = [
  1465. "机组", "批次", "周期编号", "周期开始小时", "周期结束小时", "周期有效运行小时数",
  1466. "周期基线油压", "正常波动带", "阶段演化顺序", "周期最低小时油压", "是否达到轻微", "是否达到异常",
  1467. "是否达到严重异常", "周期短时低压小时数", "正常波动段数", "正常波动小时数", "轻微阶段开始小时", "轻微阶段结束小时",
  1468. "轻微阶段运行小时数", "异常阶段开始小时", "异常阶段结束小时", "异常阶段运行小时数",
  1469. "严重异常阶段开始小时", "严重异常阶段结束小时", "严重异常阶段运行小时数",
  1470. ]
  1471. write_csv(out / "润滑油压力演化周期.csv", clean_rows(oil_cycles), cycle_fields, {field: field for field in cycle_fields})
  1472. vibration_fields = [
  1473. "机组", "批次", "周期编号", "振动阶段", "开始小时", "结束小时", "持续自然小时数", "有效运行小时数",
  1474. "近期基线振动", "正常波动带", "阶段开始小时振动", "阶段结束小时振动", "阶段最大振动值",
  1475. "阶段开始相对近期基线变化", "阶段结束相对近期基线变化", "阶段最大连续异常有效小时数",
  1476. "阶段高报警小时数", "阶段高高报警小时数", "是否形成确认等级", "阶段触发依据",
  1477. "阶段综合评分", "持续时间得分", "平均振动得分", "最大振动值得分", "峰值偏离得分",
  1478. "最强异常侧", "最强侧近期峰值基线", "最强侧近期稳健尺度", "阶段最大峰值/基线比例",
  1479. ]
  1480. write_csv(out / "振动异常段.csv", clean_rows(vibration_events), vibration_fields, {field: field for field in vibration_fields})
  1481. print(f"小时特征:{len(all_rows):,} 行")
  1482. print(f"润滑油压力小时等级:{len(oil_hourly_levels):,} 行")
  1483. print(f"润滑油压力小时阶段段:{len(oil_hour_stage_segments):,} 段")
  1484. print(f"润滑油压力周期:{len(oil_cycles):,} 个")
  1485. print(f"振动异常段:{len(vibration_events):,} 段")
  1486. print(f"输出目录:{out}")
  1487. if __name__ == "__main__":
  1488. main()