提交 ef5bdd1d authored 作者: 陈泽健's avatar 陈泽健

fix(performance): 启动回填 ResourceClosedError 修复(双 scalar 调用)+ 补提交拆表测试套件

- _migrate_legacy_executions() 中 existing.scalar() 双调用:Result.scalar() 取完首行即关闭 result,二次调用抛 ResourceClosedError——legacy 迁移一旦完成,此后每次启动回填必炸,导致拆表从表存量回填永远不执行;改为单次取值 legacy_count
- 补提交 test_performance_satellite_tables.py 14 用例(ea2a90cb 提交信息已声明但文件漏入库)
- HANDOFF_性能测试.md 记录拆表部署 5.60 闭环:12 从表建成 + 存量回填 101 行 + JSON 'null' 字面量校验口径 + 部署/验证脚本清单
Co-Authored-By: 's avatarClaude <noreply@anthropic.com>
上级 f60a2914
......@@ -2,13 +2,69 @@
> **生成时间**: 2026-09-08
> **当前分支**: `platform-auto-test`
> **最近提交**: `b21f2fe9` fix(performance): 全库扫描 order_by 兜底排查——26 处两段式改造防 MySQL 1038(已推送 origin/platform-auto-test)
> **会话窗口**: 性能测试 — 2026-09-08 P3 待办闭环:performance 大 JSON/TEXT 列垂直拆表(7 个从表 + 双写 + 读优先从表 + 启动回填 + 零回归)
> **状态**: ✅ 代码开发与测试完成(3 文件修改 + 1 新增测试,全量 427 passed 零回归),待部署 5.60 与提交推送
> **最近提交**: `ea2a90cb` perf(performance): 大JSON列垂直拆表闭环 MySQL 1038 + 登录汇总透出(已推送 origin/platform-auto-test)
> **会话窗口**: 性能测试 — 2026-09-08 P3 待办闭环:performance 大 JSON/TEXT 列垂直拆表(7 个从表 + 双写 + 读优先从表 + 启动回填 + 零回归)+ **部署 5.60 验证闭环**
> **状态**: ✅ 已提交推送 `ea2a90cb` + 已部署 5.60(12 张从表建成 + 存量回填 101 行 + 修复启动回填 ResourceClosedError,本地全量 437 passed)
---
## ⚡ 最新会话更新(2026-09-08)— P3 待办闭环:performance 大 JSON 列垂直拆表 ✅
## ⚡ 最新会话更新(2026-09-08 续)— 拆表部署 5.60 验证闭环 + 启动回填 ResourceClosedError 修复 ✅
### A. 部署状态核验(拆表代码此前未部署)
- 提交 `ea2a90cb` 推送后**未走部署步骤**,远端 5.60 五个关键文件 md5 全部匹配 `b21f2fe9`(CRLF 变体),MySQL 无任何从表;
- 本次补部署 5 个文件(`database.py` / `models/performance.py` / `services/performance_service.py` / `routers/performance.py` / `schemas/performance.py`),远端 .bak 备份 + 原子替换(.new → mv)+ md5 逐文件复核一致。
### B. 意外发现并修复:启动回填 ResourceClosedError(`_migrate_legacy_executions`)
首次部署重启后从表建成但回填 0 行,启动日志:`存量执行数据回填失败(不影响启动,下次重试): This result object is closed.`
| 项 | 说明 |
|----|------|
| 根因 | `performance_service.py` `_migrate_legacy_executions()``if existing.scalar() and existing.scalar() > 0:`**同一 Result 调用两次 `scalar()`**——SQLAlchemy 中 `scalar()`(内部 `first()`)取完第一行即关闭 result,第二次调用抛 `ResourceClosedError` |
| 为何早不炸 | 5.60 首次迁移前无 legacy 记录 → `scalar()` 返回 0 falsy → `and` 短路不调用第二次;**迁移一旦完成,之后每次启动必炸**,导致后续拆表回填永远不执行 |
| 修复 | 改为 `legacy_count = (...).scalar() or 0; if legacy_count > 0:`(单次取值);全库扫描确认无其他同类双调用 |
| 验证 | 本地最小复现(双 scalar 抛错)+ 修复后 init_db 正常;5.60 重启日志 `存量回填已执行过,跳过` + `大字段从表存量回填: 101 行` ✅ |
### C. 部署验证结果 ✅
| 项 | 结果 |
|----|------|
| 远端 5 文件 md5 | ✅ 全部 = HEAD `ea2a90cb` 工作区(含 scalar 修复) |
| 容器 | ✅ `docker compose restart app` → Up (healthy),`/health` 200 第 1 次探测 |
| 代码标记 | ✅ 容器内 `_backfill_satellite_tables` / `_satellite_upsert` / `_TASK_SATELLITES` 全命中 |
| 12 张从表建表 | ✅ MySQL 全部创建(create_all) |
| 存量回填 | ✅ 101 行(任务 13 + 执行 79 + 合并报告 9),12/12 组「主表有效非空行数 = 从表行数」逐一吻合 |
| API 冒烟 | ✅ `/api/performance/tasks``/api/performance/executions` 均 200 |
| 本地全量测试 | ✅ **437 passed**(含拆表套件 14 用例),零回归 |
| 启动日志 | ✅ 无 performance 相关 error/traceback |
> 注:回填校验需排除 MySQL JSON 字面量 `'null'`(`JSON_TYPE=NULL`)行——主表部分历史执行的大列存的是 JSON null(无有效数据),从表无行是正确行为,读路径从表缺行回退主表列语义一致。
### D. 剩余待办(更新)
| # | 任务 | 优先级 | 说明 |
|---|------|--------|------|
| 1 | **scalar 修复提交推送** | P1 | `performance_service.py` `_migrate_legacy_executions()` 双 scalar 修复(本节 B)+ 本 HANDOFF 更新,待 /GitCommit |
| 2 | 跟踪其他窗口 progress | — | recorder/projects 为其他窗口工作,合并时注意冲突 |
### E. 本次会话工具与脚本(`backend/tmp/`,gitignored)
| 脚本 | 用途 |
|------|------|
| `verify_satellite_deploy_560.py` | 部署核验:远端 5 文件 md5 + 容器健康 + 从表清单 + 启动日志(含双 scalar 修复前发现回填失败的过程) |
| `deploy_satellite_560.py` | 部署执行:5 文件 .bak 备份 → .new 原子上传 → mv → md5 复核 → `docker compose restart app` → 健康探测 → 容器内 grep 代码标记 |
| `verify_backfill_560.py` / `verify_backfill_final_560.py` | 回填对比:12 组「主表非空 vs 从表行数」(final 版排除 JSON 字面量 `'null'`,结论全部 OK) |
| `diff_missing_backfill_560.py` | 定位缺行明细:LEFT JOIN 找主表非空但从表缺行的 execution(确认全部为 `JSON_TYPE=NULL``'null'` 字面量行) |
| `check_tables_and_running_560.py` | 从表存在性 + running/pending 执行检查(部署前确认重启无中断风险) |
**远端备份**:5.60 `/data/third_party/plat-auto-test/backend/app/` 下旧文件备份为 `*.bak_20260908_*`(部署脚本自动创建),确认稳定后可清理。
**md5 对比方法论**:Windows 工作区文件为 CRLF,`git show` 输出为 LF,与 Linux 远端直接 md5 比对会全部误报不一致——需按「git blob 内容 × LF/CRLF 变体」逐一匹配远端 md5(本次远端 5 文件全部命中 `b21f2fe9` 的 CRLF 变体,由此断定未部署)。部署脚本上传本地文件原样字节(CRLF),Python 运行不受换行符影响,与历史部署习惯一致。
---
## ⚡ 会话更新(2026-09-08)— P3 待办闭环:performance 大 JSON 列垂直拆表 ✅
### A. 背景与目标
......@@ -65,8 +121,8 @@ MySQL 1038(`Out of sort memory`)根因是整行加载大 JSON/TEXT 列消耗
| # | 任务 | 优先级 | 说明 |
|---|------|--------|------|
| 1 | 部署 5.60 | P2 | 将 `database.py` / `models/performance.py` / `services/performance_service.py` 部署至 5.60 容器,并重启容器生效回填 |
| 2 | 代码提交推送 | P2 | /GitCommit 规范提交并推送到 origin/platform-auto-test |
| 1 | ~~部署 5.60~~ | ~~P2~~ | ✅ 已部署并验证(12 从表 + 101 行回填,见顶部 2026-09-08 续节 C) |
| 2 | ~~代码提交推送~~ | ~~P2~~ | ✅ 已完成(`ea2a90cb` 已推送;后续 scalar 修复待提交,见顶部待办) |
### G. 多窗口并行开发注意(提交时必读)
......
......@@ -1258,12 +1258,14 @@ class PerformanceService:
Returns:
int: 回填的执行记录数
"""
# 检查是否已有 legacy 回填记录
existing = await self.db.execute(
# 检查是否已有 legacy 回填记录(注意:Result.scalar() 只能调用一次,二次调用抛 ResourceClosedError)
legacy_count = (
await self.db.execute(
select(func.count()).select_from(PerformanceExecution)
.where(PerformanceExecution.triggered_by == "legacy")
)
if existing.scalar() and existing.scalar() > 0:
).scalar() or 0
if legacy_count > 0:
logger.info("存量回填已执行过,跳过")
return 0
......
#!/usr/bin/env python
# -*- coding: utf-8 -*-
"""
模块名称:test_performance_satellite_tables.py
模块描述:性能测试大字段从表(P3 拆表)专项测试
覆盖指标:
- 从表 upsert 语义(插入/原地更新/空值删除)
- 任务与执行记录的双写镜像(_sync_task_satellites / _sync_exec_satellites)
- 读取从表优先、主表列回退(_load_task_satellite_fields / _load_exec_satellite_fields)
- 启动存量回填幂等(_backfill_satellite_tables)
- 合并报告 task_summaries 从表优先读(_get_report_task_summaries)
- create_task / update_task CSV 从表联动
作者:czj
创建日期:2026-09-08
"""
import sys
import uuid
from pathlib import Path
import pytest
from sqlalchemy import select, func
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from sqlalchemy.pool import NullPool
# 保证可导入 backend.app
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from app.database import Base # noqa: E402
from app.models.performance import ( # noqa: E402
PerformanceTask,
PerformanceExecution,
PerformanceProjectReport,
PerformanceTaskCsv,
PerformanceTaskResourceSummary,
PerformanceExecutionApiSummary,
PerformanceProjectReportSummary,
)
from app.services.performance_service import PerformanceService # noqa: E402
# ==================== 夹具 ====================
@pytest.fixture()
async def db_session(tmp_path):
"""独立 SQLite 会话(每用例独立库文件,避免跨事件循环复用连接)"""
dbfile = tmp_path / f"sat_{uuid.uuid4().hex}.db"
engine = create_async_engine(
f"sqlite+aiosqlite:///{dbfile.as_posix()}",
poolclass=NullPool,
)
SessionLocal = async_sessionmaker(engine, expire_on_commit=False)
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
async with SessionLocal() as session:
yield session
await engine.dispose()
def _make_task(**overrides) -> PerformanceTask:
"""构造最小任务 ORM 对象(是否落库由调用方控制)"""
defaults = dict(
id=f"perf_{uuid.uuid4().hex[:12]}",
name="拆表测试任务",
target_url="https://192.168.5.44/api/login",
method="POST",
mode="concurrency",
concurrency=5,
duration=10,
csv_parameterization_enabled=True,
csv_content="username,password\nuser1,p1\nuser2,p2",
)
defaults.update(overrides)
return PerformanceTask(**defaults)
def _make_execution(task: PerformanceTask, **overrides) -> PerformanceExecution:
"""构造最小执行记录 ORM 对象(含执行机资源与接口汇总大 JSON)"""
defaults = dict(
id=f"exec_{uuid.uuid4().hex[:12]}",
task_id=task.id,
task_name=task.name,
status="completed",
resource_summary={"cpuSeries": [{"cpuPercent": 30.0}], "avgCpuPercent": 30.0},
api_summary={
"login": {"totalRequests": 100, "avgResponseTime": 12.5, "errorRate": 0.0}
},
)
defaults.update(overrides)
return PerformanceExecution(**defaults)
async def _sat_count(db, sat_cls, owner_id: str) -> int:
"""统计某从表下指定 owner 的行数"""
return (
await db.execute(
select(func.count()).select_from(sat_cls).where(sat_cls.owner_id == owner_id)
)
).scalar()
# ==================== 从表 upsert 语义 ====================
class TestSatelliteUpsert:
"""_satellite_upsert:插入 / 原地更新 / 空值删除"""
async def test_insert_then_update_no_duplicate(self, db_session):
svc = PerformanceService(db_session)
task = _make_task()
db_session.add(task)
await db_session.flush()
await svc._satellite_upsert(
PerformanceTaskResourceSummary, task.id, {"v": 1}
)
await svc._satellite_upsert(
PerformanceTaskResourceSummary, task.id, {"v": 2}
)
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskResourceSummary, task.id) == 1
data = (
await db_session.execute(
select(PerformanceTaskResourceSummary.data).where(
PerformanceTaskResourceSummary.owner_id == task.id
)
)
).scalar_one_or_none()
assert data == {"v": 2}
async def test_empty_value_deletes_row(self, db_session):
"""None 值同步应删除从表行(避免读到陈旧数据)"""
svc = PerformanceService(db_session)
task = _make_task()
db_session.add(task)
await db_session.flush()
await svc._satellite_upsert(
PerformanceTaskResourceSummary, task.id, {"v": 1}
)
await db_session.flush()
assert await _sat_count(db_session, PerformanceTaskResourceSummary, task.id) == 1
await svc._satellite_upsert(PerformanceTaskResourceSummary, task.id, None)
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskResourceSummary, task.id) == 0
async def test_blank_string_deletes_csv_row(self, db_session):
"""空字符串 CSV 视为无数据 → 删除从表行"""
svc = PerformanceService(db_session)
task = _make_task(csv_content="")
db_session.add(task)
await db_session.flush()
# 先写入再清空
await svc._satellite_upsert(PerformanceTaskCsv, task.id, "a,b\n1,2")
await db_session.flush()
assert await _sat_count(db_session, PerformanceTaskCsv, task.id) == 1
await svc._satellite_upsert(PerformanceTaskCsv, task.id, "")
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskCsv, task.id) == 0
# ==================== 双写镜像 ====================
class TestDualWriteMirror:
"""_sync_task_satellites / _sync_exec_satellites:主表 → 从表镜像"""
async def test_task_full_mirror(self, db_session):
"""全字段镜像:csv + 5 个汇总 JSON 全部落从表"""
svc = PerformanceService(db_session)
task = _make_task(
resource_summary={"cpu": [1, 2]},
target_resource_summary={"cpuPercent": 40.0},
transaction_summary={"steps": [{"name": "s1", "avg": 5.0}]},
api_summary={"login": {"totalRequests": 1}},
login_summary={"total": 5, "success": 5, "failed": 0},
)
db_session.add(task)
await db_session.flush()
await svc._sync_task_satellites(task)
await db_session.commit()
for sat_cls in (
PerformanceTaskCsv,
PerformanceTaskResourceSummary,
):
assert await _sat_count(db_session, sat_cls, task.id) == 1
csv_data = (
await db_session.execute(
select(PerformanceTaskCsv.data).where(
PerformanceTaskCsv.owner_id == task.id
)
)
).scalar_one_or_none()
assert "user1,p1" in csv_data
async def test_task_partial_mirror_keeps_other_fields(self, db_session):
"""fields 限定镜像范围时,未指定字段的从表行不受影响"""
svc = PerformanceService(db_session)
task = _make_task(
resource_summary={"cpu": [1]},
login_summary={"total": 3, "success": 2, "failed": 1},
)
db_session.add(task)
await db_session.flush()
await svc._sync_task_satellites(task)
await db_session.commit()
# 主表清空 resource_summary,仅同步该字段
task.resource_summary = None
await svc._sync_task_satellites(task, fields=["resource_summary"])
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskResourceSummary, task.id) == 0
# login_summary 从表行未被波及
assert await _sat_count(db_session, PerformanceExecutionApiSummary, task.id) == 0
from app.models.performance import PerformanceTaskLoginSummary
assert await _sat_count(db_session, PerformanceTaskLoginSummary, task.id) == 1
async def test_exec_mirror(self, db_session):
"""执行记录镜像:resource_summary / api_summary 落 exec 版从表"""
svc = PerformanceService(db_session)
task = _make_task()
db_session.add(task)
await db_session.flush()
execution = _make_execution(task)
db_session.add(execution)
await db_session.flush()
await svc._sync_exec_satellites(execution)
await db_session.commit()
assert await _sat_count(db_session, PerformanceExecutionApiSummary, execution.id) == 1
data = (
await db_session.execute(
select(PerformanceExecutionApiSummary.data).where(
PerformanceExecutionApiSummary.owner_id == execution.id
)
)
).scalar_one_or_none()
assert data["login"]["totalRequests"] == 100
# ==================== 读优先级 ====================
class TestReadPriority:
"""从表优先,主表列回退"""
async def test_satellite_wins_over_stale_main_column(self, db_session):
"""从表有值时以其为准(即使主表列陈旧/不同)"""
svc = PerformanceService(db_session)
task = _make_task()
db_session.add(task)
await db_session.flush()
await svc._satellite_upsert(
PerformanceTaskResourceSummary, task.id, {"from": "satellite"}
)
task.resource_summary = {"from": "main"}
await db_session.flush()
fields = await svc._load_task_satellite_fields(task)
assert fields["resource_summary"] == {"from": "satellite"}
async def test_fallback_to_main_column_when_no_satellite(self, db_session):
"""从表无行时回退主表列(存量未回填场景)"""
svc = PerformanceService(db_session)
task = _make_task(
csv_content="x,y",
login_summary={"total": 1, "success": 1, "failed": 0},
)
db_session.add(task)
await db_session.flush()
fields = await svc._load_task_satellite_fields(task)
assert fields["csv_content"] == "x,y"
assert fields["login_summary"] == {"total": 1, "success": 1, "failed": 0}
# 主表无值 → None
assert fields["resource_summary"] is None
async def test_exec_fallback_and_priority(self, db_session):
svc = PerformanceService(db_session)
task = _make_task()
db_session.add(task)
await db_session.flush()
execution = _make_execution(task)
db_session.add(execution)
await db_session.flush()
# 无从表行 → 回退主表列
fields = await svc._load_exec_satellite_fields(execution)
assert fields["resource_summary"] == execution.resource_summary
# 从表行存在 → 从表优先
await svc._satellite_upsert(
PerformanceExecutionApiSummary, execution.id, {"sat": True}
)
fields = await svc._load_exec_satellite_fields(execution)
assert fields["api_summary"] == {"sat": True}
assert fields["resource_summary"] == execution.resource_summary
# ==================== 存量回填幂等 ====================
class TestBackfill:
async def test_backfill_creates_rows_once(self, db_session):
"""存量行回填从表;二次执行不重复写入"""
svc = PerformanceService(db_session)
task = _make_task(
transaction_summary={"steps": [{"name": "t", "avg": 1.0}]},
)
db_session.add(task)
await db_session.flush()
execution = _make_execution(task)
db_session.add(execution)
await db_session.flush()
report = PerformanceProjectReport(
id=f"rpt_{uuid.uuid4().hex[:12]}",
project_id="proj_x",
task_ids="[]",
task_summaries='[{"task_id": "t1", "avg": 1.0}]',
)
db_session.add(report)
await db_session.flush()
first = await svc._backfill_satellite_tables()
await db_session.commit()
assert first >= 3 # task.csv + task.transaction + exec.resource + exec.api + report
second = await svc._backfill_satellite_tables()
await db_session.commit()
assert second == 0
assert await _sat_count(db_session, PerformanceTaskCsv, task.id) == 1
assert await _sat_count(db_session, PerformanceProjectReportSummary, report.id) == 1
async def test_backfill_skips_empty_values(self, db_session):
"""主表列为空的行不产生从表行"""
svc = PerformanceService(db_session)
task = _make_task(csv_content="")
db_session.add(task)
await db_session.flush()
await svc._backfill_satellite_tables()
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskCsv, task.id) == 0
# ==================== 合并报告从表读 ====================
class TestReportSummaries:
async def test_report_satellite_priority(self, db_session):
"""task_summaries 从表优先,无行回退主表列"""
svc = PerformanceService(db_session)
report = PerformanceProjectReport(
id=f"rpt_{uuid.uuid4().hex[:12]}",
project_id="proj_x",
task_ids="[]",
task_summaries='[{"from": "main"}]',
)
db_session.add(report)
await db_session.flush()
# 无从表行 → 回退主表列
text = await svc._get_report_task_summaries(report)
assert '"from": "main"' in text
# 从表行存在 → 从表优先
await svc._satellite_upsert(
PerformanceProjectReportSummary, report.id, '[{"from": "satellite"}]'
)
text = await svc._get_report_task_summaries(report)
assert '"from": "satellite"' in text
# _report_to_dict 使用 override
result = svc._report_to_dict(report, text)
assert result["task_summaries"] == [{"from": "satellite"}]
# 不传 override 时回退主表列
result_main = svc._report_to_dict(report)
assert result_main["task_summaries"] == [{"from": "main"}]
# ==================== CRUD 联动 ====================
class TestCrudIntegration:
async def test_create_and_update_task_csv_satellite(self, db_session):
"""create_task 写从表;update_task 变更同步;清空删除从表行"""
from app.schemas.performance import PerformanceTaskCreate
svc = PerformanceService(db_session)
task = await svc.create_task(
PerformanceTaskCreate(
name="CRUD联动任务",
target_url="https://192.168.5.44/api/login",
method="POST",
mode="concurrency",
concurrency=2,
duration=5,
csv_parameterization_enabled=True,
csv_content="u,p\na,b",
)
)
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskCsv, task.id) == 1
# 更新 CSV → 从表同步更新
await svc.update_task(task.id, {"csv_content": "u,p\nc,d"})
await db_session.commit()
data = (
await db_session.execute(
select(PerformanceTaskCsv.data).where(
PerformanceTaskCsv.owner_id == task.id
)
)
).scalar_one_or_none()
assert data == "u,p\nc,d"
# 清空 CSV → 从表行删除
await svc.update_task(task.id, {"csv_content": ""})
await db_session.commit()
assert await _sat_count(db_session, PerformanceTaskCsv, task.id) == 0
# 主表列同步为空(双写一致性)
refreshed = await db_session.get(PerformanceTask, task.id)
assert not (refreshed.csv_content or "").strip()
async def test_get_report_mode3_reads_satellite(self, db_session):
"""get_report 模式3(任务直构)读取从表数据"""
svc = PerformanceService(db_session)
task = _make_task()
db_session.add(task)
await db_session.flush()
await svc._sync_task_satellites(task)
await db_session.commit()
report = await svc.get_report(task_id=task.id)
assert report is not None
assert report["resource_summary"] is None # 未设置
assert "user1,p1" not in str(report) # csv 不在报告体中
# 主表列与从表一致时报告能正常取到汇总
task.resource_summary = {"cpu": [1]}
await svc._sync_task_satellites(task, fields=["resource_summary"])
await db_session.commit()
report = await svc.get_report(task_id=task.id)
assert report["resource_summary"] == {"cpu": [1]}
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论