项目背景与目标
一个文档处理服务每天执行多批任务,运行记录包含执行时间、工作节点、输入文件数、输入体积、耗时、状态和重试次数。维护者需要回答:整体成功率是多少,哪些节点处理较慢,输入体积与耗时是否相关,失败是否集中在某些日期。
输入是 CSV 运行记录,输出包括数据质量报告、节点汇总、每日汇总、相关系数和三张图。结果应满足:项目能够自行生成样例数据;重复记录、非法时间和缺失耗时可追踪;异常耗时不会被静默删除;结论明确样本边界。
这个项目的数据、字段、名称和业务设定均为自创,只用于演示完整分析流程,不代表任何真实系统。

示例结果把每日成功率、节点耗时分布和输入体积与耗时的关系放在同一张分析画布中。三类图分别回答稳定性、节点差异和变量关系问题,不能相互替代。
整体架构
生成或接收运行记录
│
▼
字段、类型与质量检查 ──► quality_report.csv
│
▼
去重、缺失处理、异常标记
│
├──► worker_summary.csv
├──► daily_summary.csv
├──► analysis_summary.txt
└──► figures/*.png
建议目录:
batch-run-analysis/
├── data/
├── output/
│ └── figures/
├── prepare_sample.py
├── analyze.py
└── requirements.txt
核心流程
- 检查必需列,固定标识符和类别字段类型。
- 记录原始行数、重复运行 ID、非法时间和缺失耗时。
- 去除重复 ID;无法确定时间、耗时和输入体积的记录不进入指标计算。
- 检查数值范围,将 IQR 范围外耗时标记为候选异常值。
- 计算成功率、节点吞吐、每日失败量以及输入体积与耗时的相关系数。
- 输出表格和图表,再把事实、推断和局限分开描述。
数据字典与分析口径
| 字段 | 含义 | 约束 | 是否参与指标 |
|---|---|---|---|
run_id | 单次运行标识 | 非空、唯一 | 用于计数和去重 |
started_at | 开始时间 | 可解析时间 | 用于每日汇总 |
worker | 工作节点 | 非空类别 | 用于节点比较 |
file_count | 文件数量 | 非负整数 | 描述任务规模 |
input_mb | 输入体积 | 非负数 | 与耗时相关分析 |
duration_sec | 执行耗时 | 非负数 | 核心性能指标 |
status | 运行结果 | success/failed | 计算成功率 |
retry_count | 重试次数 | 非负整数 | 辅助诊断 |
异常耗时只是不符合当前总体分布的候选点,可能是运行事故,也可能是真实大任务。本项目保留原记录并增加 duration_outlier,节点常规耗时比较暂时排除它;质量报告仍记录完整输入问题。
成功率定义为 成功运行数 / 有效运行数。如果无效记录与失败高度相关,直接删除会高估成功率,因此正式系统还应单独报告“无法评估率”。样例项目规模很小,统计结果只用于验证流程。
关键实现
依赖文件:
pandas
numpy
matplotlib
seaborn
生成包含数据质量问题的样例:
# prepare_sample.py
from pathlib import Path
import pandas as pd
rows = [
("run-001", "2026-07-01 09:00", "worker-a", 12, 18.0, 42.0, "success", 0),
("run-002", "2026-07-01 10:10", "worker-b", 20, 35.0, 70.0, "success", 0),
("run-003", "2026-07-01 11:20", "worker-c", 8, 10.0, 31.0, "failed", 2),
("run-004", "2026-07-02 09:15", "worker-a", 30, 52.0, 96.0, "success", 1),
("run-005", "2026-07-02 10:30", "worker-b", 16, 22.0, 55.0, "success", 0),
("run-006", "2026-07-02 13:00", "worker-c", 25, 41.0, 88.0, "failed", 3),
("run-007", "2026-07-03 08:50", "worker-a", 10, 14.0, 36.0, "success", 0),
("run-008", "2026-07-03 10:40", "worker-b", 35, 61.0, 112.0, "success", 1),
("run-009", "2026-07-03 15:20", "worker-c", 18, 29.0, 67.0, "success", 0),
("run-010", "2026-07-04 09:05", "worker-a", 40, 74.0, 138.0, "success", 1),
("run-010", "2026-07-04 09:05", "worker-a", 40, 74.0, 138.0, "success", 1),
("run-011", "invalid-time", "worker-b", 14, 19.0, 48.0, "failed", 2),
("run-012", "2026-07-04 12:10", "worker-c", 22, 38.0, None, "success", 0),
("run-013", "2026-07-04 16:30", "worker-b", 60, 90.0, 620.0, "success", 2),
]
columns = [
"run_id", "started_at", "worker", "file_count", "input_mb",
"duration_sec", "status", "retry_count",
]
output = Path("data/batch_runs.csv")
output.parent.mkdir(parents=True, exist_ok=True)
pd.DataFrame(rows, columns=columns).to_csv(output, index=False)
print(f"created {output} with {len(rows)} rows")
分析程序把数据质量、指标和绘图拆开:
# analyze.py
from pathlib import Path
import matplotlib
matplotlib.use("Agg")
import matplotlib.pyplot as plt
import pandas as pd
import seaborn as sns
DATA_PATH = Path("data/batch_runs.csv")
OUTPUT_DIR = Path("output")
FIGURE_DIR = OUTPUT_DIR / "figures"
REQUIRED_COLUMNS = {
"run_id", "started_at", "worker", "file_count", "input_mb",
"duration_sec", "status", "retry_count",
}
def load_data(path: Path) -> pd.DataFrame:
data = pd.read_csv(path, dtype={"run_id": "string", "worker": "string"})
missing = REQUIRED_COLUMNS - set(data.columns)
if missing:
raise ValueError(f"缺少必需列: {', '.join(sorted(missing))}")
return data
def clean_data(raw: pd.DataFrame) -> tuple[pd.DataFrame, pd.DataFrame]:
data = raw.copy()
data["started_at"] = pd.to_datetime(data["started_at"], errors="coerce")
quality = pd.DataFrame(
{
"metric": ["raw_rows", "duplicate_ids", "invalid_times", "missing_duration"],
"value": [
len(data),
int(data.duplicated(subset="run_id").sum()),
int(data["started_at"].isna().sum()),
int(data["duration_sec"].isna().sum()),
],
}
)
data = data.drop_duplicates(subset="run_id", keep="first")
data = data.dropna(subset=["started_at", "input_mb", "duration_sec"]).copy()
valid_status = {"success", "failed"}
if not set(data["status"]).issubset(valid_status):
raise ValueError("status 只能是 success 或 failed")
if (data[["file_count", "input_mb", "duration_sec", "retry_count"]] < 0).any().any():
raise ValueError("数值指标不能为负数")
q1 = data["duration_sec"].quantile(0.25)
q3 = data["duration_sec"].quantile(0.75)
iqr = q3 - q1
data["duration_outlier"] = ~data["duration_sec"].between(
q1 - 1.5 * iqr, q3 + 1.5 * iqr
)
data["day"] = data["started_at"].dt.date
data["is_success"] = data["status"].eq("success")
return data, quality
def summarize(data: pd.DataFrame):
regular = data.loc[~data["duration_outlier"]]
worker = (
regular.groupby("worker", as_index=False)
.agg(
run_count=("run_id", "count"),
success_rate=("is_success", "mean"),
median_duration=("duration_sec", "median"),
)
.sort_values(["success_rate", "median_duration"], ascending=[False, True])
)
daily = (
data.groupby("day", as_index=False)
.agg(run_count=("run_id", "count"), failed_count=("is_success", lambda x: (~x).sum()))
)
correlation = float(regular["input_mb"].corr(regular["duration_sec"]))
return worker, daily, correlation
def draw_figures(data: pd.DataFrame, worker: pd.DataFrame, daily: pd.DataFrame) -> None:
FIGURE_DIR.mkdir(parents=True, exist_ok=True)
sns.set_theme(style="whitegrid")
regular = data.loc[~data["duration_outlier"]]
figure, axis = plt.subplots(figsize=(7, 4))
sns.histplot(data=data, x="duration_sec", bins=8, ax=axis)
axis.set(title="Duration distribution", xlabel="Seconds", ylabel="Count")
figure.tight_layout(); figure.savefig(FIGURE_DIR / "duration_distribution.png", dpi=150)
plt.close(figure)
figure, axis = plt.subplots(figsize=(7, 4))
sns.scatterplot(data=regular, x="input_mb", y="duration_sec", hue="worker", ax=axis)
axis.set(title="Input size and duration", xlabel="Input MB", ylabel="Seconds")
figure.tight_layout(); figure.savefig(FIGURE_DIR / "size_duration.png", dpi=150)
plt.close(figure)
figure, axis = plt.subplots(figsize=(7, 4))
sns.barplot(data=worker, x="worker", y="median_duration", ax=axis)
axis.set(title="Median duration by worker", xlabel="Worker", ylabel="Seconds")
figure.tight_layout(); figure.savefig(FIGURE_DIR / "worker_duration.png", dpi=150)
plt.close(figure)
def main() -> None:
OUTPUT_DIR.mkdir(exist_ok=True)
raw = load_data(DATA_PATH)
clean, quality = clean_data(raw)
worker, daily, correlation = summarize(clean)
quality.to_csv(OUTPUT_DIR / "quality_report.csv", index=False)
worker.to_csv(OUTPUT_DIR / "worker_summary.csv", index=False)
daily.to_csv(OUTPUT_DIR / "daily_summary.csv", index=False)
draw_figures(clean, worker, daily)
summary = (
f"raw_rows={len(raw)}\nanalysis_rows={len(clean)}\n"
f"duration_outliers={int(clean['duration_outlier'].sum())}\n"
f"input_duration_correlation={correlation:.3f}\n"
)
(OUTPUT_DIR / "analysis_summary.txt").write_text(summary, encoding="utf-8")
print(summary)
print(worker.to_string(index=False))
if __name__ == "__main__":
main()
结果怎样解读
节点汇总至少同时查看运行数、成功率和中位耗时。运行数很少的节点即使成功率为 100%,也不能说明更稳定;中位数降低极端值影响,但仍需查看高分位耗时。输入体积与耗时的相关系数接近 1,只表示样例中线性共同变化明显,不排除文件类型、并发和缓存等混杂因素。
每日失败数适合发现集中故障,但不同日期任务总数不同时还应计算失败率。图表与汇总表必须使用一致的数据口径,否则“图上较慢、表中较快”可能只是一个排除了异常值、另一个没有排除。
为清洗规则补测试
核心清洗函数可以用最小 DataFrame 验证,不依赖完整 CSV:
import pandas as pd
def test_clean_data_reports_invalid_rows() -> None:
raw = pd.DataFrame([
{
"run_id": "a", "started_at": "2026-01-01", "worker": "w1",
"file_count": 1, "input_mb": 2.0, "duration_sec": 3.0,
"status": "success", "retry_count": 0,
},
{
"run_id": "b", "started_at": "bad", "worker": "w1",
"file_count": 1, "input_mb": 2.0, "duration_sec": 3.0,
"status": "failed", "retry_count": 1,
},
])
clean, quality = clean_data(raw)
invalid = quality.loc[quality["metric"] == "invalid_times", "value"].item()
assert invalid == 1
assert len(clean) == 1
实际项目还应覆盖重复 ID、负值、未知状态、全缺失和没有异常值的情况。测试固定的是业务规则,防止后续重构悄悄改变统计口径。
运行与结果检查
python3 -m venv .venv
source .venv/bin/activate
python -m pip install -r requirements.txt
python prepare_sample.py
python analyze.py
find output -maxdepth 2 -type f -print
验证时应识别一个重复 ID、一个非法时间、一个缺失耗时和一个候选异常耗时,并生成三张汇总文件、分析摘要及三张 PNG 图。原始 CSV 不被覆盖。
结果、评估与阶段复盘
样例中输入体积与运行耗时通常呈正相关,不同工作节点的中位耗时和成功率也存在差异。这些结果只能说明当前构造样本的表现,不能直接证明节点性能问题:任务难度、并发负载和缓存状态都没有记录。
相比直接删除异常行,先标记异常耗时可以同时保留运行事故线索和常规吞吐分析口径。分析可信度来自规则透明、统计口径一致和输出可复查,而不是图表数量。
常见问题与排查
- 找不到 CSV:从项目根目录运行,并打印
DATA_PATH.resolve()。 - 清洗后行数异常:逐条核对去重、时间解析和必需字段过滤数量。
- 相关系数为
NaN:有效样本太少,或某列没有变化。 - 成功率显示为小数:展示时转换为百分比,计算时保留原始比例。
- 节点比较不公平:先检查各节点输入规模和任务类型是否可比。
- 图表与表格口径不一致:明确哪些结果排除了异常耗时。
改进方向
- 增加任务类型、并发数、队列等待和错误类别字段。
- 将清洗规则写成测试,覆盖负值、未知状态和重复 ID。
- 增加分位数耗时,避免只看平均值。
- 数据量增大后按日期分区读取,控制内存占用。
- 接入真实运行日志前完成脱敏和字段授权检查。
- 将质量报告、图表和结论生成时间写入运行清单,形成可追溯分析版本。
许可协议:CC BY-NC 4.0
更新于 1 小时前
觉得文章有帮助?点个赞吧!
0 条评论


