大语料下载器 本节摘要:训练语言模型远在第一次前向之前就开始了。语料要落到磁盘、解压、去重、可寻址,断点续传的故事得在网络断在 4% 之前就讲清。本节构建流式下载器:拉取压缩分片、用 Zstandard 实时解压、用 MinHash 加局部敏感哈希(LSH)给近重复文档指纹、写出流水线其余部分可信任的分片清单。两个失败要从第一分钟就为之设计:部分下载续传(HTTP Range 对已校验字节偏移)与重复文档移除(MinHash+LSH 在亚线性成本下抓近重复)。两者都有知名解,却常被跳过,因为流水线从一行 长出牙齿开始。 对应原课程:Phase 19 · Lesson 42 · (原英文 )。本节属「预训练/分布式」赛道第一节。
本节摘要:训练语言模型远在第一次前向之前就开始了。语料要落到磁盘、解压、去重、可寻址,断点续传的故事得在网络断在 4% 之前就讲清。本节构建流式下载器:拉取压缩分片、用 Zstandard 实时解压、用 MinHash 加局部敏感哈希(LSH)给近重复文档指纹、写出流水线其余部分可信任的分片清单。两个失败要从第一分钟就为之设计:部分下载续传(HTTP Range 对已校验字节偏移)与重复文档移除(MinHash+LSH 在亚线性成本下抓近重复)。两者都有知名解,却常被跳过,因为流水线从一行
requests.get长出牙齿开始。
对应原课程:Phase 19 · Lesson 42 ·
large-corpus-downloader(原英文phases/19-capstone-projects/42-large-corpus-downloader/docs/en.md)。本节属「预训练/分布式」赛道第一节。
阅读完本节,你应当能够:
urllib 流式拉远程分片,用 zstandard 解压,不把整文件缓冲进内存。Range 请求对已校验字节偏移做断点续传。第一次在 200 GB 语料上训练,网络断在 41%,脚本以 urllib 异常退出;第二次断在 78%;到 99% 你已重写循环三次。两个失败要从第一分钟为之设计。
服务器须支持 Range,客户端须对磁盘记录追踪已校验偏移,且偏移要扛过进程死亡。偏移与文件哪怕差一字节,续传写的就是垃圾,语料以一种只在 token 化时才暴露的方式损坏。
精确哈希去重漏近重复:同一篇维基百科带着三种样板页脚、同一份代码文件换许可头、同一篇博文每个链接带追踪参数。MinHash+LSH 在亚线性成本下抓这些:每文档一个签名,每签名一次桶查。
标准库 urllib.request.urlopen 返回类文件对象。包进 zstandard.ZstdDecompressor().stream_reader,字节从网络经解压器流入文档迭代器,从不在内存里实体化压缩分片或解压分片。唯一内存成本是行缓冲、当前文档的 MinHash 签名、LSH 索引。
下载器每分片写两个文件:分片本身与 .partial.json 检查点。检查点记 verified_bytes、expected_size、sha256_prefix(对前 verified_bytes 字节算的)、源 URL。启动时下载器读检查点、对磁盘字节重算 sha256_prefix,只在重算哈希匹配时续传。哈希错就把 partial 丢弃、从字节零重启。静默损坏不可能,因为校验的是字节而非假设。
code/main.py 实现:
download_shard(url, dest):流式 GET,写 .partial.json 检查点,Range 续传,sha256 校验已下载部分。stream_zstd_jsonl(path):经 zstd 解压迭代 JSONL 文档,不实体化全文件。MinHash(seed, num_perm):对文档 shingle 集合算 MinHash 签名。LSHIndex(num_bands, rows_per_band):把签名分桶,近重复碰撞。dedupe(documents, threshold):返回新文档与被丢弃判定。emit_manifest(shards):写含内容哈希、字节、文档数、去重判定的清单。MinHash + LSH 骨架:
class MinHash: def __init__(self, num_perm=128, seed=1): self.num_perm = num_perm # 每个置换一个独立哈希族(固定种子保证可复现) self.a = [randint(1, 2**61-1, seed=seed+i) for i in range(num_perm)] self.b = [randint(0, 2**61-1, seed=seed+num_perm+i) for i in range(num_perm)] self.sig = [float("inf")] * num_perm def update(self, shingles): for s in shingles: h = hash(s) & ((1<<61)-1) for i in range(self.num_perm): v = (self.a[i]*h + self.b[i]) % (1<<61) if v < self.sig[i]: self.sig[i] = v class LSHIndex: def __init__(self, num_perm=128, num_bands=32): self.bands = num_bands; self.rows = num_perm // num_bands self.buckets = [defaultdict(set) for _ in range(num_bands)] def insert(self, doc_id, sig): for b in range(self.bands): chunk = tuple(sig[b*self.rows:(b+1)*self.rows]) self.buckets[b][chunk].add(doc_id) # 同 chunk -> 候选近重复 def is_near_dup(self, sig): for b in range(self.bands): chunk = tuple(sig[b*self.rows:(b+1)*self.rows]) if self.buckets[b][chunk]: return True return False
设计要点:MinHash 的可复现性靠固定种子的哈希族;LSH 把
num_perm维签名切成num_bands段,同段同值即候选近重复。num_bands越多召回越高、精度越低;阈值与(bands, rows)的关系是 Jaccard 相似度的 S 曲线,调参即挪曲线陡峭点。续传的sha256_prefix校验是防静默损坏的关键——不校验而假设偏移正确,一次进程死亡就毁整份语料。
HuggingFace datasets 的 load_dataset 把下载解压缓存打包,但不做去重与断点续传的细粒度控制。Common Crawl 的 cc_net、EleutherAI 的 the-pile 脚本是工业级实现,核心仍是 MinHash+LSH(用 datasketch 库)+ 分片下载。本节手写让你看清续传的 Range 握手、MinHash 的签名、LSH 的桶。生产规模上,去重常在 Spark/Dask 上分布式跑,但单机签名逻辑与本节一致。Zstandard 是 2020 年后的事实压缩标准(比 gzip 解压快 10 倍、比率高),Common Crawl、The Pile v2、RedPajama 都用它。
code/main.py:download_shard、stream_zstd_jsonl、MinHash、LSHIndex、emit_manifest 均可复用。demo 用 mock 本地服务器(或文件分片)跑通流式下载、续传、去重、清单写出。清单是流水线下游(第 41 节 token 化)信任的产物——它告诉你每个分片的内容哈希、字节、文档数、去重判定,使下游能做全局偏移与校验。
.partial.json 的偏移或磁盘字节,确认续传检测哈希不匹配、从零重启。num_bands 与 rows_per_band,绘 Jaccard 相似度与碰撞概率的 S 曲线。urllib + zstandard.stream_reader,字节网络直达文档迭代器。.partial.json 记偏移与哈希,重启重算校验才续传。下一节,我们做「HDF5 token 化语料」——把下载的语料流式 token 化进可调整大小的 HDF5 整数数据集,供训练器按行速流式读取。