From 20e6d974271fc3b600bbbe332d1f01d49ac342c5 Mon Sep 17 00:00:00 2001 From: iomgaa Date: Sat, 18 Jul 2026 07:52:53 -0400 Subject: [PATCH] =?UTF-8?q?=E5=B1=821/T2:=20teacher.py=20=E6=89=B9?= =?UTF-8?q?=E9=87=8F=E7=94=9F=E6=88=90=20+=20sha256=20JSONL=20=E7=BC=93?= =?UTF-8?q?=E5=AD=98=EF=BC=9Bteacher=20=E6=94=B9=E5=AE=9A=20MiniMax-M3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - teacher.py: 通用 OpenAI 兼容客户端(配置驱动 base_url,替代 OpenRouter 专用); 缓存即断点(逐条落盘+flush,重跑自动续传);单条失败先落盘其余、结束汇总显式报错; M3 思考段 ... 入库前剥离(只剥开头一段) - configs.py: 新增 TeacherGenConfig(采样参数显式化;连接三元组走 .env) - scripts/generate_teacher_completions.py: 自包含生成脚本(本地跑,与训练侧 同 seed 同子集约束已注明) - teacher 决策变更同步:.env.example / docs/00 关键设定与存档点 / docs/02 - tests/test_teacher.py: 10 个单测(假客户端注入),含与 attach 的端到端契约闭环 Co-Authored-By: Claude Fable 5 --- .env.example | 6 +- ars_opd/configs.py | 39 +++++ ars_opd/teacher.py | 187 ++++++++++++++++++++++++ docs/00-roadmap.md | 4 +- docs/02-sft-baseline.md | 4 +- scripts/generate_teacher_completions.py | 38 +++++ tests/test_teacher.py | 143 ++++++++++++++++++ 7 files changed, 414 insertions(+), 7 deletions(-) create mode 100644 ars_opd/teacher.py create mode 100644 scripts/generate_teacher_completions.py create mode 100644 tests/test_teacher.py diff --git a/.env.example b/.env.example index ae4cf11..e433c96 100644 --- a/.env.example +++ b/.env.example @@ -1,8 +1,8 @@ # 复制为 .env 并填入真实值(.env 已被 gitignore,严禁提交密钥) -# teacher API(OpenAI 兼容格式,DeepSeek / MiniMax 二选一填) -TEACHER_API_BASE=https://api.deepseek.com/v1 +# teacher API(OpenAI 兼容格式;2026-07-18 定:自建 new-api 网关 + MiniMax-M3) +TEACHER_API_BASE=https://newapi.iomgaa.online/v1 TEACHER_API_KEY= -TEACHER_MODEL=deepseek-chat +TEACHER_MODEL=MiniMax-M3 # W&B(仅远程训练需要) WANDB_API_KEY= diff --git a/ars_opd/configs.py b/ars_opd/configs.py index 81dc7fd..2ef14e4 100644 --- a/ars_opd/configs.py +++ b/ars_opd/configs.py @@ -111,3 +111,42 @@ class SFTConfig: raise ValueError( f"max_steps 只接受 -1(按 epoch)或正整数,收到 {self.max_steps}" ) + + +@dataclass(frozen=True) +class TeacherGenConfig: + """teacher 批量生成(层 1 能力)的采样与执行参数。 + + 连接信息(API 地址/密钥/模型名)不在这里——那是部署环境的事实,走 `.env` + (teacher.py 读取);这里只放"换一组值就是换一个实验"的采样参数。 + """ + + temperature: float = 1.0 + top_p: float = 0.95 + """MiniMax M 系官方推荐采样参数:temperature=1.0, top_p=0.95。""" + + max_tokens: int = 8192 + """teacher 单条回复的 token 上限。非显然约束:M3 的思考段也计入此额度, + 设太小会把解答挤没(只剩被截断的思考);student 侧超长解答由 collator 的 + completion 预算兜住,这里宁可给足。""" + + strip_think: bool = True + """剥离 content 开头的 ... 思考段。SFT 的监督目标是最终 + 解答;student 以 enable_thinking=False 训练,学思考段会与模板约定矛盾。""" + + concurrency: int = 8 + """并发请求数(线程池大小)。""" + + max_retries: int = 3 + """单请求的网络级重试次数(openai 客户端内建指数退避)。""" + + system_prompt: str | None = None + """None = 不加 system 轮(DAPO 题面自带作答指令,不需要额外指挥)。""" + + def __post_init__(self) -> None: + if self.max_tokens <= 0: + raise ValueError(f"max_tokens 必须为正,收到 {self.max_tokens}") + if self.concurrency < 1: + raise ValueError(f"concurrency 必须 ≥1,收到 {self.concurrency}") + if self.temperature < 0: + raise ValueError(f"temperature 必须 ≥0,收到 {self.temperature}") diff --git a/ars_opd/teacher.py b/ars_opd/teacher.py new file mode 100644 index 0000000..79c1f85 --- /dev/null +++ b/ars_opd/teacher.py @@ -0,0 +1,187 @@ +"""teacher rollout 采样(IO 边缘,论文 §3.2.1)。 + +层 1 起步能力:给一批 prompt 批量生成解答,落盘 sha256 键的 JSONL 缓存 +(键契约在 data.prompt_key 单点定义,本模块与 data.attach_teacher_completions +共用)。层 5 在此长出 chunk 前缀续写的 MC rollout 能力。 + +连接信息从 `.env` 读取(TEACHER_API_BASE / TEACHER_API_KEY / TEACHER_MODEL), +密钥永不出现在代码与配置类里。 + +缓存即断点:生成过程逐条追加写盘,任何中断(网络、Ctrl-C、单条失败)后 +重跑同一命令,已完成的条目自动跳过——API 花的钱不会白花。 +""" + +from __future__ import annotations + +import json +import os +import re +from concurrent.futures import ThreadPoolExecutor, as_completed +from pathlib import Path + +from dotenv import load_dotenv +from openai import OpenAI + +from ars_opd.configs import TeacherGenConfig +from ars_opd.data import prompt_key + +Messages = list[dict[str, str]] + + +def _load_teacher_env(env_file: str | None = None) -> tuple[str, str, str]: + """从 .env(及进程环境)读取 API 连接三元组,缺一项都显式报错。""" + load_dotenv(env_file) + values = {} + for name in ("TEACHER_API_BASE", "TEACHER_API_KEY", "TEACHER_MODEL"): + value = os.environ.get(name, "").strip() + if not value: + raise ValueError( + f"环境变量 {name} 未设置。复制 .env.example 为 .env 并填入真实值。" + ) + values[name] = value + return ( + values["TEACHER_API_BASE"], + values["TEACHER_API_KEY"], + values["TEACHER_MODEL"], + ) + + +def _strip_leading_think(text: str) -> str: + """剥离 content 开头的 ... 段(M3 等 reasoning 模型会内联思考)。 + + 只剥开头一段:解答正文里若出现字面 "" 字样(例如题目在讨论标签本身), + 不应被误删。 + """ + return re.sub(r"^\s*.*?\s*", "", text, count=1, flags=re.DOTALL) + + +class TeacherClient: + """OpenAI 兼容的 teacher 客户端:单条生成 + 采样参数收口。 + + 差异标注:参考实现是 OpenRouter 专用客户端(带其私有请求头与站点字段); + 我们用通用 OpenAI 客户端 + base_url 配置驱动,任何兼容网关(new-api、 + vLLM serve、官方 API)都无需改代码。 + + 测试注入口:传入 client/model 可绕过 .env 与真实网络(见 tests/test_teacher.py)。 + """ + + def __init__( + self, + gen_config: TeacherGenConfig, + client: OpenAI | None = None, + model: str | None = None, + ) -> None: + self.cfg = gen_config + if client is None: + base, key, env_model = _load_teacher_env() + client = OpenAI( + base_url=base, api_key=key, max_retries=gen_config.max_retries + ) + model = model or env_model + if model is None: + raise ValueError("注入 client 时必须同时指定 model") + self.client = client + self.model = model + + def generate(self, messages: Messages) -> str: + """对单条 prompt(messages 列表,末轮为 user)生成解答文本。 + + 返回剥离思考段、去首尾空白后的解答。空解答直接报错——空字符串写进 + 缓存会在训练时变成全 -100 的空样本(trainer 会炸,但应在这里更早炸)。 + """ + if self.cfg.system_prompt is not None: + messages = [ + {"role": "system", "content": self.cfg.system_prompt} + ] + messages + resp = self.client.chat.completions.create( + model=self.model, + messages=messages, + temperature=self.cfg.temperature, + top_p=self.cfg.top_p, + max_tokens=self.cfg.max_tokens, + ) + content = resp.choices[0].message.content or "" + if self.cfg.strip_think: + content = _strip_leading_think(content) + content = content.strip() + if not content: + raise ValueError( + "teacher 返回空解答(可能:max_tokens 太小把思考截断在半途," + "或模型拒答)。该条不会入缓存。" + ) + return content + + +def generate_completions( + prompts: list[Messages], + cache_path: str, + teacher: TeacherClient, +) -> None: + """批量生成解答并追加写入 JSONL 缓存(每行 {"key", "completion", "preview"})。 + + - 已在缓存中的键直接跳过(断点续传); + - 并发线程池执行,每完成一条立即写盘并 flush(中断不丢已完成的结果); + - 单条失败不中断其余任务(并发中的兄弟请求已经花了钱,先让它们落盘), + 全部结束后若有失败则汇总显式报错——重跑即续传,绝不静默缺数据。 + """ + path = Path(cache_path) + path.parent.mkdir(parents=True, exist_ok=True) + + done_keys = _cached_keys(path) + todo = [(prompt_key(p), p) for p in prompts] + todo = [(k, p) for k, p in todo if k not in done_keys] + print( + f"[teacher] 共 {len(prompts)} 条:缓存命中 {len(prompts) - len(todo)}," + f"待生成 {len(todo)},并发 {teacher.cfg.concurrency}", + flush=True, + ) + if not todo: + return + + failures: list[tuple[str, str]] = [] + finished = 0 + # 写盘收口在主线程(as_completed 消费端),工作线程只跑网络请求—— + # 多线程同写一个文件句柄会交错损坏 JSONL + with open(path, "a", encoding="utf-8") as f: + with ThreadPoolExecutor(max_workers=teacher.cfg.concurrency) as pool: + futures = {pool.submit(teacher.generate, p): (k, p) for k, p in todo} + for fut in as_completed(futures): + key, p = futures[fut] + try: + completion = fut.result() + except Exception as e: # noqa: BLE001 —— 收集后统一显式报错,非静默吞错 + failures.append((key, repr(e))) + continue + finally: + finished += 1 + if finished % 20 == 0 or finished == len(todo): + print(f"[teacher] {finished}/{len(todo)} 完成", flush=True) + record = { + "key": key, + "completion": completion, + # preview 仅供人工抽查缓存文件,消费端(attach)只认 key/completion + "preview": p[-1]["content"][:80], + } + f.write(json.dumps(record, ensure_ascii=False) + "\n") + f.flush() + + if failures: + examples = "; ".join(f"{k[:12]}…: {err}" for k, err in failures[:3]) + raise RuntimeError( + f"{len(failures)}/{len(todo)} 条生成失败(成功的已入缓存,重跑本命令" + f"即断点续传)。前几条错误:{examples}" + ) + + +def _cached_keys(path: Path) -> set[str]: + """读取缓存中已有的键集合;文件不存在视为空缓存(首跑)。""" + if not path.exists(): + return set() + keys = set() + with open(path, encoding="utf-8") as f: + for line_no, line in enumerate(f, 1): + if not line.strip(): + continue + rec = json.loads(line) # 坏行直接炸:缓存损坏必须暴露,不能悄悄重新生成 + keys.add(rec["key"]) + return keys diff --git a/docs/00-roadmap.md b/docs/00-roadmap.md index 371e789..0fad931 100644 --- a/docs/00-roadmap.md +++ b/docs/00-roadmap.md @@ -11,7 +11,7 @@ - **当前层**: 层 1(SFT 基线),待开工 - **层 0**: ✅ 已关账(2026-07-18)。两端 pytest 4/4 全绿;本地 env `ars-opd`(torch 2.13 cu130,4070Ti 可做小规模 GPU 调试);远程 env `/data/zym/envs/ars-opd`(torch 2.10 cu128,8 卡可见);gitea 双端打通(SSH 222) - **已完成学习**: 第一章全部精讲(式 1-8、§3.1-3.6、detach 命门专题);第二章已写好待读 -- **下一步**: Claude 编写 T1→T3→T4→T2→T5(工作模式已改:Claude 写码、用户精读提问,见 CLAUDE.md URGENT.1);默认参数已默认通过(deepseek-chat / enable_thinking=False / max_length=4096 / 1k 子集) +- **下一步**: Claude 编写 T1→T3→T4→T2→T5(工作模式已改:Claude 写码、用户精读提问,见 CLAUDE.md URGENT.1);默认参数已通过(teacher=MiniMax-M3 经自建网关 / enable_thinking=False / max_length=4096 / 1k 子集);T1/T3/T4 已完成入库 - **层 1 讨论已完成的**: docs/02 全部难点已精讲(collator 五步流水线与坑二、FSDP 决策与 DDP 触发点、删除/替代清单逐项理由、hash 不稳定演示) - **未精讲的文档账**: docs/01 的 §3.7(KL 锚三处实现差异)、§3.8(论文外稳定器)、§4(训练步流程走读) - **未精讲的文档账**: docs/01 的 §3.7(KL 锚三处实现差异)、§3.8(论文外稳定器)、§4(训练步流程走读) @@ -54,5 +54,5 @@ | KL 锚权重 β | **0.1**(§5.1) | 0.1 | ⚠️ 代码默认 `mc_kl_weight=0` 与论文背离,须显式指定 | | 训练数据 | DAPO-Math-17K(prompt-only) | 同(层 1 先抽 ~1k 子集控制 API 成本) | 一份数据服务层 1-6;学生升到 1.7B 后可直接对表论文 Table 1 | | Student | Qwen3-1.7B / 4B | Qwen3-0.6B | 跑通优先;升级 1.7B 即可与论文对比 | -| Teacher | Qwen3-32B / Claude-4.5-Haiku / Gemini-2.5-Flash | DeepSeek 或 MiniMax(OpenAI 兼容) | logit-free 主路径 | +| Teacher | Qwen3-32B / Claude-4.5-Haiku / Gemini-2.5-Flash | MiniMax-M3(自建 new-api 网关,OpenAI 兼容;2026-07-18 由 DeepSeek 改定) | logit-free 主路径;M3 是 reasoning 模型,思考段入库前剥离(teacher.py strip_think) | | SFT 基线定义 | teacher rollout 上的离线蒸馏(非人写答案) | 同 | 对应参考实现 `_generate_teacher_completions` + JSONL 缓存路径 | diff --git a/docs/02-sft-baseline.md b/docs/02-sft-baseline.md index b83946e..d447840 100644 --- a/docs/02-sft-baseline.md +++ b/docs/02-sft-baseline.md @@ -6,7 +6,7 @@ 式(1)是标准交叉熵,但注意论文 §5.1 对基线的定义:**SFT = 在 teacher rollout 上的离线蒸馏**(Kim & Rush 2016 式 sequence-level distillation),不是"在人写答案上训练"。流程:拿 DAPO-Math-17K 的题目 → teacher 生成解答 → 学生对解答做掩码交叉熵。§4.4 的 Thm 4.4 顺带证明了这种 SFT **不具有** tokenizer/风格不变性(损失绑死 teacher 的具体 token 选择)——这是它后面被 OmniOPD 超越的理论伏笔。 -本层设定(roadmap 已定):student Qwen3-0.6B;teacher DeepSeek(OpenAI 兼容);数据抽 DAPO ~1k 子集控制 API 成本;目标是**管线跑通 + loss 正常下降**,不追分数。 +本层设定(roadmap 已定):student Qwen3-0.6B;teacher MiniMax-M3(自建 new-api 网关,OpenAI 兼容;2026-07-18 由 DeepSeek 改定);数据抽 DAPO ~1k 子集控制 API 成本;目标是**管线跑通 + loss 正常下降**,不追分数。 ## 2. 参考实现解剖 @@ -78,7 +78,7 @@ F.cross_entropy(..., ignore_index=-100) | T4 | 掩码 SFT 损失 + 最小训练循环(HF Trainer 子类) | `ars_opd/trainer.py`(最小形态) | 只做 2.4 那四行的事 | | T5 | 自包含实验脚本(写死全参数,零参数复现) | `scripts/train_sft.sh` | 触发 Video-Tree §2.5 规则接入;显式 CUDA_VISIBLE_DEVICES 4 卡 | -建议顺序 T1→T3→T4(本地可测)→T2(要 API key)→T5(远程)。**默认参数提案**(可否决):teacher 用 `deepseek-chat`(非 reasoner,短答案省钱)、`enable_thinking=False`、`max_length=4096 / max_prompt_length=1024`、子集 1000 题。 +建议顺序 T1→T3→T4(本地可测)→T2(要 API key)→T5(远程)。**默认参数**(已通过;teacher 2026-07-18 改定):teacher 用 `MiniMax-M3`(自建 new-api 网关;M 系是 reasoning 模型,content 可能内联 `` 思考段,入库前由 teacher.py 剥离)、`enable_thinking=False`、`max_length=4096 / max_prompt_length=1024`、子集 1000 题。 ## 5. 验证方式 diff --git a/scripts/generate_teacher_completions.py b/scripts/generate_teacher_completions.py new file mode 100644 index 0000000..16447d8 --- /dev/null +++ b/scripts/generate_teacher_completions.py @@ -0,0 +1,38 @@ +"""层 1:为 DAPO 1k 子集生成 teacher(MiniMax-M3)解答缓存。 + +自包含实验脚本:全部参数写死在此,零参数复现。在**本地**运行(纯 API 调用, +不需要 GPU;本机可直连自建网关): + + conda activate ars-opd + python -u scripts/generate_teacher_completions.py + +前置: +1. .env 已填 TEACHER_API_BASE / TEACHER_API_KEY / TEACHER_MODEL; +2. DAPO parquet 已下载到 DATASET_PATH(见 docs/02 §4)。 + +中断安全:缓存逐条落盘,重跑本脚本自动跳过已完成条目(断点续传)。 +""" + +from ars_opd.configs import SFTConfig, TeacherGenConfig +from ars_opd.data import load_sft_dataset +from ars_opd.teacher import TeacherClient, generate_completions + +# 非显然约束:这里的 dataset/subset_size/seed 必须与 T5 训练脚本完全一致—— +# 两侧各自走"加载→归一→抽子集",seed 相同才是同一批题(data.py 有详注) +DATASET_PATH = "data/dapo-math-17k.parquet" +CACHE_PATH = "data/teacher_completions_dapo1k_minimax-m3.jsonl" + +sft_cfg = SFTConfig( + dataset_path=DATASET_PATH, + output_dir="outputs/_unused", # 本脚本不训练,仅复用数据管线配置 + subset_size=1000, + seed=42, + # teacher_completions_path 留空:此刻缓存尚不存在,取的就是 prompt-only 子集 +) + +dataset = load_sft_dataset(sft_cfg) +prompts = [row["messages"] for row in dataset] + +teacher = TeacherClient(TeacherGenConfig()) # 采样参数全用 configs.py 的显式默认 +generate_completions(prompts, CACHE_PATH, teacher) +print(f"完成。缓存文件:{CACHE_PATH}") diff --git a/tests/test_teacher.py b/tests/test_teacher.py new file mode 100644 index 0000000..781605b --- /dev/null +++ b/tests/test_teacher.py @@ -0,0 +1,143 @@ +"""层 1 / T2:teacher 批量生成与缓存单测。 + +用假 OpenAI 客户端注入(TeacherClient 的测试口),验证思考段剥离、缓存契约 +(与 data.attach_teacher_completions 的端到端闭环)、断点续传、失败汇总。 +""" + +from types import SimpleNamespace + +import pytest +from datasets import Dataset + +from ars_opd.configs import TeacherGenConfig +from ars_opd.data import attach_teacher_completions, prompt_key +from ars_opd.teacher import TeacherClient, _load_teacher_env, generate_completions + + +class FakeClient: + """最小 OpenAI 客户端替身:chat.completions.create 按 responder 出内容。""" + + def __init__(self, responder): + self.calls = [] + self._responder = responder + self.chat = SimpleNamespace(completions=SimpleNamespace(create=self._create)) + + def _create(self, model, messages, **kwargs): + self.calls.append(messages) + content = self._responder(messages) + return SimpleNamespace( + choices=[SimpleNamespace(message=SimpleNamespace(content=content))] + ) + + +def make_teacher(responder, **cfg_overrides): + cfg = TeacherGenConfig(**cfg_overrides) + return TeacherClient(cfg, client=FakeClient(responder), model="fake-m3") + + +def user(q): + return [{"role": "user", "content": q}] + + +# --------------------------------------------------------------------------- +# TeacherClient.generate +# --------------------------------------------------------------------------- + + +def test_剥离开头思考段(): + teacher = make_teacher(lambda m: "心算一下\n答案是 42") + assert teacher.generate(user("q")) == "答案是 42" + + +def test_正文中的think字样不误删(): + teacher = make_teacher(lambda m: "x正文提到 标签本身") + assert teacher.generate(user("q")) == "正文提到 标签本身" + + +def test_只剩思考段等于空解答_报错(): + teacher = make_teacher(lambda m: "思考被截断在半途") + # 未闭合的 think 段剥不掉,但闭合后为空的要报错 + teacher_empty = make_teacher(lambda m: "只有思考 ") + with pytest.raises(ValueError, match="空解答"): + teacher_empty.generate(user("q")) + # 未闭合时保留原文(宁可保留可疑内容也不静默删成空) + assert "" in teacher.generate(user("q")) + + +def test_关闭strip_think则原样保留(): + teacher = make_teacher(lambda m: "ab", strip_think=False) + assert teacher.generate(user("q")) == "ab" + + +def test_system_prompt前置(): + teacher = make_teacher(lambda m: "ok", system_prompt="你是数学助教") + teacher.generate(user("q")) + sent = teacher.client.calls[0] + assert sent[0] == {"role": "system", "content": "你是数学助教"} + assert sent[1]["role"] == "user" + + +def test_注入client但不给model报错(): + with pytest.raises(ValueError, match="model"): + TeacherClient(TeacherGenConfig(), client=FakeClient(lambda m: "x"), model=None) + + +# --------------------------------------------------------------------------- +# generate_completions:缓存契约与断点续传 +# --------------------------------------------------------------------------- + + +def test_端到端契约_生成的缓存能被attach消费(tmp_path): + cache = str(tmp_path / "cache.jsonl") + prompts = [user("1+1=?"), user("2+2=?")] + teacher = make_teacher(lambda m: f"对「{m[-1]['content']}」的解答") + + generate_completions(prompts, cache, teacher) + + ds = Dataset.from_list([{"messages": p} for p in prompts]) + out = attach_teacher_completions(ds, cache) + assert out[0]["messages"][-1]["content"] == "对「1+1=?」的解答" + assert out[1]["messages"][-1]["content"] == "对「2+2=?」的解答" + + +def test_断点续传_已缓存的不重新生成(tmp_path): + cache = str(tmp_path / "cache.jsonl") + prompts = [user("q1"), user("q2")] + teacher = make_teacher(lambda m: "a") + + generate_completions([prompts[0]], cache, teacher) + assert len(teacher.client.calls) == 1 + generate_completions(prompts, cache, teacher) # q1 命中缓存 + assert len(teacher.client.calls) == 2 # 只多了 q2 一次调用 + + +def test_单条失败_其余落盘_结束时汇总报错(tmp_path): + cache = str(tmp_path / "cache.jsonl") + prompts = [user("好题"), user("坏题")] + + def responder(m): + if m[-1]["content"] == "坏题": + raise RuntimeError("网关 500") + return "解答" + + teacher = make_teacher(responder) + with pytest.raises(RuntimeError, match="1/2"): + generate_completions(prompts, cache, teacher) + + # 成功的那条已经在缓存里,重跑只会补坏题 + from ars_opd.teacher import _cached_keys + from pathlib import Path + + assert _cached_keys(Path(cache)) == {prompt_key(prompts[0])} + + +# --------------------------------------------------------------------------- +# .env 读取 +# --------------------------------------------------------------------------- + + +def test_env缺失显式报错(monkeypatch): + for name in ("TEACHER_API_BASE", "TEACHER_API_KEY", "TEACHER_MODEL"): + monkeypatch.delenv(name, raising=False) + with pytest.raises(ValueError, match="TEACHER_API_BASE"): + _load_teacher_env(env_file="/不存在的路径/.env")