| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315 |
- #!/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()
|