#!/usr/bin/env python3 """Second PKS scanner focused on pressure usefulness. This is intentionally a new scanner. It does not modify the first scanner or the database. Pressure alarm limits come from the site_point configuration confirmed for YSJ_5..YSJ_9; the output keeps both absolute alarm evidence and relative, recent-baseline evidence so they can be reviewed separately. """ from __future__ import annotations import argparse import csv import math import os 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 BATCH_TO_UNIT = {30: "7号机", 31: "8号机", 32: "9号机"} EXPECTED_PER_HOUR = 720 MIN_COVERAGE = 0.80 MIN_HISTORY_HOURS = 24 # AlarmLimit1..4 in site_point. The same engineering limits are configured # for all three units, while the measured values remain unit-specific. PRESSURE_LIMITS = { "oil": {"field": "oil", "name": "润滑油压力", "low": 0.34, "low_low": 0.31, "high": 0.68, "high_high": 0.72}, "inlet": {"field": "inlet", "name": "进气压力", "low": 0.45, "low_low": 0.40, "high": 0.95, "high_high": 1.00}, "p1": {"field": "p1", "name": "一级排气压力", "high": 2.60, "high_high": 2.75}, "p2": {"field": "p2", "name": "二级排气压力", "high": 4.90, "high_high": 5.00}, "p3": {"field": "p3", "name": "三级排气压力", "high": 6.20, "high_high": 6.30}, } PRESSURE_KEYS = ("oil", "inlet", "p1", "p2", "p3") EXHAUST_KEYS = ("p1", "p2", "p3") SQL = """ SELECT FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) AS hour_start, COUNT(*) AS samples, SUM(YSJ_41 > 0) 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_6 END) AS inlet_avg, MIN(CASE WHEN YSJ_41 > 0 THEN YSJ_6 END) AS inlet_min, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_6 END) AS inlet_max, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_7 END) AS p1_avg, MIN(CASE WHEN YSJ_41 > 0 THEN YSJ_7 END) AS p1_min, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_7 END) AS p1_max, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_8 END) AS p2_avg, MIN(CASE WHEN YSJ_41 > 0 THEN YSJ_8 END) AS p2_min, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_8 END) AS p2_max, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_9 END) AS p3_avg, MIN(CASE WHEN YSJ_41 > 0 THEN YSJ_9 END) AS p3_min, MAX(CASE WHEN YSJ_41 > 0 THEN YSJ_9 END) AS p3_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, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_14 END) AS oil_temp_avg, AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_40 END) AS valve_avg, SUM(YSJ_41 > 0 AND YSJ_5 <= 0.34) AS oil_low_count, SUM(YSJ_41 > 0 AND YSJ_5 <= 0.31) AS oil_low_low_count, SUM(YSJ_41 > 0 AND YSJ_5 >= 0.68) AS oil_high_count, SUM(YSJ_41 > 0 AND YSJ_6 <= 0.45) AS inlet_low_count, SUM(YSJ_41 > 0 AND YSJ_6 >= 0.95) AS inlet_high_count, SUM(YSJ_41 > 0 AND YSJ_7 >= 2.60) AS p1_high_count, SUM(YSJ_41 > 0 AND YSJ_7 >= 2.75) AS p1_high_high_count, SUM(YSJ_41 > 0 AND YSJ_8 >= 4.90) AS p2_high_count, SUM(YSJ_41 > 0 AND YSJ_8 >= 5.00) AS p2_high_high_count, SUM(YSJ_41 > 0 AND YSJ_9 >= 6.20) AS p3_high_count, SUM(YSJ_41 > 0 AND YSJ_9 >= 6.30) AS p3_high_high_count, SUM(YSJ_41 > 0 AND YSJ_10 >= 7.1) AS coupling_alarm_count, SUM(YSJ_41 > 0 AND YSJ_11 >= 7.1) AS chain_alarm_count FROM pks_long_sample WHERE import_batch_id = %s AND sample_time >= %s AND sample_time < %s GROUP BY FROM_UNIXTIME((UNIX_TIMESTAMP(sample_time) DIV 3600) * 3600) ORDER BY hour_start """ 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): while start < end: nxt = min(start + timedelta(days=days), end) yield start, nxt start = 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 fetch_hourly(batch, start, end, chunk_days): rows = [] connection = connect() try: with connection.cursor() as cursor: for index, (lo, hi) in enumerate(chunks(start, end, chunk_days), 1): cursor.execute("EXPLAIN " + SQL, (batch, lo, hi)) plan = cursor.fetchone() if not plan or plan.get("key") not in ("PRIMARY",): raise RuntimeError(f"安全检查失败:batch={batch} chunk={lo} EXPLAIN 未使用 PRIMARY: {plan}") cursor.execute(SQL, (batch, lo, hi)) for raw in cursor.fetchall(): row = {"batch": batch, "unit": BATCH_TO_UNIT[batch], "hour": raw["hour_start"].strftime("%Y-%m-%d %H:%M:%S"), "samples": int(raw["samples"] or 0), "running_samples": int(raw["running_samples"] or 0)} row["coverage"] = row["running_samples"] / EXPECTED_PER_HOUR for key, value in raw.items(): if key not in ("hour_start", "samples", "running_samples"): row[key] = number(value) rows.append(row) print(f"batch={batch} chunk={index} {lo}..{hi}: {len(rows):,} total hourly rows") finally: connection.close() return rows def baseline(values): values = [v for v in values if v is not None] if len(values) < 12: return None, None center = median(values) mad = median(abs(v - center) for v in values) return center, max(1.4826 * mad, abs(center) * 0.01, 1e-6) def linear_slope(values): pairs = [(i, v) for i, v in enumerate(values) if v is not None] if len(pairs) < 4: return None xbar = sum(x for x, _ in pairs) / len(pairs) ybar = sum(y for _, y in pairs) / len(pairs) den = sum((x - xbar) ** 2 for x, _ in pairs) return None if not den else sum((x - xbar) * (y - ybar) for x, y in pairs) / den def add_candidate(out, row, kind, score, reason, history): out.append({"candidate_id": f"{row['unit']}-{kind}-{row['hour']}", "unit": row["unit"], "batch": row["batch"], "anomaly_type": kind, "start_time": row["hour"], "end_time": row["hour"], "peak_time": row["hour"], "score": round(min(max(score, 0), 12), 3), "history_valid_hours": history, "coverage": round(row["coverage"], 3), "running_samples": row["running_samples"], "repeat_count": 1, "candidate_duration_hours": 1, "reason": reason}) def pressure_context(row, prior): """Describe exhaust-pressure evidence without creating pressure-only events.""" changes = [] for key in EXHAUST_KEYS: current = row.get(f"{key}_avg") center, scale = baseline([x.get(f"{key}_avg") for x in prior]) if current is None or center is None: continue z = (current - center) / scale if abs(z) >= 2.5: direction = "上升" if z > 0 else "下降" changes.append(f"{PRESSURE_LIMITS[key]['name']}{direction}{abs(z):.1f}倍尺度") if len(changes) >= 2: return ";压力联合证据:多级同步变化(" + ",".join(changes) + ")" if changes: return ";压力联合证据:" + changes[0] return ";压力联合证据:排气压力无明显同步异常" def merge_candidates(candidates): """Merge hourly triggers of the same phenomenon into reviewable segments.""" candidates.sort(key=lambda x: (x["batch"], x["anomaly_type"], x["start_time"])) merged = [] for item in candidates: if merged and item["batch"] == merged[-1]["batch"] and item["anomaly_type"] == merged[-1]["anomaly_type"]: previous = datetime.strptime(merged[-1]["end_time"], "%Y-%m-%d %H:%M:%S") current = datetime.strptime(item["start_time"], "%Y-%m-%d %H:%M:%S") if current - previous <= timedelta(hours=6): merged[-1]["end_time"] = item["end_time"] merged[-1]["candidate_duration_hours"] = int( (current - datetime.strptime(merged[-1]["start_time"], "%Y-%m-%d %H:%M:%S")).total_seconds() / 3600 ) + 1 merged[-1]["repeat_count"] += 1 merged[-1]["coverage"] = min(merged[-1]["coverage"], item["coverage"]) merged[-1]["running_samples"] = max(merged[-1]["running_samples"], item["running_samples"]) if item["score"] > merged[-1]["score"]: merged[-1]["score"] = item["score"] merged[-1]["peak_time"] = item["peak_time"] if item["reason"] not in merged[-1]["reason"]: merged[-1]["reason"] += "; " + item["reason"] continue merged.append(dict(item)) for item in merged: item["candidate_id"] = f"{item['unit']}-{item['anomaly_type']}-{item['start_time']}" return sorted(merged, key=lambda x: (x["batch"], x["start_time"], x["anomaly_type"])) def detect(rows, baseline_hours): result = [] grouped = defaultdict(list) for row in rows: grouped[row["batch"]].append(row) for batch, series in grouped.items(): series.sort(key=lambda x: x["hour"]) valid = [x for x in series if x["coverage"] >= MIN_COVERAGE] for index, row in enumerate(series): if row["coverage"] < MIN_COVERAGE: continue prior = [x for x in valid if x["hour"] < row["hour"]][-baseline_hours:] if len(prior) < MIN_HISTORY_HOURS: continue history = len(prior) oil_center, oil_scale = baseline([x.get("oil_avg") for x in prior]) oil = row.get("oil_avg") oil_slope = linear_slope([x.get("oil_avg") for x in prior[-6:]] + [oil]) low_ratio = (row.get("oil_low_count") or 0) / max(row["running_samples"], 1) low_low_ratio = (row.get("oil_low_low_count") or 0) / max(row["running_samples"], 1) if low_low_ratio >= 0.05 or low_ratio >= 0.20: add_candidate(result, row, "润滑油压力异常", 5 + 5 * min(low_low_ratio * 4, 1), f"YSJ_5低报警比例{low_ratio:.1%},低低报警比例{low_low_ratio:.1%}{pressure_context(row, prior)}", history) if oil_center is not None and oil is not None and oil_slope is not None: deviation = (oil_center - oil) / oil_scale if deviation >= 2.5 and oil_slope < 0: add_candidate(result, row, "润滑油压力异常", deviation + min(abs(oil_slope) / oil_scale, 5), f"油压均值{oil:.3f},低于近期基线{deviation:.1f}倍尺度,近6小时斜率{oil_slope:.5f}{pressure_context(row, prior)}", history) abnormal_exhaust = [] for key in EXHAUST_KEYS: current = row.get(f"{key}_avg") center, scale = baseline([x.get(f"{key}_avg") for x in prior]) if current is None or center is None: continue z = abs(current - center) / scale high_ratio = (row.get(f"{key}_high_count") or 0) / max(row["running_samples"], 1) if z >= 3: abnormal_exhaust.append((key, current, center, z)) for side, avg_key, max_key in (("联轴器端", "coupling_avg", "coupling_max"), ("链轮端", "chain_avg", "chain_max")): center, scale = baseline([x.get(avg_key) for x in prior]) max_center, max_scale = baseline([x.get(max_key) for x in prior]) if center is None or max_center is None or row.get(max_key) is None: continue avg_z = (row.get(avg_key) - center) / scale if row.get(avg_key) is not None else 0 peak_z = (row[max_key] - max_center) / max_scale alarm_ratio = (row.get(f"{('coupling' if side == '联轴器端' else 'chain')}_alarm_count") or 0) / max(row["running_samples"], 1) if alarm_ratio >= 0.01 or peak_z >= 4 or avg_z >= 3: add_candidate(result, row, "振动异常(压力联合评估)", min(max(peak_z, avg_z, 0) + min(alarm_ratio * 8, 3), 12), f"{side}均值偏离{avg_z:.1f},峰值偏离{peak_z:.1f},报警比例{alarm_ratio:.1%}{pressure_context(row, prior)}", history) return merge_candidates(result) def write_csv(path, rows, fields, headers): path.parent.mkdir(parents=True, exist_ok=True) with path.open("w", encoding="utf-8-sig", newline="") as handle: writer = csv.DictWriter(handle, fieldnames=fields, extrasaction="ignore") writer.writerow(headers) writer.writerows(rows) def main(): parser = argparse.ArgumentParser(description="只读扫描 PKS 压力与报警特征(第二版)") parser.add_argument("--batch", nargs="+", type=int, choices=sorted(BATCH_TO_UNIT), default=list(BATCH_TO_UNIT)) parser.add_argument("--start", default="2025-04-01 00:00:00") parser.add_argument("--end", default="2026-05-01 00:00:00") parser.add_argument("--chunk-days", type=int, default=14) parser.add_argument("--baseline-hours", type=int, default=336) parser.add_argument("--output-dir", default=str(Path(__file__).resolve().parents[1] / "cache" / "pks_scanner" / "v2_pressure")) args = parser.parse_args() start, end = parse_dt(args.start), parse_dt(args.end) if end <= start or args.chunk_days < 1 or args.baseline_hours < MIN_HISTORY_HOURS: parser.error("时间范围或参数无效") print("只读模式:不修改任何数据库数据;第一版程序未使用") rows = [] for batch in args.batch: rows.extend(fetch_hourly(batch, start, end, args.chunk_days)) rows.sort(key=lambda x: (x["batch"], x["hour"])) candidates = detect(rows, args.baseline_hours) out = Path(args.output_dir) hourly_fields = list(rows[0]) if rows else [] candidate_fields = ["candidate_id", "unit", "batch", "anomaly_type", "start_time", "end_time", "peak_time", "score", "history_valid_hours", "candidate_duration_hours", "repeat_count", "coverage", "running_samples", "reason"] hourly_headers = {key: key for key in hourly_fields} write_csv(out / "hourly_pressure_summary.csv", rows, hourly_fields, hourly_headers) write_csv(out / "pressure_fault_candidates.csv", candidates, candidate_fields, {key: key for key in candidate_fields}) print(f"小时摘要:{len(rows):,} 行") print(f"候选记录:{len(candidates):,} 条") print(f"输出目录:{out}") if __name__ == "__main__": main()