|
|
@@ -0,0 +1,262 @@
|
|
|
+#!/usr/bin/env python3
|
|
|
+"""Scan running PKS data for interpretable fault-candidate segments.
|
|
|
+
|
|
|
+Read-only, bounded scanner. MySQL performs hourly aggregation inside each
|
|
|
+(batch, time-chunk) range; Python only receives the small hourly summary.
|
|
|
+The output is a candidate list for human review, not a confirmed diagnosis.
|
|
|
+"""
|
|
|
+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号机"}
|
|
|
+DEFAULT_COLUMNS = ("YSJ_5", "YSJ_7", "YSJ_8", "YSJ_9", "YSJ_10", "YSJ_11", "YSJ_14", "YSJ_40", "YSJ_41")
|
|
|
+EXPECTED_PER_HOUR = 720
|
|
|
+MIN_COVERAGE = 0.80
|
|
|
+
|
|
|
+HOURLY_FIELDS = [
|
|
|
+ "batch", "unit", "hour", "samples", "running_samples", "coverage",
|
|
|
+ "oil_avg", "oil_min", "oil_max", "p1_avg", "p2_avg", "p3_avg",
|
|
|
+ "coupling_avg", "coupling_max", "chain_avg", "chain_max",
|
|
|
+ "oil_temp_avg", "valve_avg",
|
|
|
+]
|
|
|
+
|
|
|
+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_7 END) AS p1_avg,
|
|
|
+ AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_8 END) AS p2_avg,
|
|
|
+ AVG(CASE WHEN YSJ_41 > 0 THEN YSJ_9 END) AS p3_avg,
|
|
|
+ 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
|
|
|
+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):
|
|
|
+ cursor = start
|
|
|
+ step = timedelta(days=days)
|
|
|
+ while cursor < end:
|
|
|
+ nxt = min(cursor + step, end)
|
|
|
+ yield cursor, nxt
|
|
|
+ cursor = nxt
|
|
|
+
|
|
|
+
|
|
|
+def to_float(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: int, start: datetime, end: datetime, chunk_days: int, limit_chunks: int | None):
|
|
|
+ rows = []
|
|
|
+ connection = connect()
|
|
|
+ try:
|
|
|
+ with connection.cursor() as cursor:
|
|
|
+ for index, (chunk_start, chunk_end) in enumerate(chunks(start, end, chunk_days), 1):
|
|
|
+ if limit_chunks is not None and index > limit_chunks:
|
|
|
+ break
|
|
|
+ cursor.execute("EXPLAIN " + SQL, (batch, chunk_start, chunk_end))
|
|
|
+ plan = cursor.fetchone()
|
|
|
+ if not plan or plan.get("key") not in ("PRIMARY",):
|
|
|
+ raise RuntimeError(f"安全检查失败:batch={batch} chunk={chunk_start} EXPLAIN 未使用 PRIMARY: {plan}")
|
|
|
+ cursor.execute(SQL, (batch, chunk_start, chunk_end))
|
|
|
+ chunk_rows = cursor.fetchall()
|
|
|
+ for row in chunk_rows:
|
|
|
+ item = {
|
|
|
+ "batch": batch, "unit": BATCH_TO_UNIT[batch],
|
|
|
+ "hour": row["hour_start"].strftime("%Y-%m-%d %H:%M:%S"),
|
|
|
+ "samples": int(row["samples"] or 0),
|
|
|
+ "running_samples": int(row["running_samples"] or 0),
|
|
|
+ }
|
|
|
+ item["coverage"] = item["running_samples"] / EXPECTED_PER_HOUR
|
|
|
+ for key in HOURLY_FIELDS[6:]:
|
|
|
+ item[key] = to_float(row[key])
|
|
|
+ rows.append(item)
|
|
|
+ print(f"batch={batch} chunk={index} {chunk_start}..{chunk_end}: {len(chunk_rows):,} hourly rows")
|
|
|
+ finally:
|
|
|
+ connection.close()
|
|
|
+ return rows
|
|
|
+
|
|
|
+
|
|
|
+def robust_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])
|
|
|
+ scale = max(1.4826 * mad, abs(center) * 0.01, 1e-9)
|
|
|
+ return center, scale
|
|
|
+
|
|
|
+
|
|
|
+def 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)
|
|
|
+ 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 add_reason(candidates, row, kind, score, reason):
|
|
|
+ candidates.append({
|
|
|
+ "candidate_id": f"{row['unit']}-{kind}-{row['hour']}",
|
|
|
+ "unit": row["unit"], "batch": row["batch"], "fault_hint": kind,
|
|
|
+ "start_time": row["hour"], "end_time": row["hour"],
|
|
|
+ "peak_time": row["hour"], "score": round(score, 3),
|
|
|
+ "coverage": round(row["coverage"], 3),
|
|
|
+ "running_samples": row["running_samples"], "reason": reason,
|
|
|
+ })
|
|
|
+
|
|
|
+
|
|
|
+def detect(rows, baseline_hours: int):
|
|
|
+ by_batch = defaultdict(list)
|
|
|
+ for row in rows:
|
|
|
+ by_batch[row["batch"]].append(row)
|
|
|
+ candidates = []
|
|
|
+ for batch, series in by_batch.items():
|
|
|
+ series.sort(key=lambda r: r["hour"])
|
|
|
+ valid = [r for r in series if r["coverage"] >= MIN_COVERAGE]
|
|
|
+ for i, row in enumerate(series):
|
|
|
+ if row["coverage"] < MIN_COVERAGE:
|
|
|
+ continue
|
|
|
+ prior = [r for r in valid if r["hour"] < row["hour"]][-baseline_hours:]
|
|
|
+ if len(prior) < 12:
|
|
|
+ continue
|
|
|
+
|
|
|
+ oil_center, oil_scale = robust_baseline([r["oil_avg"] for r in prior])
|
|
|
+ oil_values = [r["oil_avg"] for r in prior[-6:]] + [row["oil_avg"]]
|
|
|
+ oil_slope = slope(oil_values)
|
|
|
+ if oil_center is not None and row["oil_avg"] is not None and oil_slope is not None:
|
|
|
+ oil_drop = (oil_center - row["oil_avg"]) / oil_scale
|
|
|
+ if oil_drop >= 2.5 and oil_slope < 0:
|
|
|
+ add_reason(candidates, row, "润滑油压力下降", oil_drop + min(abs(oil_slope) / oil_scale, 5),
|
|
|
+ f"YSJ_5低于近期基线{oil_drop:.1f}倍稳健尺度,近6个有效小时斜率为{oil_slope:.5f}")
|
|
|
+
|
|
|
+ for side, avg_key, max_key in (("联轴器端", "coupling_avg", "coupling_max"), ("链轮端", "chain_avg", "chain_max")):
|
|
|
+ avg_center, avg_scale = robust_baseline([r[avg_key] for r in prior])
|
|
|
+ max_center, max_scale = robust_baseline([r[max_key] for r in prior])
|
|
|
+ if max_center is None or row[max_key] is None:
|
|
|
+ continue
|
|
|
+ peak_score = (row[max_key] - max_center) / max_scale
|
|
|
+ avg_score = ((row[avg_key] - avg_center) / avg_scale) if avg_center is not None and row[avg_key] is not None else 0
|
|
|
+ if peak_score >= 4:
|
|
|
+ kind = "活塞/机械冲击候选" if side == "链轮端" else "机械振动候选"
|
|
|
+ add_reason(candidates, row, kind, peak_score + max(avg_score, 0),
|
|
|
+ f"{side}峰值{row[max_key]:.3f},高于近期基线{peak_score:.1f}倍稳健尺度;均值偏离{avg_score:.1f}")
|
|
|
+ elif avg_score >= 3:
|
|
|
+ add_reason(candidates, row, "气阀相关振动候选", avg_score,
|
|
|
+ f"{side}平均振动{row[avg_key]:.3f},高于近期基线{avg_score:.1f}倍稳健尺度")
|
|
|
+
|
|
|
+ return merge_candidates(candidates)
|
|
|
+
|
|
|
+
|
|
|
+def merge_candidates(candidates):
|
|
|
+ candidates.sort(key=lambda r: (r["batch"], r["fault_hint"], r["start_time"]))
|
|
|
+ merged = []
|
|
|
+ for item in candidates:
|
|
|
+ if merged and item["batch"] == merged[-1]["batch"] and item["fault_hint"] == merged[-1]["fault_hint"]:
|
|
|
+ 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"]
|
|
|
+ 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))
|
|
|
+ return merged
|
|
|
+
|
|
|
+
|
|
|
+def write_csv(path: Path, rows, fields):
|
|
|
+ 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.writeheader()
|
|
|
+ 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, help="近期有效运行小时数,默认14天")
|
|
|
+ parser.add_argument("--limit-chunks", type=int, default=None, help="仅用于安全试跑")
|
|
|
+ parser.add_argument("--output-dir", default=str(Path(__file__).resolve().parents[1] / "cache" / "pks_scanner"))
|
|
|
+ args = parser.parse_args()
|
|
|
+ if args.chunk_days < 1 or args.baseline_hours < 12:
|
|
|
+ parser.error("chunk-days 必须>=1,baseline-hours 必须>=12")
|
|
|
+ start, end = parse_dt(args.start), parse_dt(args.end)
|
|
|
+ if end <= start:
|
|
|
+ parser.error("end 必须晚于 start")
|
|
|
+
|
|
|
+ print("只读模式:不修改任何数据库数据")
|
|
|
+ all_rows = []
|
|
|
+ for batch in args.batch:
|
|
|
+ all_rows.extend(fetch_hourly(batch, start, end, args.chunk_days, args.limit_chunks))
|
|
|
+ all_rows.sort(key=lambda r: (r["batch"], r["hour"]))
|
|
|
+ candidates = detect(all_rows, args.baseline_hours)
|
|
|
+ out = Path(args.output_dir)
|
|
|
+ write_csv(out / "hourly_summary.csv", all_rows, HOURLY_FIELDS)
|
|
|
+ write_csv(out / "fault_candidates.csv", candidates,
|
|
|
+ ["candidate_id", "unit", "batch", "fault_hint", "start_time", "end_time", "peak_time", "score", "coverage", "running_samples", "reason"])
|
|
|
+ print(f"小时摘要:{len(all_rows):,} 行")
|
|
|
+ print(f"候选段:{len(candidates):,} 段")
|
|
|
+ print(f"输出:{out / 'hourly_summary.csv'}")
|
|
|
+ print(f"输出:{out / 'fault_candidates.csv'}")
|
|
|
+
|
|
|
+
|
|
|
+if __name__ == "__main__":
|
|
|
+ main()
|