"""隔离库层的分片并行入口:发现与修复期用分片,全量只跑一次做最终判定。 为什么能并行:每个 pytest 会话自带会话前缀的克隆库(`muse_test_<会话>_<随机>`), 角色初始化走共享目录的咨询锁,模板按指纹缓存——互不共享可写状态。 分片只按测试文件切,避免同一文件的夹具被拆到两个进程。 用法: python 工具/并行数据库测试.py --片数 3 # 全部隔离库用例 python 工具/并行数据库测试.py --片数 3 tests/迁移 # 只跑指定范围 退出码:任一分片非零即非零;每片日志落在 /tmp/muse-数据库分片-<序号>.log。 """ from __future__ import annotations import argparse import os import shlex import subprocess import sys import tempfile from concurrent.futures import ThreadPoolExecutor from pathlib import Path 根 = Path(__file__).resolve().parents[1] 基础表达式 = "数据库 and not 宿主 and not 真实模型 and not 浏览器 and not 网络" def 收集分片范围(范围: list[str], 片数: int, 表达式: str) -> list[list[str]]: """按文件把用例均分到各片;同一文件的用例永不拆开。""" 命令 = [ sys.executable, "-m", "pytest", "--外部环境", "-m", 表达式, "--collect-only", "-q", *范围, ] 结果 = subprocess.run(命令, cwd=根, capture_output=True, text=True, timeout=600) if 结果.returncode != 0: print(结果.stdout[-2000:], file=sys.stderr) print(结果.stderr[-2000:], file=sys.stderr) raise SystemExit("分片前采集失败;先修采集错误再分片") 计数: dict[str, int] = {} for 行 in 结果.stdout.splitlines(): if "::" not in 行: continue 文件 = 行.split("::", 1)[0].strip() if 文件.endswith(".py"): 计数[文件] = 计数.get(文件, 0) + 1 if not 计数: raise SystemExit("没有采集到隔离库用例;检查标记表达式与范围") 桶: list[list[str]] = [[] for _ in range(片数)] 负载 = [0] * 片数 for 文件, 数量 in sorted(计数.items(), key=lambda kv: (-kv[1], kv[0])): 最轻 = 负载.index(min(负载)) 桶[最轻].append(文件) 负载[最轻] += 数量 return 桶 def 跑一片(序号: int, 文件: list[str], 表达式: str) -> tuple[int, str, str]: 日志 = Path(tempfile.gettempdir()) / f"muse-数据库分片-{序号}.log" 命令 = [ sys.executable, "-m", "pytest", "--外部环境", "-m", 表达式, "-q", "-p", "no:cacheprovider", "--durations=10", *文件, ] with 日志.open("w", encoding="utf-8") as 输出: 输出.write(" ".join(shlex.quote(x) for x in 命令) + "\n\n") 输出.flush() 进程 = subprocess.run(命令, cwd=根, stdout=输出, stderr=subprocess.STDOUT) 尾部 = "\n".join(日志.read_text(encoding="utf-8").splitlines()[-3:]) return 进程.returncode, str(日志), 尾部 def main(argv: list[str] | None = None) -> int: 解析 = argparse.ArgumentParser(description=__doc__) 解析.add_argument("--片数", type=int, default=3) 解析.add_argument( "--含慢", action="store_true", help="含单例超过三分钟的慢用例;最终全量才需要" ) 解析.add_argument("范围", nargs="*", default=[]) 值 = 解析.parse_args(argv) if 值.片数 < 1: raise SystemExit("片数必须为正整数") # 连接串先查:采集期不连库,缺配置只会在采集之后暴露成误导性的「采集失败」。 环境 = {**os.environ} if not 环境.get("MUSE_TEST_DATABASE_URL"): print("缺少 MUSE_TEST_DATABASE_URL;隔离库层拒绝隐式回退", file=sys.stderr) return 2 范围 = 值.范围 or ["tests"] 表达式 = 基础表达式 if 值.含慢 else f"{基础表达式} and not 慢" 桶 = 收集分片范围(范围, 值.片数, 表达式) print(f"分片:{值.片数} 片,共 {sum(len(b) for b in 桶)} 个测试文件;表达式 {表达式}") for 序号, 文件 in enumerate(桶, start=1): print(f" 第 {序号} 片:{len(文件)} 个文件") with ThreadPoolExecutor(max_workers=值.片数) as 池: 结果 = list(池.map(lambda 项: 跑一片(项[0], 项[1], 表达式), enumerate(桶, start=1))) 失败 = 0 for 序号, (码, 日志, 尾部) in enumerate(结果, start=1): print(f"\n=== 第 {序号} 片 退出码 {码};日志 {日志}\n{尾部}") 失败 += 码 != 0 print(f"\n分片汇总:{值.片数} 片,{失败} 片非零") return 1 if 失败 else 0 if __name__ == "__main__": raise SystemExit(main())