#!/usr/bin/env python3 """Build long-term oil-pressure and short-term vibration evolution records. This is a new, read-only scanner. It deliberately ignores inlet and exhaust pressure. The scanner reads the alarm configuration from site_point for each unit by AlarmType rather than assuming that AlarmLimit1 has the same meaning on every machine. """ from __future__ import annotations import argparse import csv import math import os import re import sys from collections import defaultdict from datetime import datetime, timedelta from pathlib import Path from statistics import median import pymysql sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from app.config import settings # noqa: E402 EXPECTED_PER_HOUR = 720 MIN_HOUR_COVERAGE = 0.80 MIN_DAILY_HOURS = 4 MIN_HISTORY_HOURS = 24 BASELINE_DAYS = 3 VIBRATION_BASELINE_DAYS = 7 RESET_GAP_HOURS = 12 RESET_RISE_ABS = 0.04 RESET_RISE_REL = 0.10 HOUR_ALARM_SAMPLE_RATIO = 0.20 HOURLY_BASELINE_HOURS = 24 HOURLY_TREND_HOURS = 6 OIL_CONFIRM_HOURS = 3 OIL_DOWNGRADE_HOURS = 2 OIL_RECOVER_HOURS = 3 VIBRATION_CONFIRM_HOURS = 3 VIBRATION_MERGE_GAP_HOURS = 6 VIBRATION_RESTART_GAP_HOURS = 12 SITE_POINT_COLUMNS = ( "ItemName, ItemDescription, Units, " "AlarmType1, AlarmType2, AlarmType3, AlarmType4, " "AlarmLimit1, AlarmLimit2, AlarmLimit3, AlarmLimit4" ) HOURLY_SQL = """ SELECT import_batch_id, FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start, COUNT(*) AS samples, SUM(CASE WHEN YSJ_41 > 0 THEN 1 ELSE 0 END) AS running_samples, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_avg, MIN(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_min, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_5 END) AS oil_max, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_avg, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_10 END) AS coupling_max, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_avg, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_11 END) AS chain_max, SUM(CASE WHEN YSJ_41 > 0 AND YSJ_5 <= %s THEN 1 ELSE 0 END) AS oil_low_count, SUM(CASE WHEN YSJ_41 > 0 AND YSJ_5 <= %s THEN 1 ELSE 0 END) AS oil_low_low_count, SUM(CASE WHEN YSJ_41 > 0 AND YSJ_10 >= %s THEN 1 ELSE 0 END) AS coupling_high_count, SUM(CASE WHEN YSJ_41 > 0 AND YSJ_10 >= %s THEN 1 ELSE 0 END) AS coupling_high_high_count, SUM(CASE WHEN YSJ_41 > 0 AND YSJ_11 >= %s THEN 1 ELSE 0 END) AS chain_high_count, SUM(CASE WHEN YSJ_41 > 0 AND YSJ_11 >= %s THEN 1 ELSE 0 END) AS chain_high_high_count FROM pks_long_sample WHERE import_batch_id = %s AND sample_time >= %s AND sample_time < %s GROUP BY import_batch_id, FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) ORDER BY hour_start, import_batch_id """ NUMERIC_HOURLY_FIELDS = ( "oil_avg", "oil_min", "oil_max", "coupling_avg", "coupling_max", "chain_avg", "chain_max", "oil_low_count", "oil_low_low_count", "coupling_high_count", "coupling_high_high_count", "chain_high_count", "chain_high_high_count", ) def connect(): return pymysql.connect( host=settings.db_host, port=settings.db_port, user=settings.db_user, password=settings.db_password, database=settings.db_name, charset="utf8mb4", cursorclass=pymysql.cursors.DictCursor, connect_timeout=settings.db_connect_timeout, read_timeout=600, write_timeout=60, autocommit=True, ) def parse_dt(value: str) -> datetime: return datetime.strptime(value, "%Y-%m-%d %H:%M:%S") def chunks(start: datetime, end: datetime, days: int): cursor = start while cursor < end: nxt = min(cursor + timedelta(days=days), end) yield cursor, nxt cursor = nxt def number(value): if value is None: return None try: value = float(value) except (TypeError, ValueError): return None return value if math.isfinite(value) else None def load_pks_batches() -> dict[int, str]: """Resolve PKS import batches to units from the import metadata table.""" connection = connect() try: with connection.cursor() as cursor: cursor.execute( """ SELECT id, directory_path FROM compressor_data_import WHERE data_type = %s AND directory_path IS NOT NULL ORDER BY id """, ("PKS",), ) result = {} for row in cursor.fetchall(): match = re.search(r"(? tuple[float | None, str]: for index in range(1, 5): configured_type = item.get(f"AlarmType{index}") if configured_type is not None and str(configured_type).strip() == alarm_type: return number(item.get(f"AlarmLimit{index}")), f"AlarmType{index}/AlarmLimit{index}" return None, "" def load_alarm_config(batch_to_unit: dict[int, str]) -> dict[int, dict]: """Read alarm limits by AlarmType, because AlarmLimit positions differ.""" result = {} connection = connect() try: with connection.cursor() as cursor: unit_numbers = sorted({unit[0] for unit in batch_to_unit.values()}) names = [f"YSJ{unit}_{point}" for unit in unit_numbers for point in (5, 10, 11)] placeholders = ",".join(["%s"] * len(names)) cursor.execute( f"SELECT {SITE_POINT_COLUMNS} FROM site_point WHERE ItemName IN ({placeholders})", names, ) items = {row["ItemName"]: row for row in cursor.fetchall()} finally: connection.close() unit_configs = {} for unit in sorted(set(batch_to_unit.values())): unit_number = unit[0] oil = items.get(f"YSJ{unit_number}_5") coupling = items.get(f"YSJ{unit_number}_10") chain = items.get(f"YSJ{unit_number}_11") if not oil or not coupling or not chain: raise RuntimeError(f"site_point 缺少 {unit} 的 YSJ_5/10/11 配置") oil_entries = { "low": alarm_entry(oil, "PVLow"), "low_low": alarm_entry(oil, "PVLowLow"), "high": alarm_entry(oil, "PVHigh"), "high_high": alarm_entry(oil, "PVHighHigh"), } coupling_entries = { "high": alarm_entry(coupling, "PVHigh"), "high_high": alarm_entry(coupling, "PVHighHigh"), } chain_entries = { "high": alarm_entry(chain, "PVHigh"), "high_high": alarm_entry(chain, "PVHighHigh"), } oil_limits = {key: value for key, (value, _source) in oil_entries.items()} coupling_limits = {key: value for key, (value, _source) in coupling_entries.items()} chain_limits = {key: value for key, (value, _source) in chain_entries.items()} if any(value is None for value in oil_limits.values()): raise RuntimeError(f"{unit} YSJ_5 缺少完整低压报警阈值") if any(value is None for value in (*coupling_limits.values(), *chain_limits.values())): raise RuntimeError(f"{unit} YSJ_10/11 缺少完整振动报警阈值") unit_configs[unit] = { "unit": unit, "oil": oil_limits, "coupling": coupling_limits, "chain": chain_limits, "alarm_sources": { "oil": {key: source for key, (_value, source) in oil_entries.items()}, "coupling": {key: source for key, (_value, source) in coupling_entries.items()}, "chain": {key: source for key, (_value, source) in chain_entries.items()}, }, "metadata": { "oil_description": oil.get("ItemDescription") or "", "oil_units": oil.get("Units") or "", "coupling_description": coupling.get("ItemDescription") or "", "coupling_units": coupling.get("Units") or "", "chain_description": chain.get("ItemDescription") or "", "chain_units": chain.get("Units") or "", }, } for batch, unit in batch_to_unit.items(): result[batch] = {**unit_configs[unit], "batch": batch} return result def load_batch_ranges(batches: list[int], start: datetime, end: datetime) -> dict[int, tuple[datetime, datetime]]: placeholders = ",".join(["%s"] * len(batches)) connection = connect() try: with connection.cursor() as cursor: cursor.execute( f""" SELECT import_batch_id, MIN(sample_time) AS data_start, MAX(sample_time) AS data_end FROM pks_long_sample WHERE import_batch_id IN ({placeholders}) AND sample_time >= %s AND sample_time < %s GROUP BY import_batch_id """, (*batches, start, end), ) return { int(row["import_batch_id"]): ( max(row["data_start"], start), min(row["data_end"] + timedelta(seconds=5), end), ) for row in cursor.fetchall() } finally: connection.close() def query_params(batches: list[int], start: datetime, end: datetime, config: dict) -> tuple: return ( config["oil"]["low"], config["oil"]["low_low"], config["coupling"]["high"], config["coupling"]["high_high"], config["chain"]["high"], config["chain"]["high_high"], *batches, start, end, ) def fetch_hourly(unit: str, batches: list[int], start: datetime, end: datetime, chunk_days: int, config: dict): rows = [] placeholders = ",".join(["%s"] * len(batches)) sql = HOURLY_SQL.replace("import_batch_id = %s", f"import_batch_id IN ({placeholders})") connection = connect() try: with connection.cursor() as cursor: for index, (chunk_start, chunk_end) in enumerate(chunks(start, end, chunk_days), 1): params = query_params(batches, chunk_start, chunk_end, config) cursor.execute("EXPLAIN " + sql, params) plan = cursor.fetchone() if not plan or plan.get("key") not in ("PRIMARY",): raise RuntimeError( f"安全检查失败:unit={unit} batches={batches} chunk={chunk_start} " f"EXPLAIN 未使用 PRIMARY: {plan}" ) cursor.execute(sql, params) chunk_count = 0 for raw in cursor.fetchall(): batch = int(raw["import_batch_id"]) row = { "批次": batch, "机组": unit, "小时": raw["hour_start"].strftime("%Y-%m-%d %H:%M:%S"), "样本数": int(raw["samples"] or 0), "运行样本数": int(raw["running_samples"] or 0), } row["运行覆盖率"] = row["运行样本数"] / EXPECTED_PER_HOUR for field in NUMERIC_HOURLY_FIELDS: row[field] = number(raw[field]) row["_dt"] = raw["hour_start"] rows.append(row) chunk_count += 1 print( f"unit={unit} batches={batches} chunk={index} {chunk_start}..{chunk_end}: " f"{chunk_count:,} hourly rows" ) finally: connection.close() return rows def median_value(values): values = [value for value in values if value is not None] return median(values) if values else None def max_consecutive(values): longest = 0 current = 0 for value in values: if value: current += 1 longest = max(longest, current) else: current = 0 return longest def robust_baseline(values, minimum: int = MIN_HISTORY_HOURS): values = [value for value in values if value is not None] if len(values) < minimum: return None, None center = median(values) mad = median(abs(value - center) for value in values) # Keep a meaningful relative scale when a sensor is very stable. scale = max(1.4826 * mad, abs(center) * 0.05, 1e-6) return center, scale def linear_slope(values): pairs = [(index, value) for index, value in enumerate(values) if value is not None] if len(pairs) < 3: return None xbar = sum(x for x, _ in pairs) / len(pairs) ybar = sum(y for _, y in pairs) / len(pairs) denominator = sum((x - xbar) ** 2 for x, _ in pairs) if denominator == 0: return None return sum((x - xbar) * (y - ybar) for x, y in pairs) / denominator def consecutive_decline_hours(history: list[dict], current: dict) -> int: """Count contiguous valid hours that have not risen from the prior hour.""" count = 1 newer = current for older in reversed(history): if newer["_dt"] - older["_dt"] != timedelta(hours=1): break if newer["oil_avg"] is None or older["oil_avg"] is None or newer["oil_avg"] > older["oil_avg"]: break count += 1 newer = older return count def consecutive_flag_hours(history: list[dict], current: dict, field: str) -> int: count = 1 if current.get(field) else 0 if not count: return 0 newer = current for older in reversed(history): if newer["_dt"] - older["_dt"] != timedelta(hours=1) or not older.get(field): break count += 1 newer = older return count def build_daily(rows: list[dict], configs: dict[int, dict]): grouped = defaultdict(list) for row in rows: if row["运行覆盖率"] >= MIN_HOUR_COVERAGE and row["oil_avg"] is not None: grouped[(row["批次"], row["_dt"].date())].append(row) daily = [] for (batch, day), values in sorted(grouped.items()): if len(values) < MIN_DAILY_HOURS: continue config = configs[batch] running_samples = sum(row["运行样本数"] for row in values) low_count = sum(row["oil_low_count"] or 0 for row in values) low_low_count = sum(row["oil_low_low_count"] or 0 for row in values) oil_values = [row["oil_avg"] for row in values] low_hours = [ (row["oil_low_count"] or 0) / max(row["运行样本数"], 1) >= HOUR_ALARM_SAMPLE_RATIO for row in values ] low_low_hours = [ (row["oil_low_low_count"] or 0) / max(row["运行样本数"], 1) >= HOUR_ALARM_SAMPLE_RATIO for row in values ] daily.append({ "批次": batch, "机组": config["unit"], "日期": day.isoformat(), "_date": day, "有效运行小时数": len(values), "日运行覆盖率": running_samples / (EXPECTED_PER_HOUR * 24), "日油压中位数": median(oil_values), "日油压最小值": min(row["oil_min"] for row in values if row["oil_min"] is not None), "日油压最大值": max(row["oil_max"] for row in values if row["oil_max"] is not None), # Raw sample ratios are retained for diagnosis. Alarm hours use # a 20% running-sample gate so isolated 5-second drops do not # become a full-day severe pressure phase. "低报警小时数": sum(low_hours), "低低报警小时数": sum(low_low_hours), "低报警最长连续小时数": max_consecutive(low_hours), "低低报警最长连续小时数": max_consecutive(low_low_hours), "低报警样本数": low_count, "低低报警样本数": low_low_count, "低报警运行样本比例": low_count / max(running_samples, 1), "低低报警运行样本比例": low_low_count / max(running_samples, 1), "低报警阈值": config["oil"]["low"], "低低报警阈值": config["oil"]["low_low"], }) return daily def assign_cycles(daily: list[dict]): by_batch = defaultdict(list) for row in daily: by_batch[row["批次"]].append(row) for batch, series in by_batch.items(): series.sort(key=lambda row: row["_date"]) cycle_number = 0 previous = None for row in series: reset_reason = "" if previous is None: cycle_number += 1 reset_reason = "扫描范围内首个有效运行日" else: gap_days = (row["_date"] - previous["_date"]).days rise = row["日油压中位数"] - previous["日油压中位数"] rise_rel = rise / max(abs(previous["日油压中位数"]), 1e-6) if gap_days > RESET_GAP_DAYS: cycle_number += 1 reset_reason = f"有效运行间隔{gap_days}天,疑似重新开机周期" elif rise >= RESET_RISE_ABS and rise_rel >= RESET_RISE_REL: cycle_number += 1 reset_reason = f"日油压回升{rise:.3f}MPa,疑似恢复后重新运行" if cycle_number == 0: cycle_number = 1 row["周期编号"] = f"{row['机组']}-P{cycle_number:02d}" row["周期起点依据"] = reset_reason previous = row def classify_oil_daily(daily: list[dict]): by_cycle = defaultdict(list) for row in daily: by_cycle[row["周期编号"]].append(row) for series in by_cycle.values(): series.sort(key=lambda row: row["_date"]) baseline_days = series[:BASELINE_DAYS] cycle_baseline = median(row["日油压中位数"] for row in baseline_days) low_alarm = series[0]["低报警阈值"] low_low_alarm = series[0]["低低报警阈值"] for index, row in enumerate(series): row["周期基线油压"] = cycle_baseline relative = row["日油压中位数"] / cycle_baseline - 1 if cycle_baseline else None recent = series[max(0, index - 4): index + 1] recent_values = [item["日油压中位数"] for item in recent] slope = linear_slope(recent_values) decline_days = 1 cursor = index while cursor > 0 and series[cursor]["日油压中位数"] <= series[cursor - 1]["日油压中位数"]: decline_days += 1 cursor -= 1 smoothed = median_value([item["日油压中位数"] for item in series[max(0, index - 2): index + 1]]) row["相对基线变化"] = relative row["近5运行日油压斜率"] = slope row["连续下降运行日数"] = decline_days row["平滑油压"] = smoothed row["是否基线期"] = "是" if index < BASELINE_DAYS else "否" downward = (slope is not None and slope < -0.002) or decline_days >= 3 recovery = index > 0 and row["日油压中位数"] > series[index - 1]["日油压中位数"] * 1.03 # Long-term phase is based on the daily level and trend. Raw # alarm samples are retained separately so isolated 5-second # drops do not promote an otherwise normal day to severe. severe_by_absolute = row["日油压中位数"] <= low_low_alarm and ( decline_days >= 2 or row["低低报警小时数"] >= 2 ) severe_by_evolution = ( relative is not None and downward and relative <= -0.20 and (decline_days >= 2 or row["日油压中位数"] <= low_alarm) ) abnormal_by_absolute = row["日油压中位数"] <= low_alarm and downward abnormal_by_evolution = relative is not None and relative <= -0.10 and downward if index < BASELINE_DAYS: phase = "正常" elif severe_by_absolute or severe_by_evolution: phase = "严重异常" elif abnormal_by_absolute or abnormal_by_evolution: phase = "异常" elif relative is not None and relative <= -0.05 and downward: phase = "轻微" else: phase = "正常" row["压力阶段"] = phase short_alarm = ( row["低报警样本数"] > 0 and row["日油压中位数"] > low_alarm ) or ( row["低低报警样本数"] > 0 and row["日油压中位数"] > low_low_alarm ) row["短时低压标记"] = "是" if short_alarm else "否" row["短时低压小时数"] = row["低报警小时数"] row["短时低低报警小时数"] = row["低低报警小时数"] reasons = [] if severe_by_absolute: reasons.append("日油压达到低低报警区间且具有持续性") if severe_by_evolution: reasons.append("相对周期基线下降至少20%且仍在下降") if abnormal_by_absolute: reasons.append("日油压进入低报警区间且仍在下降") if abnormal_by_evolution: reasons.append("相对周期基线下降至少10%且仍在下降") if phase == "轻微": reasons.append("相对周期基线下降至少5%且仍在下降") if short_alarm: reasons.append("存在短时低压样本,但日油压未进入对应长期阶段") row["阶段判定依据"] = ";".join(reasons) or "周期基线范围内" if row["短时低压标记"] == "是" and phase == "正常": shape = "短时低压波动型" elif recovery: shape = "回升型" elif downward and decline_days >= 3: shape = "渐进下降型" else: shape = "稳定型" row["变化形态"] = shape row["是否低报警"] = "是" if row["低报警小时数"] > 0 else "否" row["是否低低报警"] = "是" if row["低低报警小时数"] > 0 else "否" def build_oil_stage_segments(daily: list[dict]): """Merge daily phase labels into reviewable pressure-evolution stages.""" by_cycle = defaultdict(list) for row in daily: by_cycle[row["周期编号"]].append(row) segments = [] for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_date"]): series.sort(key=lambda row: row["_date"]) current = [] for row in series: if not current or ( row["压力阶段"] == current[-1]["压力阶段"] and (row["_date"] - current[-1]["_date"]).days <= RESET_GAP_DAYS ): current.append(row) else: segments.append(_build_oil_stage_segment(cycle_id, current, len(segments) + 1)) current = [row] if current: segments.append(_build_oil_stage_segment(cycle_id, current, len(segments) + 1)) # Stage numbers should restart within each pressure cycle. by_cycle_segments = defaultdict(list) for segment in segments: by_cycle_segments[segment["周期编号"]].append(segment) for cycle_segments in by_cycle_segments.values(): for index, segment in enumerate(cycle_segments, 1): segment["阶段序号"] = index return sorted(segments, key=lambda row: (row["批次"], row["开始日期"], row["阶段序号"])) def _build_oil_stage_segment(cycle_id: str, series: list[dict], sequence: int): first = series[0] last = series[-1] phase = first["压力阶段"] return { "机组": first["机组"], "批次": first["批次"], "周期编号": cycle_id, "阶段序号": sequence, "压力阶段": phase, "开始日期": first["日期"], "结束日期": last["日期"], "持续自然日数": (last["_date"] - first["_date"]).days + 1, "有效运行日数": len(series), "阶段运行小时数": sum(row["有效运行小时数"] for row in series), "阶段低报警小时数": sum(row["低报警小时数"] for row in series), "阶段低低报警小时数": sum(row["低低报警小时数"] for row in series), "阶段短时低压小时数": sum(row["短时低压小时数"] for row in series), "周期基线油压": first["周期基线油压"], "阶段开始日油压": first["日油压中位数"], "阶段结束日油压": last["日油压中位数"], "阶段最低日油压": min(row["日油压中位数"] for row in series), "阶段开始相对基线变化": first["相对基线变化"], "阶段结束相对基线变化": last["相对基线变化"], "阶段最大连续下降运行日数": max(row["连续下降运行日数"] for row in series), "阶段变化形态": "渐进下降型" if any(row["变化形态"] == "渐进下降型" for row in series) else "稳定/短时波动型", "阶段触发依据": "; ".join( reason for reason in (first.get("周期起点依据", ""), last.get("变化形态", "")) if reason ), } def classify_oil_hourly(rows: list[dict], daily: list[dict], configs: dict[int, dict]): """Assign a causal pressure level to each valid running hour. The cycle baseline is the pressure baseline already established for the corresponding running cycle. Recent history is limited to prior valid hours, so this output can be used to implement an online hourly updater. """ cycle_by_date = {(row["批次"], row["_date"]): row["周期编号"] for row in daily} daily_baseline = {row["周期编号"]: row["周期基线油压"] for row in daily} grouped = defaultdict(list) for row in rows: cycle_id = cycle_by_date.get((row["批次"], row["_dt"].date())) if row["运行覆盖率"] < MIN_HOUR_COVERAGE or row["oil_avg"] is None or cycle_id is None: continue item = dict(row) item["周期编号"] = cycle_id grouped[(row["批次"], cycle_id)].append(item) result = [] classified_keys = set() for (batch, cycle_id), series in grouped.items(): series.sort(key=lambda row: row["_dt"]) config = configs[batch] cycle_base = daily_baseline[cycle_id] history = [] for row in series: recent = history[-HOURLY_BASELINE_HOURS:] recent_base = median_value([item["oil_avg"] for item in recent]) trend_values = [item["oil_avg"] for item in history[-(HOURLY_TREND_HOURS - 1):]] + [row["oil_avg"]] trend = linear_slope(trend_values) decline_hours = consecutive_decline_hours(history, row) cycle_relative = row["oil_avg"] / cycle_base - 1 if cycle_base else None recent_relative = row["oil_avg"] / recent_base - 1 if recent_base else None running = max(row["运行样本数"], 1) low_ratio = (row["oil_low_count"] or 0) / running low_low_ratio = (row["oil_low_low_count"] or 0) / running low_hour = low_ratio >= HOUR_ALARM_SAMPLE_RATIO low_low_hour = low_low_ratio >= HOUR_ALARM_SAMPLE_RATIO low_low_streak = consecutive_flag_hours(history, {**row, "_low_low_hour": low_low_hour}, "_low_low_hour") downward = ( decline_hours >= 3 or (trend is not None and trend < -0.0005) or (recent_relative is not None and recent_relative <= -0.03 and trend is not None and trend < 0) ) low_alarm = config["oil"]["low"] low_low_alarm = config["oil"]["low_low"] severe = ( low_low_streak >= 3 ) or ( cycle_relative is not None and cycle_relative <= -0.20 and (downward or row["oil_avg"] <= low_alarm) ) abnormal = not severe and ( (row["oil_avg"] <= low_alarm and (low_hour or decline_hours >= 2)) or (cycle_relative is not None and cycle_relative <= -0.10 and downward) ) mild = not severe and not abnormal and ( (cycle_relative is not None and cycle_relative <= -0.05 and downward) or (recent_relative is not None and recent_relative <= -0.03 and downward) ) if severe: level = "严重异常" reason = "相对周期基线下降至少20%并持续下降" if row["oil_avg"] <= low_low_alarm: reason = "小时油压达到低低报警区间" elif abnormal: level = "异常" reason = "相对周期基线下降至少10%并持续下降" if row["oil_avg"] <= low_alarm: reason = "小时油压进入低报警区间" elif mild: level = "轻微" reason = "相对周期基线下降至少5%并持续下降" else: level = "正常" reason = "周期基线范围内" short_low = (low_hour or low_low_hour) and level == "正常" if short_low: reason += ";存在短时低压样本,未形成长期下降等级" result.append({ "机组": config["unit"], "批次": batch, "周期编号": cycle_id, "小时": row["小时"], "运行样本数": row["运行样本数"], "运行覆盖率": row["运行覆盖率"], "小时油压均值": row["oil_avg"], "小时油压最小值": row["oil_min"], "小时油压最大值": row["oil_max"], "周期基线油压": cycle_base, "近期24有效小时基线油压": recent_base, "相对周期基线变化": cycle_relative, "相对近期基线变化": recent_relative, "近6有效小时油压斜率": trend, "连续下降有效小时数": decline_hours, "历史有效运行小时数": len(history), "低报警样本比例": low_ratio, "低低报警样本比例": low_low_ratio, "低报警阈值": low_alarm, "低低报警阈值": low_low_alarm, "是否低报警小时": "是" if low_hour else "否", "是否低低报警小时": "是" if low_low_hour else "否", "压力等级": level, "短时低压标记": "是" if short_low else "否", "是否纳入等级判定": "是", "小时判定依据": reason, "_dt": row["_dt"], }) classified_keys.add((batch, row["_dt"])) history.append(row) # Keep partial-running hours visible, but do not let them affect the # long-term level or its historical baseline. for row in rows: key = (row["批次"], row["_dt"]) if row["运行样本数"] <= 0 or key in classified_keys: continue result.append({ "机组": row["机组"], "批次": row["批次"], "周期编号": cycle_by_date.get((row["批次"], row["_dt"].date()), ""), "小时": row["小时"], "运行样本数": row["运行样本数"], "运行覆盖率": row["运行覆盖率"], "小时油压均值": row["oil_avg"], "小时油压最小值": row["oil_min"], "小时油压最大值": row["oil_max"], "周期基线油压": "", "近期24有效小时基线油压": "", "相对周期基线变化": "", "相对近期基线变化": "", "近6有效小时油压斜率": "", "连续下降有效小时数": "", "历史有效运行小时数": "", "低报警样本比例": (row["oil_low_count"] or 0) / max(row["运行样本数"], 1), "低低报警样本比例": (row["oil_low_low_count"] or 0) / max(row["运行样本数"], 1), "低报警阈值": configs[row["批次"]]["oil"]["low"], "低低报警阈值": configs[row["批次"]]["oil"]["low_low"], "是否低报警小时": "", "是否低低报警小时": "", "低低报警连续小时数": "", "压力等级": "数据不足", "短时低压标记": "", "正常波动标记": "", "是否纳入等级判定": "否", "小时判定依据": "运行覆盖率不足,未参与长期小时等级判定", "_dt": row["_dt"], }) return sorted(result, key=lambda row: (row["批次"], row["_dt"])) def build_oil_hour_stage_segments(hourly_levels: list[dict]): """Merge confirmed levels and short unconfirmed triggers into segments.""" grouped = defaultdict(list) for row in hourly_levels: if row["压力等级"] in ("轻微", "异常", "严重异常"): level = row["压力等级"] elif row["压力等级"] == "轻微" and row.get("是否形成确认等级") == "否": level = "轻微" else: continue grouped[(row["批次"], row["周期编号"], level)].append(row) segments = [] for (batch, cycle_id, level), series in grouped.items(): series.sort(key=lambda row: row["_dt"]) current = [] for row in series: if not current or ( row["_dt"] - current[-1]["_dt"] <= timedelta(hours=VIBRATION_MERGE_GAP_HOURS) and row["压力等级"] == current[-1]["压力等级"] ): current.append(row) else: segments.append(_build_oil_hour_stage_segment(current, level)) current = [row] if current: segments.append(_build_oil_hour_stage_segment(current, level)) return sorted(segments, key=lambda row: (row["批次"], row["开始小时"])) def assign_hourly_cycles(rows: list[dict]): """Create oil-pressure cycles directly from valid running hours.""" grouped = defaultdict(list) for row in rows: if row["运行覆盖率"] >= MIN_HOUR_COVERAGE and row["oil_avg"] is not None: grouped[row["批次"]].append(row) for batch, series in grouped.items(): series.sort(key=lambda row: row["_dt"]) cycle_number = 0 previous = None for row in series: gap_hours = None if previous is None else (row["_dt"] - previous["_dt"]).total_seconds() / 3600 # A one-hour pressure recovery is not a restart. It can be a # sensor dropout or part of the same fault episode. Start a new # cycle only after a long gap in valid running hours. if previous is None or (gap_hours is not None and gap_hours > RESET_GAP_HOURS): cycle_number += 1 row["_cycle_id"] = f"{row['机组']}-P{cycle_number:02d}" previous = row return grouped def classify_oil_hourly_state_machine(rows: list[dict], configs: dict[int, dict]): """Classify oil pressure with upgrade/downgrade hysteresis per cycle.""" cycles = assign_hourly_cycles(rows) result = [] for batch, batch_series in cycles.items(): config = configs[batch] by_cycle = defaultdict(list) for row in batch_series: by_cycle[row["_cycle_id"]].append(row) for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_dt"]): series.sort(key=lambda item: item["_dt"]) cycle_base = median_value([row["oil_avg"] for row in series[:HOURLY_BASELINE_HOURS]]) stable_values = [row["oil_avg"] for row in series[:HOURLY_BASELINE_HOURS]] stable_center = median_value(stable_values) stable_mad = median(abs(value - stable_center) for value in stable_values) if stable_center is not None else 0 stable_scale = max(1.4826 * stable_mad, abs(stable_center or 0) * 0.01, 1e-6) history = [] state = "正常" pending = None pending_count = 0 downgrade_count = 0 for row in series: recent = history[-HOURLY_BASELINE_HOURS:] recent_base = median_value([item["oil_avg"] for item in recent]) trend_values = [item["oil_avg"] for item in history[-(HOURLY_TREND_HOURS - 1):]] + [row["oil_avg"]] trend = linear_slope(trend_values) decline_hours = consecutive_decline_hours(history, row) cycle_relative = row["oil_avg"] / cycle_base - 1 if cycle_base else None recent_relative = row["oil_avg"] / recent_base - 1 if recent_base else None running = max(row["运行样本数"], 1) low_ratio = (row["oil_low_count"] or 0) / running low_low_ratio = (row["oil_low_low_count"] or 0) / running low_hour = low_ratio >= HOUR_ALARM_SAMPLE_RATIO low_low_hour = low_low_ratio >= HOUR_ALARM_SAMPLE_RATIO low_low_streak = consecutive_flag_hours( history, {**row, "_low_low_hour": low_low_hour}, "_low_low_hour" ) downward = decline_hours >= 3 or (trend is not None and trend < -0.0005) immediate = "严重预警" if low_low_hour else "低压预警" if low_hour else "" low_alarm = config["oil"]["low"] low_low_alarm = config["oil"]["low_low"] deviation_z = ((stable_center - row["oil_avg"]) / stable_scale) if stable_center is not None else 0.0 if ( (row["oil_avg"] <= low_low_alarm and low_low_streak >= 2) or (deviation_z >= 8 and decline_hours >= OIL_CONFIRM_HOURS) ): raw = "严重异常" elif ( (row["oil_avg"] <= low_alarm and decline_hours >= OIL_CONFIRM_HOURS) or (deviation_z >= 4 and decline_hours >= OIL_CONFIRM_HOURS) ): raw = "异常" elif deviation_z >= 2 or low_hour: raw = "轻微" else: raw = "正常" rank = {"正常": 0, "轻微": 1, "异常": 2, "严重异常": 3} if rank[raw] > rank[state]: if pending == raw: pending_count += 1 else: pending, pending_count = raw, 1 required = 1 if raw == "轻微" else OIL_CONFIRM_HOURS if pending_count >= required: state = raw pending, pending_count = None, 0 downgrade_count = 0 elif rank[raw] < rank[state]: if pending == raw: pending_count += 1 else: pending, pending_count = raw, 1 required = OIL_DOWNGRADE_HOURS if raw != "正常" else OIL_RECOVER_HOURS if pending_count >= required: state = raw pending, pending_count = None, 0 downgrade_count += 1 elif raw == "正常" and state != "正常": downgrade_count += 1 pending, pending_count = None, 0 if downgrade_count >= 6: state = "正常" downgrade_count = 0 else: pending, pending_count = None, 0 downgrade_count = 0 level_label = {"正常": "正常", "轻微": "轻微", "异常": "异常", "严重异常": "严重"}[state] reason = ( f"当前油压{row['oil_avg']:.3f}MPa,近期基线{recent_base:.3f}MPa," f"偏离稳定基线波动尺度{deviation_z:.1f}倍,连续下降{decline_hours}小时" if recent_base is not None else "近期基线不足" ) if row["oil_avg"] <= low_low_alarm: reason += f",低于低低报警阈值{low_low_alarm:.3f}MPa" elif row["oil_avg"] <= low_alarm: reason += f",低于低报警阈值{low_alarm:.3f}MPa" if state != "正常": reason += f",{level_label}状态已连续确认" if low_low_streak >= 3: reason += f";低低报警连续{low_low_streak}小时" if immediate: reason += f";{immediate}" unconfirmed = state == "正常" and raw != "正常" if unconfirmed: reason += ";连续时间不足3小时,记录为轻微" output_level = "轻微" if unconfirmed else state result.append({ "机组": config["unit"], "批次": batch, "周期编号": cycle_id, "小时": row["小时"], "运行样本数": row["运行样本数"], "运行覆盖率": row["运行覆盖率"], "小时油压均值": row["oil_avg"], "小时油压最小值": row["oil_min"], "小时油压最大值": row["oil_max"], "周期基线油压": cycle_base, "近期24有效小时基线油压": recent_base, "相对周期基线变化": cycle_relative, "相对近期基线变化": recent_relative, "近6有效小时油压斜率": trend, "连续下降有效小时数": decline_hours, "历史有效运行小时数": len(history), "低报警样本比例": low_ratio, "低低报警样本比例": low_low_ratio, "低报警阈值": config["oil"]["low"], "低低报警阈值": config["oil"]["low_low"], "是否低报警小时": "是" if low_hour else "否", "是否低低报警小时": "是" if low_low_hour else "否", "低低报警连续小时数": low_low_streak, "压力等级": output_level, "原始小时等级": raw, "立即报警": immediate, "正常波动标记": "否", "是否形成确认等级": "否" if unconfirmed else "是", "是否纳入等级判定": "是", "小时判定依据": reason, "_dt": row["_dt"], }) row["_low_hour"] = low_hour row["_low_low_hour"] = low_low_hour history.append(row) return sorted(result, key=lambda row: (row["批次"], row["_dt"])) def classify_oil_hourly_segments(rows: list[dict], configs: dict[int, dict]): """Classify pressure by amplitude zone first, then confirm by duration. The configured low/low-low limits define the engineering zones. A short excursion is retained as mild; only a run of three or more consecutive valid hours in the low zone is confirmed abnormal, and the low-low zone is confirmed severe. No fixed 5/10/20 percent thresholds are used. """ cycles = assign_hourly_cycles(rows) output = [] for batch, batch_series in cycles.items(): config = configs[batch] by_cycle = defaultdict(list) for row in batch_series: by_cycle[row["_cycle_id"]].append(row) for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_dt"]): series.sort(key=lambda item: item["_dt"]) baseline_values = [row["oil_avg"] for row in series[:HOURLY_BASELINE_HOURS]] cycle_base = median_value(baseline_values) low = config["oil"]["low"] low_low = config["oil"]["low_low"] candidates = [] history = [] for row in series: current = row["oil_avg"] trend = linear_slope([item["oil_avg"] for item in history[-5:]] + [current]) decline = consecutive_decline_hours(history, row) deviation = (cycle_base - current) / cycle_base if cycle_base else 0.0 low_ratio = (row["oil_low_count"] or 0) / max(row["运行样本数"], 1) low_low_ratio = (row["oil_low_low_count"] or 0) / max(row["运行样本数"], 1) low_zone = current <= low low_low_zone = current <= low_low downward = ( decline >= 3 or (trend is not None and trend < -0.0005) ) # Relative decline is a data-derived evolution signal. The # engineering alarm limits remain higher priority, while the # 5/10/20% bands describe progression before the alarm limit. severe_candidate = low_low_zone or (deviation >= 0.20 and downward) abnormal_candidate = low_zone or (deviation >= 0.10 and downward) mild_candidate = deviation >= 0.05 and downward if severe_candidate: raw = "严重异常" elif abnormal_candidate: raw = "异常" elif mild_candidate: raw = "轻微" else: raw = "正常" candidates.append({ "row": row, "raw": raw, "trend": trend, "decline": decline, "deviation": deviation, "low_ratio": low_ratio, "low_low_ratio": low_low_ratio, "low_zone": low_zone, "low_low_zone": low_low_zone, }) history.append(row) index = 0 while index < len(candidates): raw = candidates[index]["raw"] if raw == "正常": final = "正常" end = index else: end = index while end + 1 < len(candidates) and candidates[end + 1]["raw"] == raw: end += 1 run_length = end - index + 1 # Alarm limits have highest priority for the hourly # candidate and immediate alarm fields. Long-term stage # labels still require persistence: a one/two-hour # low-low excursion is a short severe warning, not a # confirmed long-term severe stage. if run_length >= OIL_CONFIRM_HOURS: final = raw elif raw == "严重异常": final = "正常" else: final = "轻微" for item in candidates[index:end + 1]: row = item["row"] level_reason = { "轻微": "相对周期基线下降至少5%并持续下降,或普通异常持续不足3小时", "异常": "相对周期基线下降至少10%并持续下降,或低报警区间连续达到3小时", "严重异常": "相对周期基线下降至少20%并持续下降,或低低报警区间连续达到3小时", "正常": "未形成长期确认等级" if raw != "正常" else "未达到压力异常条件", }[final] if item["low_low_zone"]: level_reason += f";低低报警样本比例{item['low_low_ratio']:.1%}" elif item["low_zone"]: level_reason += f";低报警样本比例{item['low_ratio']:.1%}" output.append({ "机组": config["unit"], "批次": batch, "周期编号": cycle_id, "小时": row["小时"], "运行样本数": row["运行样本数"], "运行覆盖率": row["运行覆盖率"], "小时油压均值": row["oil_avg"], "小时油压最小值": row["oil_min"], "小时油压最大值": row["oil_max"], "周期基线油压": cycle_base, "近期24有效小时基线油压": "", "相对周期基线变化": (row["oil_avg"] / cycle_base - 1) if cycle_base else None, "相对近期基线变化": "", "近6有效小时油压斜率": item["trend"], "连续下降有效小时数": item["decline"], "历史有效运行小时数": index, "低报警样本比例": item["low_ratio"], "低低报警样本比例": item["low_low_ratio"], "低报警阈值": low, "低低报警阈值": low_low, "是否低报警小时": "是" if item["low_zone"] else "否", "是否低低报警小时": "是" if item["low_low_zone"] else "否", "低低报警连续小时数": 0, "压力等级": final, "原始小时等级": raw, "立即报警": "严重预警" if item["low_low_zone"] else "低压预警" if item["low_zone"] else "", "正常波动标记": "是" if raw != "正常" and final == "轻微" else "否", "是否形成确认等级": "是" if final in ("异常", "严重异常") else "否", "小时判定依据": f"油压{row['oil_avg']:.3f}MPa,周期基线{cycle_base:.3f}MPa," f"相对周期基线变化{((row['oil_avg'] / cycle_base - 1) if cycle_base else 0):+.1%}," f"{level_reason}", "_dt": row["_dt"], }) index = end + 1 return sorted(output, key=lambda row: (row["批次"], row["_dt"])) def _build_oil_hour_stage_segment(series: list[dict], level: str): first = series[0] last = series[-1] first_reason = series[0]["小时判定依据"] return { "机组": first["机组"], "批次": first["批次"], "周期编号": first["周期编号"], "压力阶段": level, "开始小时": first["小时"], "结束小时": last["小时"], "持续自然小时数": int((last["_dt"] - first["_dt"]).total_seconds() / 3600) + 1, "有效运行小时数": len(series), "周期基线油压": first["周期基线油压"], "正常波动带": first.get("正常波动带", ""), "正常波动带": first.get("正常波动带", ""), "阶段开始小时油压": first["小时油压均值"], "阶段结束小时油压": last["小时油压均值"], "阶段最低小时油压": min(row["小时油压均值"] for row in series), "阶段开始相对周期基线变化": first["相对周期基线变化"], "阶段结束相对周期基线变化": last["相对周期基线变化"], "阶段最大连续下降有效小时数": max(row["连续下降有效小时数"] for row in series), "阶段低报警小时数": sum(row["是否低报警小时"] == "是" for row in series), "阶段低低报警小时数": sum(row["是否低低报警小时"] == "是" for row in series), "阶段触发依据": ( f"油压{first['小时油压均值']:.3f}~{last['小时油压均值']:.3f}MPa," f"最低{min(row['小时油压均值'] for row in series):.3f}MPa," f"连续{len(series)}个有效运行小时;{first_reason}" ), "是否形成确认等级": "是" if level != "轻微" or any(row.get("是否形成确认等级") == "是" for row in series) else "否", } def summarize_hourly_oil_cycles(hourly_levels: list[dict], hourly_segments: list[dict]): by_cycle = defaultdict(list) for row in hourly_levels: if row["压力等级"] != "数据不足": by_cycle[row["周期编号"]].append(row) segments_by_cycle = defaultdict(list) for segment in hourly_segments: segments_by_cycle[segment["周期编号"]].append(segment) result = [] for cycle_id, series in sorted(by_cycle.items(), key=lambda item: item[1][0]["_dt"]): series.sort(key=lambda row: row["_dt"]) segments = segments_by_cycle[cycle_id] levels = ["正常波动", "轻微", "异常", "严重异常"] fields = { "机组": series[0]["机组"], "批次": series[0]["批次"], "周期编号": cycle_id, "周期开始小时": series[0]["小时"], "周期结束小时": series[-1]["小时"], "周期有效运行小时数": len(series), "周期基线油压": series[0]["周期基线油压"], "正常波动带": series[0].get("正常波动带", ""), "正常波动带": series[0].get("正常波动带", ""), "阶段演化顺序": "→".join(segment["压力阶段"] for segment in segments), "周期最低小时油压": min(row["小时油压均值"] for row in series), "是否达到轻微": "是" if any(row["压力等级"] == "轻微" for row in series) else "否", "是否达到异常": "是" if any(row["压力等级"] == "异常" for row in series) else "否", "是否达到严重异常": "是" if any(row["压力等级"] == "严重异常" for row in series) else "否", "周期短时低压小时数": sum(row["立即报警"] != "" and row["压力等级"] == "正常" for row in series), "正常波动段数": sum(segment["压力阶段"] == "正常波动" for segment in segments), "正常波动小时数": sum(segment["有效运行小时数"] for segment in segments if segment["压力阶段"] == "正常波动"), } for level in levels: selected = [segment for segment in segments if segment["压力阶段"] == level] fields[f"{level}阶段开始小时"] = selected[0]["开始小时"] if selected else "" fields[f"{level}阶段结束小时"] = selected[-1]["结束小时"] if selected else "" fields[f"{level}阶段运行小时数"] = sum(segment["有效运行小时数"] for segment in selected) result.append(fields) return result def summarize_oil_cycles(daily: list[dict], segments: list[dict], hourly_segments: list[dict]): by_cycle = defaultdict(list) for row in daily: by_cycle[row["周期编号"]].append(row) segments_by_cycle = defaultdict(list) for segment in segments: segments_by_cycle[segment["周期编号"]].append(segment) hourly_segments_by_cycle = defaultdict(list) for segment in hourly_segments: hourly_segments_by_cycle[segment["周期编号"]].append(segment) result = [] batches = sorted({row["批次"] for row in daily}) for batch in batches: batch_cycles = sorted( ((cycle_id, series) for cycle_id, series in by_cycle.items() if series[0]["批次"] == batch), key=lambda item: item[1][0]["_date"], ) previous_summary = None for cycle_id, series in batch_cycles: series.sort(key=lambda row: row["_date"]) cycle_segments = segments_by_cycle[cycle_id] cycle_hourly_segments = hourly_segments_by_cycle[cycle_id] phase_segments = { phase: [segment for segment in cycle_segments if segment["压力阶段"] == phase] for phase in ("正常", "轻微", "异常", "严重异常") } start = series[0]["_date"] end = series[-1]["_date"] previous_gap = "" recurrence = "否" if previous_summary is not None: previous_gap = (start - previous_summary["_end"]).days if ( previous_summary["达到异常阶段"] == "是" and any(row["压力阶段"] in ("异常", "严重异常") for row in series) ): recurrence = "疑似恢复后再次出现" abnormal = any(row["压力阶段"] in ("异常", "严重异常") for row in series) severe = any(row["压力阶段"] == "严重异常" for row in series) order = "→".join(segment["压力阶段"] for segment in cycle_segments) def phase_value(phase: str, field: str): values = phase_segments[phase] return values[0][field] if values else "" def phase_hours(phase: str): return sum(segment["阶段运行小时数"] for segment in phase_segments[phase]) def hourly_phase_value(phase: str, field: str): values = [segment for segment in cycle_hourly_segments if segment["压力阶段"] == phase] if not values: return "" values.sort(key=lambda segment: segment["开始小时"]) return values[0][field] if field == "开始小时" else values[-1][field] def hourly_phase_hours(phase: str): return sum( segment["有效运行小时数"] for segment in cycle_hourly_segments if segment["压力阶段"] == phase ) summary = { "机组": series[0]["机组"], "批次": series[0]["批次"], "周期编号": cycle_id, "周期开始日期": series[0]["日期"], "周期结束日期": series[-1]["日期"], "周期有效运行天数": len(series), "周期有效运行小时数": sum(row["有效运行小时数"] for row in series), "周期基线油压": series[0]["周期基线油压"], "正常阶段开始": phase_value("正常", "开始日期"), "正常阶段结束": phase_value("正常", "结束日期"), "正常阶段运行小时数": phase_hours("正常"), "轻微阶段开始": phase_value("轻微", "开始日期"), "轻微阶段结束": phase_value("轻微", "结束日期"), "轻微阶段运行小时数": phase_hours("轻微"), "异常阶段开始": phase_value("异常", "开始日期"), "异常阶段结束": phase_value("异常", "结束日期"), "异常阶段运行小时数": phase_hours("异常"), "严重异常阶段开始": phase_value("严重异常", "开始日期"), "严重异常阶段结束": phase_value("严重异常", "结束日期"), "严重异常阶段运行小时数": phase_hours("严重异常"), "周期最低日油压": min(row["日油压中位数"] for row in series), "周期最低压力阶段": max( series, key=lambda row: {"正常": 0, "轻微": 1, "异常": 2, "严重异常": 3} .get(row["压力阶段"], 0), )["压力阶段"], "达到异常阶段": "是" if abnormal else "否", "达到严重异常阶段": "是" if severe else "否", "阶段演化顺序": order, "周期变化形态": "渐进下降型" if any(row["变化形态"] == "渐进下降型" for row in series) else "稳定/短时波动型", "周期起点依据": series[0]["周期起点依据"], "与上一周期结束间隔天数": previous_gap, "是否疑似恢复后复发": recurrence, "周期短时低压小时数": sum(row["短时低压小时数"] for row in series), "小时分级轻微开始": hourly_phase_value("轻微", "开始小时"), "小时分级轻微结束": hourly_phase_value("轻微", "结束小时"), "小时分级轻微运行小时数": hourly_phase_hours("轻微"), "小时分级异常开始": hourly_phase_value("异常", "开始小时"), "小时分级异常结束": hourly_phase_value("异常", "结束小时"), "小时分级异常运行小时数": hourly_phase_hours("异常"), "小时分级严重异常开始": hourly_phase_value("严重异常", "开始小时"), "小时分级严重异常结束": hourly_phase_value("严重异常", "结束小时"), "小时分级严重异常运行小时数": hourly_phase_hours("严重异常"), "_end": end, } result.append(summary) previous_summary = summary for row in result: row.pop("_end", None) return result def vibration_hour_signals(rows: list[dict], configs: dict[int, dict]): signals = [] grouped = defaultdict(list) for row in rows: if row["运行覆盖率"] >= MIN_HOUR_COVERAGE: grouped[row["批次"]].append(row) for batch, series in grouped.items(): series.sort(key=lambda row: row["_dt"]) config = configs[batch] valid = [] last_valid_dt = None for row in series: if row["运行覆盖率"] < MIN_HOUR_COVERAGE: continue restart = last_valid_dt is None or row["_dt"] - last_valid_dt > timedelta(hours=VIBRATION_RESTART_GAP_HOURS) if restart: # The first complete running hour after a restart is a # transition hour. Keep older valid hours for the long-term # robust reference, but do not evaluate this transition hour. last_valid_dt = row["_dt"] continue baseline_start = row["_dt"] - timedelta(days=VIBRATION_BASELINE_DAYS) prior = [item for item in valid if item["_dt"] >= baseline_start] side_signals = [] for side in ("coupling", "chain"): avg_center, avg_scale = robust_baseline([item[f"{side}_avg"] for item in prior]) max_center, max_scale = robust_baseline([item[f"{side}_max"] for item in prior]) current_avg = row[f"{side}_avg"] current_max = row[f"{side}_max"] if current_max is None: continue avg_z = ((current_avg - avg_center) / avg_scale) if avg_center is not None and current_avg is not None else 0.0 peak_z = (current_max - max_center) / max_scale if max_center is not None and max_scale else 0.0 running = max(row["运行样本数"], 1) high_ratio = (row[f"{side}_high_count"] or 0) / running high_high_ratio = (row[f"{side}_high_high_count"] or 0) / running limits = config[side] local_prior = valid[:6] if len(valid) >= 6 else [] local_avg = median_value([item[f"{side}_avg"] for item in local_prior]) local_max = median_value([item[f"{side}_max"] for item in local_prior]) local_rise = False severe = ( high_high_ratio >= 0.01 or current_max >= limits["high_high"] or peak_z >= 8 ) abnormal = ( high_ratio >= 0.01 or current_max >= limits["high"] or peak_z >= 5 or avg_z >= 4 or local_rise ) mild = peak_z >= 4 or avg_z >= 3 or local_rise if severe or abnormal or mild: side_signals.append({ "side": side, "avg_z": avg_z, "peak_z": peak_z, "high_ratio": high_ratio, "high_high_ratio": high_high_ratio, "severe": severe, "abnormal": abnormal, "value": current_max, "avg": current_avg, "prior_max_center": max_center, "prior_max_scale": max_scale, "prior_avg_center": avg_center, }) if side_signals: strongest = max(side_signals, key=lambda item: (item["severe"], item["peak_z"], item["avg_z"])) 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 reasons = [] for item in side_signals: label = "联轴器端" if item["side"] == "coupling" else "链轮端" reasons.append( f"{label}均值偏离{item['avg_z']:.1f}倍、峰值偏离{item['peak_z']:.1f}倍、" f"高报警比例{item['high_ratio']:.1%}" ) signals.append({ "批次": batch, "机组": config["unit"], "_dt": row["_dt"], "_level": signal_level, "_severe": signal_level == 3, "_side_signals": side_signals, "_strongest": strongest, "_reason": ";".join(reasons), }) valid.append(row) last_valid_dt = row["_dt"] return signals def merge_vibration_signals(signals: list[dict], configs: dict[int, dict]): grouped = defaultdict(list) for signal in signals: grouped[signal["批次"]].append(signal) events = [] for batch, series in grouped.items(): series.sort(key=lambda item: item["_dt"]) current = [] for signal in series: if not current or signal["_dt"] - current[-1]["_dt"] <= timedelta(hours=VIBRATION_MERGE_GAP_HOURS): current.append(signal) else: events.append(build_vibration_event(current, configs[batch])) current = [signal] if current: events.append(build_vibration_event(current, configs[batch])) events.sort(key=lambda row: (_unit_sort_key(row["机组"]), row["开始小时"], row["批次"])) unit_indexes = defaultdict(int) for event in events: unit_indexes[event["机组"]] += 1 event["周期编号"] = f'{event["机组"]}-V{unit_indexes[event["机组"]]:02d}' filtered = [] for event in events: # A single post-start relative-rise signal is not sufficient evidence; # absolute severe signals remain eligible on their own. if event["_严重信号小时数"] > 0 or event["_触发小时数"] >= 2: filtered.append(event) return sorted(filtered, key=_output_sort_key) 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): duration_score = min(40.0, max(0.0, (duration_hours - 1) / 12 * 40)) avg_score = min(20.0, max(0.0, (avg_change - 0.10) / 0.40 * 20)) if avg_change >= 0.10 else 0.0 high = config[side]["high"] high_high = config[side]["high_high"] if max_value < high: max_score = min(10.0, max(0.0, (max_value / max(high, 1e-6) - 0.8) / 0.2 * 10)) elif max_value < high_high: max_score = 10.0 + (max_value - high) / max(high_high - high, 1e-6) * 10.0 else: max_score = 20.0 peak_score = min(20.0, max(0.0, (peak_z - 4.0) / 8.0 * 20.0)) if peak_z >= 4.0 else 0.0 total = duration_score + avg_score + max_score + peak_score hard_severe = max_value >= high_high or high_high_hours > 0 or (peak_z >= 8.0 and valid_hours >= 2) if hard_severe or total >= 60.0: level = "严重异常" elif total >= 30.0: level = "异常" else: level = "轻微" return total, level, duration_score, avg_score, max_score, peak_score def build_vibration_event(signals: list[dict], config: dict): start = signals[0]["_dt"] end = signals[-1]["_dt"] all_side = [item for signal in signals for item in signal["_side_signals"]] strongest = max(all_side, key=lambda item: (item["severe"], item["peak_z"], item["value"])) main_side = "联轴器端" if strongest["side"] == "coupling" else "链轮端" max_value = max(item["value"] for item in all_side if item["value"] is not None) max_peak_z = max(item["peak_z"] for item in all_side) high_hours = sum( 1 for signal in signals if any(item["high_ratio"] > 0 for item in signal["_side_signals"]) ) high_high_hours = sum( 1 for signal in signals if any(item["high_high_ratio"] > 0 for item in signal["_side_signals"]) ) severe_hours = sum(1 for signal in signals if signal["_severe"]) duration_hours = int((end - start).total_seconds() / 3600) + 1 running_days = len({signal["_dt"].date() for signal in signals}) side_avg_changes = [] for side in ("coupling", "chain"): values = [item["avg"] for item in all_side if item["side"] == side and item["avg"] is not None] prior_values = [item["prior_avg_center"] for item in all_side if item["side"] == side and item["prior_avg_center"] is not None] if values and prior_values and median(prior_values) != 0: side_avg_changes.append((median(values) / median(prior_values) - 1, side)) avg_change, _ = max(side_avg_changes, key=lambda item: item[0], default=(0.0, strongest["side"])) peak_ratio = max_value / max(strongest["prior_max_center"], 1e-6) if max_value >= config[strongest["side"]]["high_high"] or severe_hours > 0 and (high_high_hours > 0 or max_peak_z >= 8): level = "严重异常" elif max_value >= config[strongest["side"]]["high"] and high_hours >= 2: level = "严重异常" elif max_value >= config[strongest["side"]]["high"] or max_peak_z >= 5 or avg_change >= 0.10: level = "异常" else: level = "轻微" if peak_ratio >= 1.5 and avg_change < 0.25: shape = "间歇冲击型" elif avg_change >= 0.10 and len(signals) >= 3: shape = "持续抬升型" else: shape = "混合型" peak_signal = max(signals, key=lambda signal: signal["_strongest"]["peak_z"]) first_signal = signals[0] last_signal = signals[-1] first_value = first_signal["_strongest"]["value"] last_value = last_signal["_strongest"]["value"] baseline = strongest["prior_max_center"] scale = max(strongest.get("prior_max_scale", 0.0), 0.0) normal_band = f"{baseline - 2 * scale:.3f}~{baseline + 2 * scale:.3f} mm/s" if baseline is not None else "" consecutive = 0 longest_consecutive = 0 for signal in signals: if signal["_level"] >= 2: consecutive += 1 longest_consecutive = max(longest_consecutive, consecutive) else: consecutive = 0 score, level, duration_score, avg_score, max_score, peak_score = vibration_score( duration_hours, len(signals), avg_change, max_value, max_peak_z, high_hours, high_high_hours, config, strongest["side"] ) confirmed = "是" if level != "轻微" else "否" if shape == "持续抬升型": basis = f"{main_side}振动均值持续抬升,阶段内{longest_consecutive or len(signals)}个有效运行小时触发" elif shape == "间歇冲击型": basis = f"{main_side}重复出现振动冲击,阶段内{len(signals)}个有效运行小时触发" else: basis = f"{main_side}振动峰值或均值偏离近期基线,阶段内{len(signals)}个有效运行小时触发" if high_high_hours: basis += f";高高报警小时{high_high_hours}个" elif high_hours: basis += f";高报警小时{high_hours}个" basis += f";最大振动值{max_value:.3f}mm/s,峰值高于近期基线{max_peak_z:.2f}个稳健尺度(峰值/基线{peak_ratio:.2f}倍);综合评分{score:.1f}分" return { "机组": config["unit"], "批次": config["batch"], "周期编号": "", "振动阶段": level, "开始小时": start.strftime("%Y-%m-%d %H:%M:%S"), "结束小时": end.strftime("%Y-%m-%d %H:%M:%S"), "持续自然小时数": duration_hours, "有效运行小时数": len(signals), "近期基线振动": baseline, "正常波动带": normal_band, "阶段开始小时振动": first_value, "阶段结束小时振动": last_value, "阶段最大振动值": max_value, "阶段开始相对近期基线变化": (first_value / baseline - 1) if baseline else None, "阶段结束相对近期基线变化": (last_value / baseline - 1) if baseline else None, "阶段最大连续异常有效小时数": longest_consecutive, "阶段综合评分": score, "持续时间得分": duration_score, "平均振动得分": avg_score, "最大振动值得分": max_score, "峰值偏离得分": peak_score, "最强异常侧": main_side, "最强侧近期峰值基线": baseline, "最强侧近期稳健尺度": scale, "阶段最大峰值/基线比例": peak_ratio, "阶段高报警小时数": high_hours, "阶段高高报警小时数": high_high_hours, "_严重信号小时数": severe_hours, "_触发小时数": len(signals), "是否形成确认等级": confirmed, "阶段触发依据": basis, } def _unit_sort_key(value): match = re.search(r"[789]", str(value or "")) return int(match.group()) if match else 999 def _output_sort_key(row: dict): """Keep exported records grouped by unit, then ordered by their start time.""" time_field = next( (field for field in ("小时", "开始小时", "周期开始小时", "开始日期", "日期") if row.get(field)), None, ) time_value = str(row.get(time_field, "")) if time_field else "" return ( _unit_sort_key(row.get("机组")), time_value, int(row.get("批次") or 0), str(row.get("周期编号", "")), ) def write_csv(path: Path, rows: list[dict], fields: list[str], headers: dict[str, str]): path.parent.mkdir(parents=True, exist_ok=True) rows = sorted(rows, key=_output_sort_key) with path.open("w", encoding="utf-8-sig", newline="") as handle: writer = csv.DictWriter(handle, fieldnames=fields, extrasaction="ignore") writer.writerow({field: headers.get(field, field) for field in fields}) writer.writerows(rows) def format_number(value): if value is None or value == "": return "" if isinstance(value, float): return round(value, 6) return value def clean_rows(rows): cleaned = [] for row in rows: item = {} for key, value in row.items(): if key.startswith("_"): continue item[key] = format_number(value) cleaned.append(item) return cleaned def write_thresholds(path: Path, configs: dict[int, dict]): rows = [] for batch in sorted(configs): config = configs[batch] rows.extend([ { "机组": config["unit"], "批次": batch, "点位": "YSJ_5", "点位描述": config["metadata"]["oil_description"], "单位": config["metadata"]["oil_units"], "低报警阈值": config["oil"]["low"], "低低报警阈值": config["oil"]["low_low"], "高报警阈值": config["oil"]["high"], "高高报警阈值": config["oil"]["high_high"], "低报警配置": config["alarm_sources"]["oil"]["low"], "低低报警配置": config["alarm_sources"]["oil"]["low_low"], "高报警配置": config["alarm_sources"]["oil"]["high"], "高高报警配置": config["alarm_sources"]["oil"]["high_high"], }, { "机组": config["unit"], "批次": batch, "点位": "YSJ_10", "点位描述": config["metadata"]["coupling_description"], "单位": config["metadata"]["coupling_units"], "低报警阈值": "", "低低报警阈值": "", "高报警阈值": config["coupling"]["high"], "高高报警阈值": config["coupling"]["high_high"], "低报警配置": "", "低低报警配置": "", "高报警配置": config["alarm_sources"]["coupling"]["high"], "高高报警配置": config["alarm_sources"]["coupling"]["high_high"], }, { "机组": config["unit"], "批次": batch, "点位": "YSJ_11", "点位描述": config["metadata"]["chain_description"], "单位": config["metadata"]["chain_units"], "低报警阈值": "", "低低报警阈值": "", "高报警阈值": config["chain"]["high"], "高高报警阈值": config["chain"]["high_high"], "低报警配置": "", "低低报警配置": "", "高报警配置": config["alarm_sources"]["chain"]["high"], "高高报警配置": config["alarm_sources"]["chain"]["high_high"], }, ]) fields = [ "机组", "批次", "点位", "点位描述", "单位", "低报警阈值", "低低报警阈值", "高报警阈值", "高高报警阈值", "低报警配置", "低低报警配置", "高报警配置", "高高报警配置", ] write_csv(path, rows, fields, {field: field for field in fields}) def main(): parser = argparse.ArgumentParser(description="只读生成 PKS 润滑油压力长期演化和振动短期异常段") parser.add_argument("--batch", nargs="+", type=int, help="PKS 导入批次 ID,默认扫描 compressor_data_import 中可识别的全部批次") parser.add_argument("--start", default="2025-01-01 00:00:00") parser.add_argument("--end", default="2026-09-01 00:00:00") parser.add_argument("--chunk-days", type=int, default=14) parser.add_argument("--output-dir", default=str(Path(__file__).resolve().parents[1] / "cache" / "pks_scanner" / "evolution")) args = parser.parse_args() start, end = parse_dt(args.start), parse_dt(args.end) if end <= start or args.chunk_days < 1: parser.error("时间范围或 chunk-days 无效") print("只读模式:不修改数据库;只分析 YSJ_5、YSJ_10、YSJ_11、YSJ_41") batch_to_unit = load_pks_batches() selected_batches = args.batch or sorted(batch_to_unit) unknown_batches = sorted(set(selected_batches) - set(batch_to_unit)) if unknown_batches: parser.error(f"以下批次不是 compressor_data_import 中可识别的 PKS 批次:{unknown_batches}") configs = load_alarm_config({batch: batch_to_unit[batch] for batch in selected_batches}) print(f"已从 compressor_data_import 解析 {len(selected_batches)} 个 PKS 批次的机组:{batch_to_unit}") print(f"振动基线最多使用最近 {VIBRATION_BASELINE_DAYS} 天的有效运行小时") batch_ranges = load_batch_ranges(selected_batches, start, end) grouped_batches = defaultdict(list) for batch in selected_batches: if batch in batch_ranges: grouped_batches[batch_to_unit[batch]].append(batch) all_rows = [] for unit, batches in sorted(grouped_batches.items(), key=lambda item: _unit_sort_key(item[0])): unit_start = min(batch_ranges[batch][0] for batch in batches) unit_end = max(batch_ranges[batch][1] for batch in batches) config = configs[batches[0]] all_rows.extend(fetch_hourly(unit, sorted(batches), unit_start, unit_end, args.chunk_days, config)) all_rows.sort(key=lambda row: (_unit_sort_key(row["机组"]), row["_dt"], row["批次"])) oil_hourly_levels = classify_oil_hourly_segments(all_rows, configs) oil_hour_stage_segments = build_oil_hour_stage_segments(oil_hourly_levels) oil_cycles = summarize_hourly_oil_cycles(oil_hourly_levels, oil_hour_stage_segments) signals = vibration_hour_signals(all_rows, configs) vibration_events = merge_vibration_signals(signals, configs) out = Path(args.output_dir) write_thresholds(out / "报警阈值配置.csv", configs) hourly_fields = [ "批次", "机组", "小时", "样本数", "运行样本数", "运行覆盖率", "oil_avg", "oil_min", "oil_max", "coupling_avg", "coupling_max", "chain_avg", "chain_max", "oil_low_count", "oil_low_low_count", "coupling_high_count", "coupling_high_high_count", "chain_high_count", "chain_high_high_count", ] hourly_headers = { "oil_avg": "润滑油压力均值", "oil_min": "润滑油压力最小值", "oil_max": "润滑油压力最大值", "coupling_avg": "联轴器端振动均值", "coupling_max": "联轴器端振动最大值", "chain_avg": "链轮端振动均值", "chain_max": "链轮端振动最大值", "oil_low_count": "润滑油低报警样本数", "oil_low_low_count": "润滑油低低报警样本数", "coupling_high_count": "联轴器端高报警样本数", "coupling_high_high_count": "联轴器端高高报警样本数", "chain_high_count": "链轮端高报警样本数", "chain_high_high_count": "链轮端高高报警样本数", "批次": "批次", "机组": "机组", "小时": "小时", "样本数": "样本数", "运行样本数": "运行样本数", "运行覆盖率": "运行覆盖率", } hourly_level_fields = [ "机组", "批次", "周期编号", "小时", "运行样本数", "运行覆盖率", "小时油压均值", "小时油压最小值", "小时油压最大值", "周期基线油压", "近期24有效小时基线油压", "相对周期基线变化", "相对近期基线变化", "近6有效小时油压斜率", "连续下降有效小时数", "历史有效运行小时数", "低报警样本比例", "低低报警样本比例", "低报警阈值", "低低报警阈值", "是否低报警小时", "是否低低报警小时", "压力等级", "原始小时等级", "立即报警", "正常波动标记", "低低报警连续小时数", "是否纳入等级判定", "小时判定依据", ] write_csv(out / "润滑油压力小时等级.csv", clean_rows(oil_hourly_levels), hourly_level_fields, {field: field for field in hourly_level_fields}) hour_stage_fields = [ "机组", "批次", "周期编号", "压力阶段", "开始小时", "结束小时", "持续自然小时数", "有效运行小时数", "周期基线油压", "正常波动带", "阶段开始小时油压", "阶段结束小时油压", "阶段最低小时油压", "阶段开始相对周期基线变化", "阶段结束相对周期基线变化", "阶段最大连续下降有效小时数", "阶段低报警小时数", "阶段低低报警小时数", "是否形成确认等级", "阶段触发依据", ] write_csv(out / "润滑油压力小时阶段段.csv", clean_rows(oil_hour_stage_segments), hour_stage_fields, {field: field for field in hour_stage_fields}) cycle_fields = [ "机组", "批次", "周期编号", "周期开始小时", "周期结束小时", "周期有效运行小时数", "周期基线油压", "正常波动带", "阶段演化顺序", "周期最低小时油压", "是否达到轻微", "是否达到异常", "是否达到严重异常", "周期短时低压小时数", "正常波动段数", "正常波动小时数", "轻微阶段开始小时", "轻微阶段结束小时", "轻微阶段运行小时数", "异常阶段开始小时", "异常阶段结束小时", "异常阶段运行小时数", "严重异常阶段开始小时", "严重异常阶段结束小时", "严重异常阶段运行小时数", ] write_csv(out / "润滑油压力演化周期.csv", clean_rows(oil_cycles), cycle_fields, {field: field for field in cycle_fields}) vibration_fields = [ "机组", "批次", "周期编号", "振动阶段", "开始小时", "结束小时", "持续自然小时数", "有效运行小时数", "近期基线振动", "正常波动带", "阶段开始小时振动", "阶段结束小时振动", "阶段最大振动值", "阶段开始相对近期基线变化", "阶段结束相对近期基线变化", "阶段最大连续异常有效小时数", "阶段高报警小时数", "阶段高高报警小时数", "是否形成确认等级", "阶段触发依据", "阶段综合评分", "持续时间得分", "平均振动得分", "最大振动值得分", "峰值偏离得分", "最强异常侧", "最强侧近期峰值基线", "最强侧近期稳健尺度", "阶段最大峰值/基线比例", ] write_csv(out / "振动异常段.csv", clean_rows(vibration_events), vibration_fields, {field: field for field in vibration_fields}) print(f"小时特征:{len(all_rows):,} 行") print(f"润滑油压力小时等级:{len(oil_hourly_levels):,} 行") print(f"润滑油压力小时阶段段:{len(oil_hour_stage_segments):,} 段") print(f"润滑油压力周期:{len(oil_cycles):,} 个") print(f"振动异常段:{len(vibration_events):,} 段") print(f"输出目录:{out}") if __name__ == "__main__": main()