把 data-process/ 目录下的每一段代码,与论文《ZGCM-1: A Fully Open and Extremely Efficient Foundation Model for Math and Agentic Search》(arXiv:2609.13356v1)的数据章节逐项对齐:谁定义契约、谁执行清洗、谁做决策、谁负责训练侧打包。
用户问的是“如何操作 data-process 过程”。要回答这个问题,必须先把两份材料的分工讲清楚,否则很容易误读成“代码就是论文的实现”或“论文描述的就是这段代码”。事实不是这样,二者是互补但边界明确的:
unittest + smoke test。QualityPolicy、DATASET_PROFILES、prefix_sample_mix 正是论文可以注入自己数值的插槽;而论文的规模、消融与去污染协议,是代码不下场实现、但必须由训练环境补齐的部分。
decision / reasons / flags / buckets / metrics / record_sha256 的清洗后 JSONL + summary.json + manifest.json,可直接交给下游分片与训练。论文把数据旅程分为三段:General Pre-Training(两个 data stage)、Mid-Training(三个 context stage)、Post-Training(SFT,其后 RL)。代码用三个 stage 目录一一对应,另加一个通用 CLI。
stages/pretrain 与 stages/midtrain。| 阶段 | 论文规模(tokens / 样本) | Context 课程 | 代码入口 | 章节 |
|---|---|---|---|---|
| General Pre-Training Stage 1 | 0.99T | 16K | stages/pretrain/cli.py源族 code / pdf_ocr / general_text | §9 |
| General Pre-Training Stage 2 | 3.20T | 16K | ||
| Mid-Training 16K | 180B | ≤16K | stages/midtrain/cli.py类别 code / web / instruction / agentic / math / reasoning | §10 |
| Mid-Training 64K | 240B | 180.89B@≤16K + 59.11B@16–64K | ||
| Mid-Training 256K | 180.51B | 127.81B@≤16K + 21.72B@16–64K + 30.98B@>64K | ||
| Supervised Fine-Tuning | 4,921,933 条 | 64K / 256K 两档 | stages/posttrain/cli.py | §11 |
| 通用(任意类别) | — | — | zgcm_data_pipeline/cli.py | §14 |
论文附录 A.5 给出:先做池级去重(pool-wide deduplication)与数学错误过滤,得到候选池约 5.3B 条 / 2.86T tokens,再采样成三段累积式 context 混合。代码里这条链的“采样器”是 mix.prefix_sample_mix,“候选池 → 三段”的长度划分器是 midtrain/buckets.assign_length_bucket。
data-process/
├── pyproject.toml # 4 个 console script 入口
├── src/zgcm_data_pipeline/
│ ├── models.py # Stage / DataCategory / Decision 枚举 + Sample 记录
│ ├── core.py # Step / Pipeline / PipelineContext(流式引擎)
│ ├── adapters.py # DatasetAdapter:原始字段 → 规范 Sample
│ ├── quality.py # QualityPolicy + apply_quality_policy
│ ├── dedup.py # simhash64 / hamming_distance / NearDeduper
│ ├── io.py # JSONL(.gz/.zst) 读写
│ ├── pipelines.py # 按类别装配:类别算子 + 公共前后缀
│ ├── training.py # token 化/打包/loss mask/配额/分片 接口
│ ├── cli.py # 通用类别 CLI
│ ├── common/ # 稳定门面:把上述能力重新导出为“可复用算子”
│ ├── steps/ # 6 类类别算子 + common.py(公共前后缀)
│ └── stages/
│ ├── pretrain/ # 3 源族 + gharchive/finepdfs/olmocr profile
│ ├── midtrain/ # 6 类别 + buckets + mix + 3 profile
│ └── posttrain/ # 消息/工具格式 + packing 门面 + 3 profile
├── examples/ # 最小 JSONL/JSON 输入
├── scripts/smoke_test.py # 四个 CLI 的端到端冒烟
└── tests/ # unittest:管道、profile、阶段装配
Pipeline.run 返回 Iterator[Sample],逐条处理,不把全量读进内存——这是能处理 TB 级 JSONL 的前提。DROP,后续算子直接跳过(省算力,也让“丢弃原因”只记第一条决定性的)。cli.py 重新读回输出文件再统计,确保 manifest 记录数与下游 loader 真正读到的完全一致,避免“统计口径 vs 落盘口径”漂移。| 枚举 | 取值 | 含义 |
|---|---|---|
Stage | pretrain / midtrain / posttrain / sft | sft 是历史兼容别名,等价于 posttrain。 |
DataCategory | code / web / agentic / instruction / math / reasoning / pdf_ocr / general_text | 前 6 个是 Midtrain/SFT 的“能力类别”;后 2 个是 Pretrain 的“粗粒度源族”,只走共享管道。 |
Decision | keep / review / drop | 三态决策,是整套管道的核心输出。 |
@dataclass(slots=True)
class Sample:
dataset: str # 数据集名
source_id: str # 源内唯一 ID(缺失时用 "dataset:index" 兜底)
stage: Stage
category: DataCategory
data: dict[str, Any] # 原始字段 + 规范化字段(不丢源记录)
text_fields: tuple[str, ...] # 哪些字段构成“训练文本”
decision: Decision = KEEP
reasons: list[str] = [] # 决策理由(去重、保序)
flags: set[str] = set() # 审计标注(问题标记)
buckets: dict[str, str] = {} # 分桶标签(长度、语言、任务、质量…)
metrics: dict[str, float] = {} # 数值指标(分数、token 数、轮次…)
三个方法决定了它的行为:
| 方法 | 行为 | 注意点 |
|---|---|---|
text() | 把 text_fields 中所有非空字符串字段用换行连接 | “训练文本”由字段声明决定,不同类别可不同(如 math 用 question+solution)。 |
mark(decision, reason) | 单调提升决策等级:DROP 永远生效;REVIEW 仅在当前仍是 KEEP 时生效 | 决策只升不降,保证“一旦判死不再复活”;reason 去重后追加。 |
as_dict() | 序列化为 10 个字段(flags 排序后输出) | data 原样带出,这就是“保留原始记录”的落地方式。 |
@dataclass(slots=True)
class PipelineContext:
config: dict[str, object] # 所有阈值/策略都从这里读
counters: Counter[str] # 每步 seen / keep / review / drop 计数
seen_hashes: set[str] # 精确去重状态(跨整轮)
near_dedup_state: object | None # 惰性创建的 NearDeduper
@dataclass(frozen=True, slots=True)
class Step:
name: str; fn: StepFn; description: str
def __call__(self, sample, context):
before = sample.decision
sample = self.fn(sample, context)
context.count("step." + self.name + ".seen")
if sample.decision != before:
context.count("step." + self.name + "." + sample.decision.value)
return sample
step.<name>.seen 与 step.<name>.<decision> 两类计数,最后汇总进 summary.json。这使得“哪一步丢了多少条”可被直接审计——与论文“AI 迭代脚本直到数据分布满足预设标准”的闭环(§6.2 Data cleaning)天然契合。@dataclass(frozen=True, slots=True)
class DatasetAdapter:
dataset: str
stage: Stage
category: DataCategory
id_field: str = "id"
text_fields: tuple[str, ...] = ("text",)
field_map: Mapping[str, str] | None = None # 源字段名 → 规范字段名
hooks: tuple[Hook, ...] = () # 数据集专属修复,先于共享管道
def adapt(self, raw, index=0):
mapped = dict(raw)
if self.field_map:
for src, canonical in self.field_map.items():
if src in mapped and canonical not in mapped:
mapped[canonical] = mapped[src] # 不覆盖已有规范字段
source_id = str(mapped.get(self.id_field) or self.dataset + ":" + str(index))
sample = Sample(...)
for hook in self.hooks: # 按顺序执行 hook
sample = hook(sample)
return sample
适配器是论文所说“AI 为每个数据集生成独立 adapter/profile”的落点:字段映射与数据专属修复都在这里,而共享的 schema / 质量 / 去重 / 长度 / manifest 行为保持复用。
pipelines.build_pipeline 是整份代码里最需要理解的一行:
steps = category_steps[:1] + COMMON_PREFIX + category_steps[1:] + COMMON_SUFFIX
因为类别算子的第一步做的是字段规范化 / 别名映射(例如把 content 映成 text、把 instruction/output 映成 prompt/response)。只有先完成这层映射,后续 COMMON_PREFIX 里的“空值 / 长度检查”才能作用在规范字段上,而不是误判空文本。这正是注释 “map raw dataset fields into the category's canonical schema … so common checks can safely operate on canonical fields” 的含义。
instruction 为例)| # | 步骤名 | 来源 | 作用 |
|---|---|---|---|
| 1 | normalize_instruction_schema | 类别首步 | 字段别名与 messages→text 组装 |
| 2 | normalize_schema_text | 前缀 | Unicode / 换行 / 空白规范化 |
| 3 | validate_basic_integrity | 前缀 | 空值、过短、过长 |
| 4 | validate_provenance | 前缀 | 来源与许可证元数据 |
| 5 | scan_sensitive_content | 前缀 | 密钥 / 敏感词扫描 |
| 6 | validate_instruction_pair | 类别其余 | 指令-回答完整性、退化样本 |
| 7 | filter_instruction_templates | 类别其余 | 模型套话 / 拒答模板 |
| 8 | bucket_instruction | 类别其余 | 任务 / 风险 / 轮次分桶 |
| 9 | apply_quality_policy | 后缀 | 按源 / 类别阈值决策 |
| 10 | exact_deduplicate | 后缀 | 规范化文本哈希精确去重 |
| 11 | near_deduplicate | 后缀 | 可选 SimHash 近似去重 |
| 12 | estimate_tokens_and_length_bucket | 后缀 | token 估计 + 长度桶 |
| 13 | finalize_quality_bucket | 后缀 | 统一 main/review/drop 质量桶 |
DROP 并短路,不会浪费哈希计算;而长度分桶(步 12)放在最后,保证只有最终 KEEP 的记录才需要计算 token。这个次序由性能与语义共同决定。每条落盘记录至少包含(README 与 release_record 共同定义):
{
"dataset": "example",
"source_id": "row-1",
"stage": "midtrain",
"category": "instruction",
"decision": "keep", // keep | review | drop
"reasons": [], // 决策理由(如 quality_score_below_drop_threshold)
"flags": [], // 审计标注(如 missing_license, possible_secret)
"buckets": {}, // 长度 / 质量 / 语言 / 任务… 分桶
"metrics": {}, // quality_score / estimated_tokens / turns …
"data": {}, // 源记录 + 规范化字段
"record_sha256": "..." // 对 data 排序 JSON 的 sha256
}
默认丢弃记录不写出,但计入 summary;用 --keep-dropped 可保留以便排查。这正是论文“AI 检查数据组成与分布统计,若不符预设标准则自动迭代脚本”所依赖的可观测面。
这些算子在 steps/common.py,是所有类别共享的“公共前后缀”。
normalize_text_value(文本规范化)value = unicodedata.normalize("NFKC", value.replace("\r\n","\n").replace("\r","\n"))
value = _CONTROL.sub("", value) # 去掉控制字符
value = "\n".join(_WHITESPACE.sub(" ", line).strip() for line in value.splitlines())
value = _BLANK_LINES.sub("\n\n", value).strip() # 3+ 连续空行压成 2
| 步骤 | 细节 |
|---|---|
| 换行统一 | 先把 \r\n、\r 统一为 \n。 |
| Unicode 归一 | NFKC 折叠全角 / 兼容字符(如全角数字、连字),减少“看似不同实则相同”的重复。 |
| 控制字符 | _CONTROL 正则 [\x00-\x08\x0b\x0c\x0e-\x1f\x7f] 删除(保留 \t、\n)。 |
| 行内空白 | 每行把 [\t\f\v ]+ 压成单空格,并 strip 行首尾。 |
| 空行压缩 | \n{3,} → \n\n,最后整体 strip。 |
validate_basic_integrity| 条件 | 默认阈值(config 键) | 决策 | reason |
|---|---|---|---|
| 训练文本为空 | — | DROP | empty_training_text |
| 长度 < 最小值 | min_chars = 16 | DROP | too_short |
| 长度 > 最大值 | max_chars = 2,000,000 | REVIEW | too_long |
注意:过短是丢弃,过长是送审——超长记录通常不是噪声,而是需要单独处理的特殊样本。
validate_provenance(来源与许可)require_license=True 且记录无 license(也会查 metadata.license)→ flags += missing_license,REVIEW。data.source 缺失 → flags += missing_source,REVIEW。--require-license,或配置文件里的 require_license。这对应论文“来源、修订、许可证、校验和的登记”(docs/dataset_category_inventory.md 的 cross-category 清单),也是可复现 release 的硬要求。
scan_sensitive_content(敏感内容)_SECRET = re.compile(r"(?i)(?:api[_-]?key|secret|password|access[_-]?token)\s*[:=]\s*[\"']?[A-Za-z0-9_\-/.]{12,}")
# 命中 → flags += possible_secret, REVIEW
# 另支持 config["sensitive_terms"] 列表,大小写不敏感子串匹配 → flags += sensitive_term, REVIEW
exact_deduplicate(精确去重)normalized = normalize_text_value(sample.text()).casefold()
digest = sha256(normalized)
sample.data.setdefault("text_sha256", digest)
if digest in context.seen_hashes: mark(DROP, "exact_duplicate")
else: context.seen_hashes.add(digest)
去重键是规范化 + 小写化后的文本,因此对大小写、全半角、换行差异免疫。状态存在 PipelineContext.seen_hashes,一次 run 内跨所有记录生效。
estimate_tokens_and_length_bucketconfig["token_estimator"] 是可调用对象 → 用它;否则回退到 max(1, (len(text)+3)//4)(约 4 字符/token 的可复现估计)。metrics["estimated_tokens"]。(2048, 8192, 32768, 65536),标签形如 1_2048、2049_8192、32769_65536、gt_65536;边界会被排序去重。(2048, 8192, 32768, 65536) 的细粒度桶(标签 a_b);而 Midtrain 的 buckets.py 用的是 (16384, 65536, 262144) 的 B16/B64/B256 桶,专门服务于论文的三段 context 采样。finalize_quality_bucket 与 release_record# finalize_quality_bucket
bucket = "main" if KEEP else "review" if REVIEW else "drop"
sample.buckets.setdefault("quality", bucket)
# release_record(落盘前最后一步)
result = sample.as_dict()
result["record_sha256"] = sha256(json.dumps(result["data"], ensure_ascii=False, sort_keys=True, default=str))simhash64:确定性 64 位指纹def simhash64(text, *, ngram=5):
tokens = _TOKEN.findall(text.casefold()) # \w+(Unicode)
if not tokens: return 0
width = max(1, int(ngram))
features = (" ".join(tokens[i:i+width]) for i in range(max(1, len(tokens)-width+1)))
weights = [0]*64
for feature in features:
digest = blake2b(feature.encode(), digest_size=8) # 64 位哈希
for bit in range(64):
weights[bit] += 1 if digest & (1<<bit) else -1
result = 0
for bit, w in enumerate(weights):
if w >= 0: result |= (1<<bit)
return result
算法要点:把文本切成 词 5-gram;每个 gram 用 BLAKE2b-64 得到 64 位;对每一位做 ±1 投票(位为 1 记 +1,否则 −1);最终按符号定该位(≥0 取 1)。相同文本必然得到相同指纹,相似文本得到汉明距离很小的指纹。
NearDeduper:带索引的流式近似去重(附交互演示)朴素做法是“新指纹与历史全部两两比较”,复杂度 O(n²)。这里用 band 索引把候选集缩到极小:
threshold, bands = 3, 4 # 不变量:bands 必须整除 64,且 threshold < bands
width = 64 // bands # 16 位/段
values = [(fp >> (band*width)) & mask for band in range(bands)]
# check_and_add(text):先算指纹 → 取 4 段 → 只在“共享任一段”的候选里比汉明距离
candidates = set().union(*[indexes[b].get(values[b], ()) for b in range(bands)])
duplicate = any(hamming_distance(fp, fingerprints[i]) <= threshold for i in candidates)
# 无论是否重复,都把新指纹登记进 4 个 band 索引
基准指纹 A 固定;对 B 翻转若干位,观察:汉明距离、四段中哪些段相同、以及 NearDeduper 是否会把 B 判为重复(阈值 3、4 段)。
配置:near_dedup_threshold(默认 3)、near_dedup_bands(默认 4)、near_dedup_ngram(默认 5)。近似去重默认关闭(near_dedup=False),因为它对阈值敏感、必须按类别校准。开启方式:examples/quality-policies.json 里写 "near_dedup": true。
seen_hashes 是单次 run 作用域)。quality.py 把“质量判断”抽象成一份可按数据源或类别配置的策略,而不写死任何数据集。
| 字段 | 默认 | 作用 |
|---|---|---|
score_field | quality_score | 从 data 里读哪个字段作为质量分 |
keep_min | None | 低于它 → REVIEW(quality_score_below_keep_threshold) |
review_min | None | 低于它 → DROP(quality_score_below_drop_threshold)(命名上是“送审下限”,实为丢弃线) |
hard_error_field / hard_error_max | None | 硬错误分超过上限 → DROP(hard_error_score_exceeded) |
require_score | False | 缺分时是否标 missing_quality_score 并 REVIEW |
key = data["source_key"] or data["source"] or dataset or category
policy = policies[key] or policies[category] or policies["default"]
即:源特定 → 类别 → default。这让同一类别下的不同来源可以有不同的严格度。
score = _number(data[score_field]) # bool 被显式排除;字符串会尝试转 float
if score is None:
if require_score: flags.add("missing_quality_score"); mark(REVIEW, "missing_quality_score")
else:
metrics["quality_score"] = score
if review_min is not None and score < review_min: mark(DROP, "quality_score_below_drop_threshold")
elif keep_min is not None and score < keep_min: mark(REVIEW, "quality_score_below_keep_threshold")
if hard_error_field and hard_error_max is not None:
he = _number(data[hard_error_field])
if he is not None:
metrics["hard_error_score"] = he
if he > hard_error_max: mark(DROP, "hard_error_score_exceeded")
examples/quality-policies.json){
"near_dedup": true,
"near_dedup_threshold": 3,
"quality_policies": {
"web": { "score_field": "quality_mean", "keep_min": 4.0, "review_min": 3.0, "require_score": false },
"math": { "score_field": "quality_score", "keep_min": 0.8, "review_min": 0.5,
"hard_error_field": "hard_error_score", "hard_error_max": 3 }
}
}
apply_quality_policy(阈值 → keep/review/drop)与 web_knowledge.quality_band(ge4/ge3/lt3)正是这层“分档”的可复用实现;evaluator 本身属于模型资产,不在代码内。入口:stages/pretrain/cli.py,装配在 pipeline.py。支持三个源族,装配顺序与通用管道不同:
steps = dataset_steps(profile, source) + source_steps + COMMON_PREFIX + COMMON_SUFFIX
# ① 数据集档案(可选) ② 源族归一化 ③④ 公共前后缀
text,公共检查必须在其后。以 gharchive 为例共 1+5+4+5 = 15 步。--source | 归一化要点 | 默认过滤 |
|---|---|---|
| code | 复用 code 类别的 5 步(见 §13.1) | vendor / 生成物 / FIM 完整性 |
| pdf_ocr | text 缺失时依次尝试 content → ocr_text → document;缺失 text → DROP missing_extracted_text | 字符数 < pdf_min_chars(32) → REVIEW short_ocr_text;ocr_quality < pdf_min_ocr_quality(0.0) → REVIEW low_ocr_quality |
| general_text | 依次尝试 content → document → body → raw_text;缺失 → DROP missing_text | — |
gharchive(代码源,最完整的一套过滤)reject_path_reason 是逐条路径的分类器,返回第一个命中的原因:
| 判定 | 规则 | DROP reason |
|---|---|---|
| 压缩产物 | 文件名以 .min.js/.min.css/.min.html 结尾 | gharchive_minified_asset |
| vendor / 构建目录 | 路径任一段命中 REJECT_DIRS(.git .hg .svn .cache .next .nuxt .venv __pycache__ bin build coverage dependencies deps dist node_modules obj out site-packages target third_party vendor vendors) | gharchive_vendor_or_build_path |
| 锁文件 | 命中 LOCK_FILES(cargo.lock / composer.lock / gemfile.lock / go.sum / package-lock.json / pnpm-lock.yaml / poetry.lock / yarn.lock) | gharchive_lockfile |
| 二进制后缀 | 命中 BINARY_EXTENSIONS(37 项:图片 / 音视频 / 压缩包 / office / 模型权重 / 可执行…) | gharchive_binary_extension |
| 测试夹具 / 快照 | 路径含 /__snapshot__/ /__snapshots__/ /fixtures/ /golden/ /test/snapshot/ /testdata/ … | gharchive_fixture_or_snapshot |
| 机器报告 | 路径含 /html-report/ /reports/ /test-output/ | gharchive_machine_report |
| 体积上限 | size_bytes > gharchive_max_file_bytes(>0 时启用) | gharchive_configured_size_limit |
| 不支持类型 | language_for_path==None | gharchive_unsupported_file_type |
此外还有内容层过滤:\x00(前 8192 字符内)→ DROP gharchive_binary_nul;前 2000 字符命中 GENERATED_MARKERS(auto-generated / autogenerated / code generated / do not edit / do not modify / generated by)→ DROP gharchive_generated_file;is_minified(长度 ≥2000 且最长行 >200k,或 <10 行且平均行 >1000)→ DROP gharchive_minified_or_long_line。
normalize_notebook 只保留 code/markdown cell 的 source,丢弃 execution outputs 与 metadata,并用 # code cell / <!-- markdown cell --> 标记。LANGUAGE_BY_EXTENSION(约 30 种后缀)+ SPECIAL_NAMES(dockerfile / makefile / cmakelists.txt / rakefile / gemfile)。metadata.blob_id = "sha256:<text sha256>",并写回 path/language,buckets["gharchive_bucket"]="core"。finepdfs(PDF/OCR 语言 + 重叠 + 质量)| 指标 | 字段别名 | 默认阈值(config 键) | 决策 |
|---|---|---|---|
| 语言分 LID | lid_score / language_score / language_probability / language_confidence | ≥ 0.85(finepdfs_lid_threshold) | 低于 → DROP finepdfs_low_language_score |
| MinHash 重叠 | minhash_count / minhash_matches / minhash_overlap | ≤ 2.0(finepdfs_minhash_max) | 高于 → DROP finepdfs_excessive_minhash_overlap |
| Primary 质量分 | primary_score / quality_score / edu_score / pdf_quality_score | ≥ 0.8(finepdfs_primary_threshold) | 低于 → DROP finepdfs_low_primary_score |
| 任一指标缺失 | — | — | flags finepdfs_missing_quality_metric + REVIEW |
另生成 data.finepdfs_dedup_key:优先 id:<id>,否则用 url/dump/file_path/offset 组成的 provenance 键,再否则用文本 sha256;buckets["finepdfs_policy"]="relaxed_lid085_mh2_p08" 把“用的是哪版策略”写进记录,便于追溯。
olmocr(科学 PDF 压缩比过滤)ratio = len(zlib.compress(text.encode("utf-8"), 6)) / len(raw) # 0.0 若空
# ratio < olmocr_min_compression_ratio(0.08) → DROP olmocr_low_compression_ratio
# edu_score < olmocr_min_edu_score(0.6) → DROP olmocr_low_edu_score
# (非法值 → REVIEW olmocr_invalid_edu_score)
# buckets["olmocr_policy"] = "compression_ratio_ge_0p08"
≥0.08 规则就是其“清洗伪影”的可复现表达。入口:stages/midtrain/cli.py;类别 = code / web / instruction / agentic / math / reasoning。装配:
steps = profile_steps + category_pipeline.steps # profile 前置,类别管道已含公共前后缀
def assign_length_bucket(token_count, boundaries=(16_384, 65_536, 262_144)):
for upper in boundaries:
if value <= upper: return "B" + str(upper // 1024) # B16 / B64 / B256
return "gt_" + str(lower // 1024) # gt_256
| Token 数 | 桶 | 论文语义 |
|---|---|---|
| ≤ 16,384 | B16 | ≤16K 常规上下文 |
| 16,384 < n ≤ 65,536 | B64 | 16K–64K 长文档 / 长推理 |
| 65,536 < n ≤ 262,144 | B256 | 64K–256K 超长 / 跨文档 / 长程 agent |
| > 262,144 | gt_256 | 超出 256K,一般排除 |
DEFAULT_ELIGIBLE = {
"mix16": ("B16",),
"mix64": ("B16", "B64"),
"mix256": ("B16", "B64", "B256"),
}
prefix_sample_mix 按 mix16 → mix64 → mix256 的规范顺序处理;每段先从“允许的基础桶”里剔除已被前序 mix 选走的 ID,再用 blake2b(seed:final_name) 稳定排序取前 target 个。性质:
"zgcm-midtrain";候选不足会抛错,绝不静默少采。| 逻辑 | 细节 |
|---|---|
| 质量分 | quality_mean(也会查 meta/metadata);缺失 → flags + REVIEW web_knowledge_missing_quality_score,并置 quality_missing=True |
| 质量档 | ≥4.0 → ge4;≥3.0 → ge3;否则 lt3 |
| 长度桶 | le_16k / gt_16k_le_64k / gt_64k_le_256k / gt_256k(token_count 缺失 → flags + REVIEW) |
| 最低档门限 | web_knowledge_min_quality_band(默认 lt3,可设 ge3/ge4);低于门限 → DROP web_knowledge_below_quality_band |
record_drop_key 归一出行级谱系键 (source, row_index)(支持 source_rel/source_file/source_path 与 row_index/file_row_index)。normalize_drop_keys 接受字符串("src:12" 或 tab 分隔)、序列、映射三种配置形式。hard_error_score / math_hard_error_score / error_risk_score,默认阈值 4.0;命中配置键或分数 ≥ 阈值 → DROP nemotron_math_strict_hard_error;两者都缺 → REVIEW nemotron_math_missing_hard_error_signal。buckets["hard_error_risk"] = "high" | "accepted"。rough_tokens = chars / 3.6;桶 le16k / gt16k_le64k / gt64k_le256k / gt256k。agent_coding_same_repo_multiissue_pack=0 < agent_coding_trace64k_packed_full=1 < agent_coding_nonzerohero_singletrace=2 < 其它=3(dedup_priority)。near_dedup_sample_sha256),避免只比前缀。text_chars / rough_tokens_chars_div_3_6 / length_bucket / dedup_text_sha256 / pre_rendered_agent_trace=True / training_tags / completion_status / source_trace_keys。agent_coding_unverified_completion。nemotron_math 的硬错误过滤、assign_length_bucket 与 prefix_sample_mix。steps/reasoning.py 提供“推理信号 / 最终答案 / 控制残留”的结构化检查骨架。--keep-dropped 排查”正是这套闭环的工程接口。pipelines.build_pipeline(Common Cleaning = 公共前后缀)+ steps/*(Type-Specific)+ buckets.py + mix.py 的合并视图。入口:stages/posttrain/cli.py;类别 = instruction / agentic / math / reasoning / code / web。装配:
steps = profile_steps + POSTTRAIN_FORMAT_STEPS + category_pipeline.steps
# ① 数据集档案 ② 消息/工具格式规范化 ③ 类别管道(含公共前后缀)
_ROLE_ALIASES = {"human":"user", "gpt":"assistant", "model":"assistant", "function":"tool"}
def split_think_content(content):
m = _THINK.search(content) # 匹配 think 标签包裹的推理块
if not m: return None, content.strip()
reasoning = m.group(1).strip()
visible = _THINK.sub("", content).strip()
return reasoning or None, visible
normalize_messages:角色别名归一、非字符串 content 转字符串、把 assistant 的第一个 think 块抽到 reasoning_content,并记录 metrics["turns"]。normalize_tool_calls:若 call 里嵌套 function{name,arguments} 则上提 name/arguments;arguments 非字符串则转字符串;不强制任何厂商 schema。dolci_think(思考数据,8 个来源白名单)| 能力标签 | 来源(dataset_source) |
|---|---|
| openthoughts_math | saumyamalik/OpenThoughts3-full-filtered-math-decontam-v2 |
| dolci_python_code | …/correct-python-sft-187k-x16-thoughts-filtered-decontam-v2 |
| persona_precise_if | allenai/persona-precise-if-r1-final-content-filtered-chinese-filtered |
| qwq_verified_if | …/if_qwq_reasoning_verified_filtered_decontam-v2 |
| nemotron_code_no_tool | allenai/nemotron-post-training-dataset-subset-ngram-filtered-no-tool-calls |
| synthetic2_verified | allenai/SYNTHETIC-2-SFT-cn-fltrd-final-ngram-filtered-chinese-filtered |
| openthoughts_science | …/OpenThoughts3-full-filtered-science-decontam-v2 |
| openthoughts_code | …/OpenThoughts3-full-filtered-code-subsampled-decontam-v2 |
parse_assistant 对 assistant 内容做严格解析,返回 (reasoning, answer, error):无完整 think 块 → missing_complete_think;think 为空 → empty_think;可见回答为空 → empty_visible_answer;出现多个 think 块 → multiple_think_blocks。此外:
TEMPLATE_TOKEN_RE 命中外部模板 token → strict 时 DROP dolci_think_template_token,否则 REVIEW。has_repetition:≥120 个词且 5-gram 重复率 > 20% → DROP dolci_think_repetition。dolci_think_max_rough_tokens(32768) → DROP dolci_think_over_length;超 16000 → REVIEW dolci_think_long_record。dolci_tooluse(legacy 工具调用 → 原生 tools/messages)load_functions:把 system 里的 functions(JSON 字符串或列表)转成 {"type":"function","function":{...}}。strip_legacy_system:剥离 “You are a helpful function-calling AI assistant…” 这段 legacy 系统提示;若残留 <functions>/<function_calls> 则整段清空。split_calls:用一个带引号 / 转义状态机的扫描器把 weather(city='Beijing') 这类文本按顶层括号切分;parse_arguments 用 ast.parse("f(" + text + ")", mode="eval") 解析实参(位置参数 → argN,关键字参数 → 原名)。system 的 functions 上提到顶层 tools;environment → tool;带 function_calls 的 assistant → 带 tool_calls 的 assistant。validate_dolci_tooluse_record:校验角色集合、system 必须在最前、assistant 前一条必须是 user/tool、tool 前一条必须是 assistant/tool、assistant 不能既无 content 又无 tool_calls、每个 call 必须有 name + dict arguments。dolci_tooluse 的两步(convert + validate)就是这条协议的实例化。ultradata(消息修复)FOREIGN_TOKENS(<|end_header_id|>、<|start_header_id|>、<|model|>、<|messages|>、<|im_start|>、<|im_end|>、<|eot_id|>)出现在 content / reasoning_content / reasoning → DROP ultradata_foreign_chat_template_token。flags += ultradata_removed_empty_system)。repair_visible_assistant_content:若 assistant 的可见 content 里漏了 think 结束标签,取最后一个结束标签之后的内容作为可见答案;修完仍为空 → DROP ultradata_empty_after_think_repair。POSTTRAIN_FORMAT_STEPS(Schema Formatting)+ profile/category 算子(Filter & Verify)+ 训练环境的 mixture / decontamination / 打包。这些是“清洗完成到可训练”的桥梁,都在 training.py(并通过 common/ 与 posttrain/packing.py 重新导出)。它们是与 tokenizer 无关的纯算法接口。
| 函数 | 算法 | 用途 |
|---|---|---|
balance_work_units(units, workers) | LPT(最长优先,配给小负载工人;有序列表 + bisect.insort) | 把重分片先派给最空的 worker |
stable_partition(key, partitions, seed) | blake2b(seed\0key, 8) 取模 | 把样本确定性分到 N 个分区 |
deterministic_shuffle_key(key, seed) | blake2b(seed\0key, 16) | 可复现的样本级 shuffle 排序键 |
plan_row_microshards(source, counts, target_tokens) | 连续行区间贪心切块 | 按 token 预算把源切成 microshard,行域半开 [start,end) |
def allocate_mix_quotas(weights, total_samples):
exact = {k: total*w/sum(w)}
quotas = {k: floor(v)}; remaining = total - sum(quotas)
for k in sorted(k, key=(-(exact[k]-quotas[k]), k))[:remaining]: quotas[k] += 1
return quotas # 最大余数法,保证配额之和精确等于 total
def select_mixed_sample_ids(source_to_ids, quotas, seed="zgcm"):
for source: ranked = sort(ids, key=blake2b(seed + ":" + source)); take ranked[:quota]
selected.sort(key=blake2b(seed + ":mix")) # 确定性交织
配额不足会抛 ValueError(source X has N but quota is M),杜绝“悄悄少训”。这一层直接对应论文 §3.1.1 “before each source is sampled according to the final data recipe”。
IGNORE_INDEX = -100
@dataclass(frozen=True, slots=True)
class TokenizedSample:
tokens: tuple[int, ...]
source_id: str
targets: tuple[int, ...] | None = None # None = 纯预训练;有值 = SFT
def validate_tokenized_sample(s):
if not s.tokens: raise ValueError("tokenized sample is empty")
if s.targets is not None:
if len(s.tokens) != len(s.targets): raise ValueError("tokens and targets must have the same length")
if not any(t != IGNORE_INDEX for t in s.targets): raise ValueError("SFT sample has no trainable target token")
def build_loss_mask(targets): # 0 = 不计损失,1 = 计损失
return tuple(0 if t == IGNORE_INDEX else 1 for t in targets)
def pack_tokenized_samples(samples, sequence_length, allow_overlength=False):
# 1) 校验每条;超长进 overlength;不允许则跳过
# 2) accepted 按 (-length, source_id) 排序(first-fit-decreasing)
# 3) free 列表 + bisect 找“剩余空间足够”的 bin;否则新开 bin
# 4) 追加 tokens/targets,记录 cu_seqlens(累积边界),更新 free
# 5) 禁止在一次调用里混装 pretrain 与 SFT 样本
return bins, overlength
| 要点 | 含义 |
|---|---|
cu_seqlens | 打包后每条子序列的累计边界([0, len1, len1+len2, …]),供变长注意力使用。 |
| first-fit-decreasing | 先放大样本,减少碎片;free 用有序列表 + bisect 快速定位。 |
| targets 一致性 | 同一 bin 内要么全是预训练样本(targets=None),要么全是 SFT;混装直接报错。 |
| overlength 分流 | 超长样本单独返回,不静默截断——与论文 B.3 的 accepted / overlength / rendering-error 三路一致。 |
targets:system / user / tool observation 位置填 -100,assistant 的推理 / 回答 / tool_calls 位置填真实 token id。build_loss_mask 再把它转成 0/1。论文 §4.1.3 明确:“Assistant reasoning, responses, and structured tool calls contribute to the loss; system and user messages and tool observations remain masked context”。steps/code.py)| 算子 | 逻辑 |
|---|---|
normalize_code_schema | content → text(若 text 缺失);buckets["subtype"] 默认 code |
classify_code_language | 读 language 或 metadata.language → buckets["language"] |
filter_non_source_assets | 路径命中 code_denied_path_parts(默认 node_modules/ vendor/ .min.js .map)→ DROP non_core_or_vendored_file;前 4096 字符命中生成标记 → REVIEW possibly_generated_code |
validate_code_signal | 无代码符号({ } ( ) ; 或 def/class/function/import/package/SELECT/FROM)且 <3 行 → REVIEW weak_code_signal;FIM 三元组 prefix/middle/suffix 不齐 → DROP incomplete_fim_triplet,齐 → subtype=fim |
bucket_code_task | 有 patch/diff → subtype=patch;有 prompt+completion → subtype=code_qa;buckets["verified"]=verified/unverified |
steps/web.py)normalize_web_schema:question+answer 都是字符串 → text_fields=(question,answer)、subtype=web_qa;否则 subtype=web_text。filter_web_boilerplate:模板词(cookie policy / accept all cookies / privacy policy / sign in / subscribe now)命中 ≥ web_boilerplate_limit(4) → REVIEW web_boilerplate;非空行去重率 < 0.45 → DROP repetitive_web_text。validate_web_qa:question <4 或 answer 空 → DROP incomplete_qa_pair;answer 命中 unknown / n/a / i don't know / 无法回答 / 不知道 → DROP unknown_answer;grounded=False → REVIEW ungrounded_answer。bucket_web:topic 分桶;quality_score ≥0.8 / ≥0.5 → high / medium / low。steps/instruction.py)normalize_instruction_schema:instruction→prompt、output→response;messages 列表则拼成 role: content 的 text。validate_instruction_pair:messages 缺 user/assistant → DROP incomplete_instruction_messages;prompt/response 缺失 → DROP incomplete_instruction_pair;两者完全相同 → DROP prompt_response_identical。filter_instruction_templates:response 命中 as an ai language model / i cannot assist with / here is the requested response → REVIEW generic_model_template。bucket_instruction:task / risk / turns(>2 条 messages 记 multi)。steps/agentic.py)_messages:从 messages 或 trajectory 取字典列表。normalize_agentic_schema:无 text 时把 messages 拼成 text;有 text 但无 messages 且非 pre-rendered → REVIEW unstructured_agent_trace;metrics["turns"];subtype=sample_type。validate_role_order:首角色非 system/user → REVIEW invalid_first_role;空角色 → DROP missing_message_role。validate_tool_call_pairs(pre-rendered 跳过):统计 tool_calls 与 observations,配对未闭合 → REVIEW unpaired_tool_call,零调用 → REVIEW no_tool_call。detect_trajectory_outcome:失败标记(exit due to / permission denied / tool error / timed out)与成功标记(submitted / tests passed / task completed / success);buckets["outcome"]=verified/failed/weak;failed 且未 verified → REVIEW。bucket_agentic:tool_family 由 tool_names/target_tools 排序去重拼接。steps/math.py)missing_problem_or_solution;出现图片依赖(<img / 图片后缀 / includegraphics / see figure)但无 image_text → REVIEW missing_multimodal_context。solution_too_short;无最终答案信号(final answer / answer: / therefore / thus / boxed / 答案 / 所以)且无 answer 字段 → REVIEW missing_final_answer_signal。verifier_failed。steps/reasoning.py)[system prompt]、api_error)→ REVIEW control_or_prompt_artifact。missing_reasoning;<64 字符或无推理信号 → REVIEW weak_reasoning_signal;无最终答案 → REVIEW missing_final_answer。source_row_hash → lineage_hash(保留谱系)。| console script | 模块 | 必需参数 |
|---|---|---|
| zgcm-clean | zgcm_data_pipeline.cli | --input --output --dataset --stage --category |
| zgcm-pretrain-clean | …stages.pretrain.cli | --input --output --dataset --source |
| zgcm-midtrain-clean | …stages.midtrain.cli | --input --output --dataset --category |
| zgcm-posttrain-clean | …stages.posttrain.cli | --input --output --dataset --category |
共同选参:--dataset-profile、--id-field、--text-fields、--config、--require-license、--manifest、--summary、--keep-dropped。
# ① 通用 instruction(无 profile)
PYTHONPATH=src python3 -m zgcm_data_pipeline \
--input examples/instruction.jsonl \
--output outputs/instruction.cleaned.jsonl \
--dataset example_instruction --stage midtrain --category instruction \
--text-fields instruction,output \
--config examples/quality-policies.json \
--manifest outputs/instruction.manifest.json
# ② Pretrain code + gharchive 档案
PYTHONPATH=src python3 -m zgcm_data_pipeline.stages.pretrain.cli \
--input examples/datasets/gharchive.jsonl \
--output outputs/gharchive.cleaned.jsonl \
--dataset gharchive --source code --dataset-profile gharchive
# ③ Posttrain agentic + dolci_tooluse 档案
PYTHONPATH=src python3 -m zgcm_data_pipeline.stages.posttrain.cli \
--input examples/datasets/dolci_tooluse.jsonl \
--output outputs/dolci_tooluse.cleaned.jsonl \
--dataset dolci_tooluse --category agentic --dataset-profile dolci_tooluse
| 示例 | 内容 | 预期 |
|---|---|---|
| instruction.jsonl | 2 条:一条完整、一条 output 为空 | 完整条 keep;空条 drop(incomplete_instruction_pair) |
| datasets/gharchive.jsonl | src/example.py 与 vendor/example.py | src 保留(lang=Python、写 blob_id);vendor 丢弃(gharchive_vendor_or_build_path) |
| datasets/web_knowledge.jsonl | quality_mean 4.2 / 2.4 | 默认 lt3 门限两条都留;ge3/ge4 时低分条被丢 |
| datasets/dolci_tooluse.jsonl | weather 函数 + weather(city='Beijing') + environment | functions 上提为 tools;call 转 tool_calls;environment → tool 角色 |
PYTHONPATH=src python3 -m unittest discover -s tests -v
python3 -m compileall -q src
python3 scripts/smoke_test.py # 跑四个 CLI,校验 manifest.records == 输出行数| 论文概念 | 论文位置 | 代码实现 | 关系 |
|---|---|---|---|
| 共同记录 schema | §3.1.1 | models.Sample / adapters.DatasetAdapter | 直接对应;代码额外带决策 / 审计字段 |
| 源特定语言与质量过滤 | §3.1.1 | stages/*/source_cleaning.py + datasets/* profile | 直接对应(记录级) |
| 仓库清洗 / 文本抽取 | §3.1.1 | gharchive.clean_gharchive_record / normalize_pdf_ocr | 代码给规则;集群抽取在外 |
| 版本化分片 + source manifest | §3.1.1 | common/manifest.build_manifest / file_sha256 | 代码给 manifest;分片由训练环境 |
| GLM-5.1 分词与索引 | §3.1.1 | training.TokenizedSample + estimate_tokens… | 代码是接口 / 估计;真实 tokenizer 在外 |
| 课程预训练(词法复杂度排序) | §3.1.2 | deterministic_shuffle_key / stable_partition | 代码给确定性排序原语;复杂度打分在外 |
| 池级去重 | §3.2.1 / A.5 | dedup.NearDeduper / steps.common.exact_deduplicate | 代码给算子;“池级”作用域由环境组织 |
| 跨阶段去重(Stage1 ⊄ Stage2) | §3.1.1 | —(代码 seen_hashes 为单 run 作用域) | 论文有、代码未内建 |
| 长度分桶 B16/B64/B256 + 前缀采样 | §3.2.1 / 图 7 | midtrain/buckets.py + midtrain/mix.py | 直接对应 |
| 数据配方采样(按源权重) | §3.1.1 / 图 6 | training.allocate_mix_quotas + select_mixed_sample_ids | 直接对应 |
| 模型评分 + 固定多准则分档 | §4.1.1 | quality.QualityPolicy + web_knowledge.quality_band | 代码给策略 / 分档;evaluator 是模型资产 |
| 8-gram >50% 去污染 | §4.1.1 | —(代码内无 n-gram 去污染算子) | 论文有、代码未内建 |
| SFT 三路 schema 格式化 | §4.1 / 图 10 | posttrain/format.py + 三个 profile | 直接对应 |
| tool-protocol 校验与修复 | §4.1.2 | dolci_tooluse convert + validate | 直接对应 |
| think / no-think 双模 | §4.1.1 / B.4 | format.split_think_content / dolci_think.parse_assistant | 直接对应 |
| assistant-only loss mask | §4.1.3 / B.3 | training.build_loss_mask + IGNORE_INDEX | 直接对应 |
| 打包与 cu_seqlens / overlength 分流 | B.3 | training.pack_tokenized_samples | 直接对应 |
| 混合配额(normalized weights) | §3.1.1 / §4.1.1 | training.allocate_mix_quotas | 直接对应 |
| 确定性 shuffle / partition | §3.1.1 | training.stable_partition / deterministic_shuffle_key | 直接对应 |
| 数据治理闭环(图 8) | §3.2.4 / 图 8 | profile 可插拔 + 每步 counters + reasons/flags + --keep-dropped | 代码是该闭环的工程载体 |
| AI 自治等级(Data cleaning L3) | §6.2 / E.2.1 | AI 为每个数据集生成 adapter/profile,共享决策接口 | 形态一致 |
| 去污染 / 配额物化 / 索引写盘 / RL 采样 | §4.1.1 / §4.2 | —(由训练环境与训练框架承担) | 代码明确不覆盖 |
| 能力 | 谁负责 | 说明 |
|---|---|---|
| 大规模抽取 / OCR | 训练环境 | 代码只处理“已抽好文本”的记录;PDF/OCR 源是 record-level 清洗 |
| 集群、分片、worker、PVC | 训练环境 | balance_work_units / plan_row_microshards 只给规划原语 |
| 真实 tokenizer 与 chat template | 训练环境 | 代码提供 token_estimator 注入点与 TokenizedSample 接口 |
| 索引数据集写入 / loader 校验 | 训练环境 | manifest 记录数与输出一致,供 loader 对齐 |
| benchmark 去污染(8-gram) | 训练环境 | 代码无该算子 |
| 质量 evaluator 模型 | 模型资产 | 代码只消费 quality_score / quality_mean 等字段 |
| RL 数据采样 / reward | 训练框架 | 论文 §4.2 描述;代码未覆盖 |
(len+3)//4 只是可复现的近似;长度桶在真实 tokenizer 下可能位移。64 % bands == 0 且 threshold < bands 会在构造时报错。review_min 语义:名字像“送审线”,实际是丢弃线(below_drop_threshold);配置时不要按字面理解。DROP 不会再被后续算子“救回”;若希望某步先于质量过滤,应调整步骤顺序而非期待覆盖。1_2048… 与 Midtrain B16/B64/B256 用途不同,不要混用标签。data-process/;确认 Python ≥ 3.10。PYTHONPATH=src python3 -m unittest discover -s tests -v —— 跑通全部单元测试。python3 scripts/smoke_test.py —— 四个 CLI 端到端 + manifest 计数一致性。decision/reasons/buckets/metrics。clean_<name>_record(sample, context) -> Sample,只做该数据集特有的字段修复 / 过滤 / 指标,并返回 Step(...) 元组。datasets/__init__.py 注册 DATASET_PROFILES["<name>"] = <Profile>(<required source/category>, STEPS)。--dataset-profile <name> 运行;源族 / 类别不匹配会立即报错(严格校验)。context.config 键(如 gharchive_max_file_bytes),以便不改代码调参。examples/datasets/<name>.jsonl 并加一条 unittest,锁定预期决策。论文则告诉你第 3 步与第 5 步该填什么数值、为什么这样填,以及整套流程在 7.39B 规模下的最终效果(如 4.2× 预训练 time-to-loss、质量优先的 SFT 分档)。