Files
pipeline-lifetime/app/prediction.py
T

579 lines
21 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
import logging
import os
import uuid
from dataclasses import dataclass
from pathlib import Path
from typing import Any
import joblib
import matplotlib
matplotlib.use("Agg")
import matplotlib.pyplot as plt
import numpy as np
import pandas as pd
from matplotlib import font_manager, rcParams
from werkzeug.datastructures import FileStorage
from werkzeug.utils import secure_filename
from .config import IMAGE_DIR, UPLOAD_DIR
from .time_utils import utc_now
CHINESE_FONT_CANDIDATES = [
"Noto Sans CJK SC",
"Noto Sans SC",
"Noto Sans CJK JP",
"Noto Sans CJK TC",
"Source Han Sans SC",
"WenQuanYi Micro Hei",
"SimHei",
"Microsoft YaHei",
"Arial Unicode MS",
]
CHINESE_FONT_FILES = [
"/usr/share/fonts/opentype/noto/NotoSansCJK-Regular.ttc",
"/usr/share/fonts/truetype/noto/NotoSansSC-Regular.ttf",
"/usr/share/fonts/truetype/noto/NotoSansCJK-Regular.ttc",
"/usr/local/share/fonts/NotoSansCJK-Regular.ttc",
"/mnt/c/Windows/Fonts/NotoSansSC-VF.ttf",
"/mnt/c/Windows/Fonts/msyh.ttc",
"/mnt/c/Windows/Fonts/simhei.ttf",
"/mnt/c/Windows/Fonts/simsun.ttc",
]
FEATURES = [
"管材",
"管径",
"流速",
"压力",
"温度",
"降雨量",
"位置",
"结构缺陷",
"功能缺陷",
]
ID_COLUMN = "管道编号"
PIPE_AGE_COLUMN = "管龄(年)"
LEGACY_PIPE_AGE_COLUMN = "管龄"
STATUS_COLUMN = "状态"
EVENT_AGE_COLUMN = "事件/观察管龄(年)"
COLUMN_ALIASES = {
"管径(mm": "管径",
"管径(mm)": "管径",
"流速(m/s": "流速",
"压力(MPa)": "压力",
"压力(MPa": "压力",
"温度(℃)": "温度",
"年均降雨量(mm": "降雨量",
"降雨量(mm": "降雨量",
}
INPUT_COLUMNS = [
ID_COLUMN,
PIPE_AGE_COLUMN,
STATUS_COLUMN,
EVENT_AGE_COLUMN,
*FEATURES,
]
CHART_DISPLAY_LIMIT = 10
MATERIAL_COLUMN = "管材"
MATERIAL_CODE_OPTIONS = [
(1, "镀锌"),
(2, "钢塑"),
(3, "铝塑"),
(4, "PPR"),
(5, "PE"),
(6, "UPVC"),
(7, "铸铁"),
(8, "预应力"),
(9, "自应力"),
(10, "玻璃钢夹砂"),
(11, "钢管"),
(12, "钢套混凝土管"),
(13, "球墨铸铁"),
(14, "其他"),
]
MATERIAL_ALIAS_TO_CODE = {
str(code): code
for code, _ in MATERIAL_CODE_OPTIONS
}
MATERIAL_ALIAS_TO_CODE.update(
{
name.casefold(): code
for code, name in MATERIAL_CODE_OPTIONS
}
)
MATERIAL_ALIAS_TO_CODE.update(
{
f"{code}-{name}".casefold(): code
for code, name in MATERIAL_CODE_OPTIONS
}
)
SUPPORTED_EXTENSIONS = {".csv", ".xls", ".xlsx"}
CHINESE_FONT_PROP = None
DEFECT_GRADE_VALUES = {"无": 0.0, "轻度": 1.0, "中度": 3.0, "严重": 5.0}
class PredictionError(Exception):
def __init__(self, message: str, status_code: int = 400) -> None:
super().__init__(message)
self.message = message
self.status_code = status_code
@dataclass
class PredictionArtifacts:
original_filename: str
saved_path: Path
excel_path: Path
image_path: Path
image_filename: str
sample_count: int
summary_rows: list[dict[str, Any]]
analysis_text: str
class ModelBundleAdapter:
"""Expose a bundled preprocessor and survival model as one predictor."""
def __init__(self, preprocessor, model) -> None:
self.preprocessor = preprocessor
self.model = model
def _transform(self, frame: pd.DataFrame):
values = frame[FEATURES].copy()
values[MATERIAL_COLUMN] = values[MATERIAL_COLUMN].map(
lambda value: f"M{int(value)}" if pd.notna(value) else value
)
values["位置"] = values["位置"].map(
lambda value: f"L{int(value)}" if pd.notna(value) else value
)
return self.preprocessor.transform(values.to_numpy(dtype=object)).astype(np.float32)
def predict_survival_function(self, frame: pd.DataFrame):
return self.model.predict_survival_function(self._transform(frame))
def predict(self, frame: pd.DataFrame):
return self.model.predict(self._transform(frame))
def configure_matplotlib_fonts():
for font_path in CHINESE_FONT_FILES:
path = Path(font_path)
if path.exists():
font_manager.fontManager.addfont(str(path))
prop = font_manager.FontProperties(fname=str(path))
rcParams["font.family"] = [prop.get_name(), "sans-serif"]
rcParams["font.sans-serif"] = [prop.get_name(), *CHINESE_FONT_CANDIDATES, "DejaVu Sans"]
rcParams["axes.unicode_minus"] = False
return prop
available_fonts = {font.name for font in font_manager.fontManager.ttflist}
selected_fonts = [font for font in CHINESE_FONT_CANDIDATES if font in available_fonts]
if selected_fonts:
rcParams["font.family"] = [selected_fonts[0], "sans-serif"]
rcParams["font.sans-serif"] = selected_fonts + ["DejaVu Sans"]
rcParams["axes.unicode_minus"] = False
if selected_fonts:
return font_manager.FontProperties(family=selected_fonts[0])
logging.warning("未找到中文字体,生成的图表中文可能无法显示。")
return None
CHINESE_FONT_PROP = configure_matplotlib_fonts()
def chinese_font_kwargs() -> dict[str, Any]:
return {"fontproperties": CHINESE_FONT_PROP} if CHINESE_FONT_PROP else {}
def load_model(model_path: str):
if not os.path.exists(model_path):
raise FileNotFoundError(f"未找到模型文件: {model_path}")
model = joblib.load(model_path)
if isinstance(model, dict):
if not {"preprocessor", "model"}.issubset(model):
raise ValueError("模型包缺少 preprocessor 或 model。")
model = ModelBundleAdapter(model["preprocessor"], model["model"])
feature_names = getattr(model, "feature_names_in_", None)
if feature_names is not None and list(feature_names) != FEATURES:
raise ValueError(
"模型输入特征与系统配置不一致: "
f"model={list(feature_names)}, app={FEATURES}"
)
n_features = getattr(model, "n_features_in_", None)
if n_features is not None and int(n_features) != len(FEATURES):
raise ValueError(f"模型特征数量不一致: model={n_features}, app={len(FEATURES)}")
return model
def safe_unlink(path: Path) -> None:
try:
path.unlink(missing_ok=True)
except OSError as exc:
logging.warning("删除文件失败 %s: %s", path, exc)
def secure_upload_name(original_filename: str, run_id: str) -> tuple[str, str]:
suffix = Path(original_filename).suffix.lower()
if suffix not in SUPPORTED_EXTENSIONS:
raise PredictionError("不支持的格式,仅支持逗号分隔值文件或电子表格文件")
safe_full_name = secure_filename(original_filename)
safe_stem = Path(safe_full_name).stem if safe_full_name else ""
if not safe_stem:
safe_stem = "upload"
return f"{safe_stem}_{run_id}{suffix}", suffix
def read_input_file(path: Path, suffix: str) -> pd.DataFrame:
try:
if suffix == ".csv":
df = pd.read_csv(path, dtype={ID_COLUMN: "string"})
else:
try:
df = pd.read_excel(path, sheet_name="Template", dtype={ID_COLUMN: "string"})
except ValueError:
df = pd.read_excel(path, dtype={ID_COLUMN: "string"})
df["_excel_row"] = df.index + 2
except Exception as exc:
logging.exception("文件解析失败: %s", exc)
raise PredictionError("文件解析失败,请检查编码或表格格式。")
df = normalize_input_columns(df)
if suffix in {".xls", ".xlsx"}:
df = fill_defects_from_detail_sheet(path, df)
return df
def normalize_input_columns(df: pd.DataFrame) -> pd.DataFrame:
aliases = dict(COLUMN_ALIASES)
if PIPE_AGE_COLUMN not in df.columns and LEGACY_PIPE_AGE_COLUMN in df.columns:
aliases[LEGACY_PIPE_AGE_COLUMN] = PIPE_AGE_COLUMN
return df.rename(columns={name: aliases[name] for name in df.columns if name in aliases})
def normalize_id_value(value: Any, fallback: str) -> str:
if pd.isna(value):
return fallback
if isinstance(value, str):
return value.strip()
return str(value)
def defect_grade_value(value: Any) -> float | None:
if pd.isna(value):
return None
text = str(value).strip()
if text == "":
return None
if text in DEFECT_GRADE_VALUES:
return DEFECT_GRADE_VALUES[text]
try:
return float(text)
except ValueError:
return None
def weighted_defect_score(values: list[Any], weights: list[float]) -> float | None:
scores = [defect_grade_value(value) for value in values]
if any(score is None for score in scores):
return None
return round(sum(score * weight for score, weight in zip(scores, weights)), 3)
def numeric_or_none(value: Any) -> float | None:
try:
if pd.isna(value):
return None
return float(value)
except (TypeError, ValueError):
return None
def fill_defects_from_detail_sheet(path: Path, df: pd.DataFrame) -> pd.DataFrame:
if "结构缺陷" not in df.columns or "功能缺陷" not in df.columns or "_excel_row" not in df.columns:
return df
structure = pd.to_numeric(df["结构缺陷"], errors="coerce")
function = pd.to_numeric(df["功能缺陷"], errors="coerce")
needs_fill = structure.isna() | function.isna()
df["结构缺陷"] = structure
df["功能缺陷"] = function
if not needs_fill.any():
return df
try:
from openpyxl import load_workbook
workbook = load_workbook(path, data_only=False, read_only=True)
defect_sheet = workbook["缺陷计算"]
except Exception as exc:
logging.info("未能读取缺陷计算工作表,跳过缺陷自动计算: %s", exc)
return df
for index in df.index[needs_fill]:
excel_row = int(df.at[index, "_excel_row"])
if pd.isna(df.at[index, "结构缺陷"]):
value = numeric_or_none(defect_sheet.cell(excel_row, 5).value)
if value is None:
value = weighted_defect_score(
[defect_sheet.cell(excel_row, column).value for column in (2, 3, 4)],
[0.5, 0.4, 0.1],
)
df.at[index, "结构缺陷"] = value
if pd.isna(df.at[index, "功能缺陷"]):
value = numeric_or_none(defect_sheet.cell(excel_row, 10).value)
if value is None:
value = weighted_defect_score(
[defect_sheet.cell(excel_row, column).value for column in (6, 7, 8, 9)],
[0.5, 0.2, 0.2, 0.1],
)
df.at[index, "功能缺陷"] = value
return df
def validate_input_frame(df: pd.DataFrame) -> None:
if df.empty:
raise PredictionError("上传文件没有可预测的数据。")
missing = [col for col in INPUT_COLUMNS if col not in df.columns]
if missing:
raise PredictionError(f"缺少必要字段: {', '.join(missing)}")
def normalize_material_code(value: Any) -> int:
if pd.isna(value):
raise PredictionError("管材不能为空。")
if isinstance(value, str):
key = value.strip().casefold()
else:
numeric_value = float(value)
if not numeric_value.is_integer():
raise PredictionError(f"管材编码无效: {value}")
key = str(int(numeric_value))
code = MATERIAL_ALIAS_TO_CODE.get(key)
if code is None:
raise PredictionError(f"管材编码无效: {value}")
return code
def prepare_model_features(df: pd.DataFrame) -> pd.DataFrame:
x_test = df[FEATURES].copy()
x_test[MATERIAL_COLUMN] = x_test[MATERIAL_COLUMN].map(normalize_material_code)
return x_test
def grade_info(probability: float) -> tuple[str, str, str]:
if probability <= 0.2:
return ("I级", "管道安全风险十分严重,需立刻进行抢修或更新改造", "bg-dangerSoft text-dangerText")
if probability <= 0.4:
return ("II级", "管道安全风险较为严重,需尽快安排检修及加频巡检", "bg-orange-50 text-orange-600")
if probability <= 0.6:
return ("III级", "管道安全风险较低,需安排定期巡检", "bg-amber-50 text-amber-600")
if probability <= 0.8:
return ("IV级", "管道安全风险较小,维持常规巡视", "bg-blue-50 text-blue-600")
return ("V级", "管道安全,维持常规巡视", "bg-blueSoft text-primary")
def interpolate_probability(times: list[float], probs: list[float], target: float) -> float:
if not times:
return 0.0
if target <= times[0]:
return float(probs[0])
for idx in range(1, len(times)):
if times[idx] >= target:
return float(probs[idx])
return float(probs[-1])
def estimate_remaining_life(times: list[float], probs: list[float]) -> float:
for t, p in zip(times, probs):
if p <= 0.5:
return float(t)
return float(times[-1]) if times else 0.0
def make_analysis_text(summary_rows: list[dict[str, Any]]) -> str:
if not summary_rows:
return "当前结果为空,暂无可供解释的样本。"
worst = min(summary_rows, key=lambda x: x["health_probability"])
best = max(summary_rows, key=lambda x: x["health_probability"])
return (
"阶梯状曲线表示模型对不同管道随时间推移维持在安全健康状态概率的动态预测。"
f"当前样本中风险最高管道为 {worst['pipe_id']}{worst['grade_label']}),"
f"健康概率最高管道为 {best['pipe_id']}{best['health_probability']:.1%})。"
)
def run_prediction(uploaded: FileStorage, user_id: int, model) -> PredictionArtifacts:
original_filename = uploaded.filename or ""
if not original_filename:
raise PredictionError("未选择文件")
timestamp = utc_now().strftime("%Y%m%d%H%M%S")
run_id = f"{timestamp}_{uuid.uuid4().hex[:8]}"
saved_filename, suffix = secure_upload_name(original_filename, run_id)
user_dir = UPLOAD_DIR / f"user_{user_id}"
user_dir.mkdir(parents=True, exist_ok=True)
original_path = user_dir / saved_filename
uploaded.save(original_path)
try:
df = read_input_file(original_path, suffix)
validate_input_frame(df)
x_test = prepare_model_features(df)
try:
curves = model.predict_survival_function(x_test)
except Exception as exc:
logging.exception("预测失败: %s", exc)
raise PredictionError("模型预测失败,请检查输入字段类型是否正确。", 500)
image_filename = f"plot_{user_id}_{run_id}.png"
image_path = IMAGE_DIR / image_filename
summary_rows, summary_sheet_rows = render_survival_chart(df, curves, image_path)
safe_stem = Path(saved_filename).stem
excel_path = user_dir / f"{safe_stem}_pre.xlsx"
write_prediction_workbook(excel_path, curves, summary_rows, summary_sheet_rows)
return PredictionArtifacts(
original_filename=original_filename,
saved_path=original_path,
excel_path=excel_path,
image_path=image_path,
image_filename=image_filename,
sample_count=len(summary_rows),
summary_rows=summary_rows,
analysis_text=make_analysis_text(summary_rows),
)
except PredictionError:
safe_unlink(original_path)
raise
def render_survival_chart(df: pd.DataFrame, curves, image_path: Path) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]:
plt.figure(figsize=(10, 5.6))
summary_rows: list[dict[str, Any]] = []
summary_sheet_rows: list[dict[str, Any]] = []
display_note = f"说明:图中仅展示前{CHART_DISPLAY_LIMIT}条管道的示例数据,不足{CHART_DISPLAY_LIMIT}条则全部展示;完整结果请下载电子表格。"
for i, curve in enumerate(curves):
times = [float(x) for x in list(curve.x)]
probs = [float(y) for y in list(curve.y)]
pipe_id = normalize_id_value(df.iloc[i][ID_COLUMN], f"Pipe_{i+1:03d}")
raw_pipe_age = df.iloc[i][PIPE_AGE_COLUMN] if PIPE_AGE_COLUMN in df.columns else None
pipe_age_value = float(raw_pipe_age) if pd.notna(raw_pipe_age) else 0.0
pipe_age = f"{raw_pipe_age} 年" if pd.notna(raw_pipe_age) else "-"
health_probability = interpolate_probability(times, probs, pipe_age_value)
health_risk = 1.0 - health_probability
remaining_life = estimate_remaining_life(times, probs)
grade_label, grade_desc, grade_class = grade_info(health_probability)
summary_rows.append(
{
"pipe_id": pipe_id,
"pipe_age": pipe_age,
"health_probability": health_probability,
"health_risk": health_risk,
"remaining_life": remaining_life,
"grade_label": grade_label,
"grade_desc": grade_desc,
"grade_class": grade_class,
}
)
summary_sheet_rows.append(
{
ID_COLUMN: pipe_id,
PIPE_AGE_COLUMN: pipe_age,
"健康风险值": health_risk,
}
)
if i < CHART_DISPLAY_LIMIT:
plt.step(times, probs, where="post", linewidth=2, label=pipe_id)
font_kwargs = chinese_font_kwargs()
plt.xlabel("管龄(年)", **font_kwargs)
plt.ylabel("健康风险", **font_kwargs)
plt.title("管道的剩余寿命分析图", **font_kwargs)
plt.figtext(0.5, 0.02, display_note, ha="center", fontsize=9, color="#475569", **font_kwargs)
plt.grid(alpha=0.18)
if min(len(summary_rows), CHART_DISPLAY_LIMIT) <= 12:
plt.legend(loc="best", fontsize=8, prop=CHINESE_FONT_PROP)
ax = plt.gca()
if CHINESE_FONT_PROP:
for label in [*ax.get_xticklabels(), *ax.get_yticklabels()]:
label.set_fontproperties(CHINESE_FONT_PROP)
plt.tight_layout(rect=(0, 0.07, 1, 1))
plt.savefig(image_path, dpi=160, bbox_inches="tight")
plt.close()
return summary_rows, summary_sheet_rows
def write_prediction_workbook(
excel_path: Path,
curves,
summary_rows: list[dict[str, Any]],
summary_sheet_rows: list[dict[str, Any]],
) -> None:
chart_times: set[float] = set()
chart_series: list[tuple[str, dict[float, float]]] = []
for i, curve in enumerate(curves):
times = [float(x) for x in list(curve.x)]
probs = [float(y) for y in list(curve.y)]
pipe_id = summary_rows[i]["pipe_id"]
chart_times.update(times)
chart_series.append((pipe_id, dict(zip(times, probs))))
sorted_times = sorted(chart_times)
sample_columns = ["管龄(年)", *[pipe_id for pipe_id, _ in chart_series]]
sample_data_rows = [
[time, *[series.get(time) for _, series in chart_series]]
for time in sorted_times
]
with pd.ExcelWriter(excel_path, engine="openpyxl") as writer:
summary_df = pd.DataFrame(summary_sheet_rows, columns=[ID_COLUMN, PIPE_AGE_COLUMN, "健康风险值"])
sample_df = pd.DataFrame(sample_data_rows, columns=sample_columns)
if ID_COLUMN in summary_df.columns:
summary_df[ID_COLUMN] = summary_df[ID_COLUMN].astype("string")
summary_df.to_excel(writer, sheet_name="结果摘要", index=False)
sample_df.to_excel(writer, sheet_name="样本数据", index=False)
sample_worksheet = writer.book["样本数据"]
summary_worksheet = writer.book["结果摘要"]
for cell in summary_worksheet["C"][1:]:
cell.number_format = "0.0%"
sample_worksheet.freeze_panes = "A2"
sample_worksheet.auto_filter.ref = sample_worksheet.dimensions
sample_worksheet.column_dimensions["A"].width = 12
for column_cells in sample_worksheet.iter_cols(min_col=2, max_col=sample_worksheet.max_column):
header_cell = column_cells[0]
sample_worksheet.column_dimensions[header_cell.column_letter].width = max(12, len(str(header_cell.value)) + 2)
for cell in column_cells[1:]:
cell.number_format = "0.0%"
for sheet_name in ("结果摘要", "样本数据"):
worksheet = writer.book[sheet_name]
header_cells = next(worksheet.iter_rows(min_row=1, max_row=1), [])
id_column_index = None
for cell in header_cells:
if cell.value == ID_COLUMN:
id_column_index = cell.column
break
if id_column_index is None:
continue
for cell in worksheet.iter_cols(
min_col=id_column_index,
max_col=id_column_index,
min_row=2,
max_row=worksheet.max_row,
):
for item in cell:
if item.value is not None:
item.value = str(item.value)
item.number_format = "@"