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