scan_pks_fault_evolution.py 75 KB

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