前言
在大模型预训练、微调和 RAG 建库之前,原始网页通常要经过正文抽取、规范化、质量过滤、去重、模型标注和向量化。算法并不神秘,难点在于让不同类型的计算在大规模数据上稳定运行,并且在调整阈值、更换模型或新增字段时不用重写整条管线。
阿里云 EMR Serverless Spark 产品提供了针对 Serverless 环境集成和优化的 Daft。Daft 是面向多模态数据的分布式 DataFrame 引擎;EMR Serverless Daft 进一步提供文本清洗、媒体处理、文档解析等算子,并通过ai_query、ai_embedding将大模型生成与 embedding 服务调用纳入 DataFrame 执行计划。因此,“读 WARC → 抽取正文 → 评估质量 → 去重 → LLM 标注 → 生成向量 → 写入 OSS”可以在同一条 DataFrame 管线中表达。下表先展示这些能力在整条语料管线中的位置。
这条管线覆盖的产品能力
| 管线环节 | EMR Serverless Daft 能力 | 本文示例 | 主要产物 |
|---|---|---|---|
| 数据读取与正文抽取 | OSS 数据访问、WARC/HTML 解析 | CommonCrawlContentExtractor、HtmlTagRemover | URL、正文、来源文件 |
| 清洗与质量信号 | 批量文本算子、模型复用 | 空白规范化、正则替换、重复度与困惑度 | 清洗文本、质量特征 |
| 去重 | 哈希表达式、MinHash、shuffle 与图计算管线 | 精确去重 + MinHash/LSH 近似去重 | 唯一文档或代表样本 |
| 大模型调用 | 并发、批处理、限流退让与 token 统计 | ai_query | 分类、标签、JSON 抽取结果 |
| 向量化 | 批量 embedding 服务调用 | ai_embedding | 向量、token 用量、模型名 |
| 交付 | DataFrame 写出 OSS Parquet | write_parquet | 训练数据集或向量库导入文件 |
下面顺着一条具体管线走一遍:输入是 OSS 上的 Common Crawl WARC 文件,输出是可用于训练或向量检索的英文语料。示例阈值只是起点,生产环境应按语种、来源和下游任务用抽样与离线评估重新校准。
主流公开语料管线都在做什么
公开的大规模网页语料工程虽然细节不同,但普遍包含启发式过滤、去重、语种或模型质量评估等环节。
C4 对 Common Crawl 应用启发式清洗规则;Gopher 论文的 MassiveWeb 数据处理给出了文档长度、项目符号行、重复行和重复 n-gram 等过滤信号;RefinedWeb 展示了大规模网页数据过滤与去重的价值;CCNet 组合了文档去重、语种识别和基于维基百科语言模型的质量分档;FineWeb 与 FineWeb-Edu 则系统评估了过滤与去重策略,并用教育质量分类器筛选子集。
把这些做法叠在一起,就是一条六段式的流水线:
正文抽取:从 WARC / HTML 中提取人类可读的正文,去掉导航栏、广告、脚本、样式。
规范化:统一空白、脱敏、清理版权头和模板文本。
启发式规则过滤:用一组低成本的统计量过滤掉明显的低质量页面。
去重:精确去重去除完全重复,MinHash 近似去重识别 n-gram 高度重叠的近重复。
模型质量打分:用统计语言模型或分类器评估语言规范性、信息密度或领域适配度。
向量化与入库:检索路径产出向量并载入 Milvus 一类的向量数据库;训练路径直接把文本落成 Parquet。
这些公开工程提供的是方法参考,不是一组能无条件复用的通用阈值。同一条规则在英文新闻、中文论坛和代码文档上会呈现不同分布。真正消耗工程投入的,是在大规模数据上稳定执行、保留中间信号并快速迭代策略。
为什么选择 EMR Serverless Daft 来承载这条管线
规模问题。Common Crawl 完整数据集是 PB 级,WARC 文件中又包含大量 HTTP 响应记录。正文抽取是 CPU 密集型处理,单机脚本的处理时间会随数据规模快速增长。
生态割裂问题。正文抽取要用 trafilatura / jusText / goose3 这些 Python 库,质量打分要加载 fastText 模型,困惑度要加载 KenLM 和 SentencePiece,打标和向量化要调大模型服务,并发、限流、重试都得自己写。它们各有自己的依赖和加载成本,若简单封装为一个 UDF,模型会在每个 task 反复初始化。
容错问题。网页数据的质量参差不齐:截断的 HTML、错误的编码、伪装成 HTML 的二进制内容都可能出现。一条异常数据就可能导致整个 partition 重跑,损失数小时的计算进度。
组合与治理问题。每个统计量都需要稳定的输入输出类型、空值语义、模型路径和错误边界。如果每条管线都自行封装 UDF,算法版本、阈值和异常处理很容易分叉。
调度问题。清洗阶段是 CPU 密集、要高并发,打标和向量化这类大模型调用走网络、要控制并发和应对限流。写在同一个 DAG 里,资源如何分配?
EMR Serverless Daft 的内置算子库,解决的正是这一层工程问题。下面顺着管线看具体怎么写。
管线第一段:从 WARC 到一行一篇正文
假设 OSS 上有一批.warc.gz文件,可能是从 Common Crawl 同步下来的,也可能是自建爬虫落盘的。CommonCrawlContentExtractor负责把它们展开成结构化记录。
importdaftfromdaftimportcolfromdaft.emr.functionsimportCommonCrawlContentExtractor,emr_udf EXTRACT_CONCURRENCY=64# 示例值,请根据 WARC 文件数和集群可用 CPU 调整。docs=(daft.from_glob_path("oss://my-bucket/crawl/CC-MAIN-2025-33/*.warc.gz")# 显式拆分输入分区,使 WARC 文件可由多个 extractor 并行处理。.into_partitions(EXTRACT_CONCURRENCY).with_column("records",emr_udf(CommonCrawlContentExtractor,construct_args={"warc_src_type":"warc_url",# warc_binary / warc_url / warc_base64"extractor_type":"trafilatura",# trafilatura / justext / goose3},num_cpus=2,concurrency=EXTRACT_CONCURRENCY,batch_size=1,)(warc_files=col("path")),).explode("records").select(col("records").get("url").alias("url"),col("records").get("content").alias("raw_text"),col("records").get("warc_file").alias("warc_file"),).where(col("raw_text").not_null()))这个算子内部用warcio迭代 WARC 记录,只保留WARC-Type为response的条目,按 UTF-8 解码 HTTP 响应体(解码错误用替换字符兜底),再交给选定的抽取后端。一行输入(一个 WARC 文件)产出一个list[struct],struct 里有url、content、warc_file、extractor四个字段,explode之后就是一行一篇正文。其中warc_file记的是来源文件名(warc_src_type="warc_url"时取路径的 basename,其他输入模式下为空串),需要完整路径就自己保留一列原始path。
三个抽取后端的取舍值得说一句。trafilatura是该算子的默认值,RefinedWeb 也使用过它处理网页正文;justext基于停用词密度做段落级样板文本判别;goose3偏向文章型页面的主体抽取。不存在对所有站点都最优的后端,建议按域名和页面类型抽样,对比正文召回、模板残留和空结果率后再选型。
如果输入本来就是 HTML 列——站点快照、内部文档库、API 抓回来的页面——那就跳过 WARC 这层,直接用HtmlTagRemover:
fromdaft.emr.functionsimportHtmlTagRemover,emr_udf df=df.with_column("raw_text",emr_udf(HtmlTagRemover,construct_args={"separator":"\n","strip":True},num_cpus=1,concurrency=32,batch_size=512,)(texts=col("html")),)HtmlTagRemover的行为是围绕"给大模型看的纯文本"设计的,不是简单地把尖括号删掉:<script>/<style>/<head>/<title>/<noscript>/<template>连同内部文本整体剔除,注释和 DOCTYPE 一并忽略;<div>/<p>/<h1>到<h6>/<li>/<table>这类块级元素切分成独立文本段,段间用separator连接,传"\n"就是每个块独占一行;<b>/<strong>/<a>这类行内标签不产生切分,<p>First <b>bold</b> end</p>出来还是完整的一句话;<pre>内部的缩进和换行原样保留,只裁剪最外沿空白,代码块的原有排版不会被破坏。解析用的是宽容模式,标签没闭合、缺<html>/<body>包裹的碎片也能正常抽取。
管线第二段:清理噪声,同时保留行结构
正文抽取后,先做正则替换和版权注释清理,将结果保存为text_lines。这一列保留换行,专门用来计算重复行和项目符号行占比。然后再用WhitespaceNormalizer把 Unicode 空白字符(包括换行)转为 ASCII 空格,得到用于哈希、n-gram、模型打分和 embedding 的单行text。
[!warning] 不要先压平换行再计算行级信号
WhitespaceNormalizer会把\n也转为空格。如果先规范化再调用RepeatedLinesCalculator或BulletLineRatioCalculator,文档通常只剩一行,结果就失去意义。
fromdaft.emr.functionsimport(CopyrightCleaner,RegexReplacer,WhitespaceNormalizer,emr_udf,)cleaned=(docs.with_column("text_lines",emr_udf(RegexReplacer,construct_args={"patterns":[r"[\w.+-]+@[\w-]+\.[\w.]+",# 邮箱r"(?<!\d)\+?\d[\d\-\s]{7,}\d(?!\d)",# 电话号],"replacements":["<EMAIL>","<PHONE>"],},num_cpus=1,concurrency=32,batch_size=512,)(texts=col("raw_text")),).with_column("text_lines",emr_udf(CopyrightCleaner,num_cpus=1,concurrency=32,batch_size=512)(texts=col("text_lines")),).with_column("text",emr_udf(WhitespaceNormalizer,num_cpus=1,concurrency=32,batch_size=512)(texts=col("text_lines")),))RegexReplacer支持"多个模式共用一个替换串"的广播写法:replacements只给一个元素时会自动铺开到所有patterns。写错的正则会被跳过并输出警告,不会导致整个作业失败。
[!important] 清洗不等于完成合规治理 上面的邮箱和电话正则只是管线示例,不能覆盖所有个人信息。
CopyrightCleaner只根据标记清理特定注释,也不代表已获得内容授权。生产语料还需要独立的数据来源审核、隐私脱敏、内容安全、删除请求和许可证治理。
管线第三段:一次算出全部质量信号
这一段对应 Gopher 和 C4 的启发式规则。EMR Serverless Daft 把每一类统计量做成了独立算子,可以按需组合:
| 算子 | 实际输出 | 用途 | 示例阈值 |
|---|---|---|---|
TextLengthCalculator | Int64,Unicode 字符数 | 过滤过短或过长文本;与按“词数”过滤不等价 | 200 至 50 万 |
RepeatedLinesCalculator | Float64,重复行额外出现次数 / 非空行数 | 识别模板行、重复菜单与重复段落 | < 0.30 |
WordRepetitionCalculator | Float64,反复出现的 n-gram 实例占比 | 识别词组堆叠和模板重复 | 5-gram < 0.15 |
BulletLineRatioCalculator | Float64,项目符号行占比 | 识别菜单、导航或列表型页面 | < 0.90 |
UrlRatioCalculator | Float64,URL 字符覆盖率 | 识别链接密集页面 | < 0.20 |
AlphanumericRatioCalculator | Float64,Unicode 字母和数字的字符占比 | 识别符号或乱码密集文本;不是“含字母的词占比” | > 0.70 |
MaximumWordLengthCalculator | Int64,最长英文词长度 | 识别 base64、哈希串或乱码残片 | ≤ 40 |
fromdaft.emr.functionsimport(AlphanumericRatioCalculator,BulletLineRatioCalculator,MaximumWordLengthCalculator,RepeatedLinesCalculator,TextLengthCalculator,UrlRatioCalculator,WordRepetitionCalculator,emr_udf,)POOL=dict(num_cpus=1,concurrency=32,batch_size=1024)signals=(cleaned.with_column("n_chars",emr_udf(TextLengthCalculator,**POOL)(texts=col("text"))).with_column("dup_line_ratio",emr_udf(RepeatedLinesCalculator,**POOL)(texts=col("text_lines"))).with_column("dup_5gram_ratio",emr_udf(WordRepetitionCalculator,construct_args={"repetition":5,"lang":"en","tokenization":False},**POOL,)(texts=col("text")),).with_column("bullet_ratio",emr_udf(BulletLineRatioCalculator,**POOL)(texts=col("text_lines"))).with_column("url_ratio",emr_udf(UrlRatioCalculator,**POOL)(texts=col("text"))).with_column("alnum_ratio",emr_udf(AlphanumericRatioCalculator,**POOL)(texts=col("text"))).with_column("max_word_len",emr_udf(MaximumWordLengthCalculator,**POOL)(texts=col("text"))))candidate=signals.where((col("n_chars")>=200)&(col("n_chars")<=500_000)&(col("dup_line_ratio")<0.30)&(col("dup_5gram_ratio")<0.15)&(col("bullet_ratio")<0.90)&(col("url_ratio")<0.20)&(col("alnum_ratio")>0.70)&(col("max_word_len")<=40))生产环境更稳妥的做法是:先全量计算并写出质量信号,再单独应用过滤条件。先观察不同来源、语种和时间分区的特征分布,再结合人工抽样和下游效果校准阈值。上面的candidate只用于演示筛选表达式;如果要保留可解释性,应先把signals持久化。
WordRepetitionCalculator有个参数需要留意:tokenization=True时会加载 SentencePiece 模型分词,中文必须开;纯英文语料设tokenization=False直接按空白切分,省掉模型依赖,速度也更快。
管线第四段:先去重,再把成本花在唯一内容上
去重放在远程大模型调用和向量化之前,理由很直接:同一份内容不应重复消耗 token 和 embedding 请求。实际削减比例取决于 crawl 范围、时间跨度、规范化规则和近似阈值,不应在没有基准数据时承诺固定降幅。去重分两步:先精确去重,再用 MinHash + LSH 发现词汇层面的近似重复。
精确去重:对规范化后的文本做哈希,再按指纹 distinct。哈希要在规范化之后的text上算,而不是原始抽取文本——经过前面的空白归一、邮箱电话掩码和版权信息清理,只剩格式差异的同源副本会落到同一个指纹上。先哈希、再按指纹去重,也是 CCNet 和 FineWeb 采用的做法。Daft 的哈希表达式为每条文本算出一个 64 位整数指纹(默认 xxhash3),配合行级 distinct 两步完成:
exact=candidate.with_column("text_hash",col("text").hash()).distinct("text_hash")hash()用于快速分组,不是加密哈希。对审计或严格正确性要求较高的数据,可在相同哈希组内再比较规范化文本,并按时间、来源优先级或文本完整度显式选择代表行,而不是依赖distinct的任意保留结果。
精确去重消掉的是规范化后完全一致的副本。剩下的近重复可能只是替换了少量模板、广告或句子,可以用 MinHash 签名和 LSH 候选集识别。MinHash 估计的是 n-gram 集合相似度,不等价于语义去重,对大幅改写或跨语言转述未必有效。Daft 内置 MinHash 表达式:
sig=exact.with_column("min_hashes",col("text").minhash(num_hashes=128,ngram_size=5,seed=1,hash_function="xxhash",),)拿到签名后,还需要执行 LSH 分桶、候选对构建、连通分量和代表样本选择。这部分请参考 Daft 官方 MinHash 去重教程。为了让下文变量含义明确,约定deduped是完成近似去重后的 DataFrame;如果当前只需要精确去重,可以用下面的基线设置继续运行:
# 基线:暂不执行 LSH 近似去重# 生产管线中,将这一行替换为官方教程产出的代表样本 DataFrame。deduped=exact两步的顺序来自成本考虑:精确去重是哈希加 shuffle,通常比 MinHash 签名、LSH 分桶和连通分量更便宜。先用它降低候选集规模,可减少近似去重阶段的计算和 shuffle 开销。建议把每阶段的输入行数、输出行数、保留率和最大重复簇写入运行报告,用真实数据评估去重价值。
管线第五段:模型质量打分
规则能过滤明显的低质量内容,但过滤不了“语法正确却缺少信息量”的文本。这一步交给两个轻量模型。走到这里的数据已经完成去重,只需要对代表样本打分,可避免为重复内容反复支付推理成本。
EnTextQualityScorer使用 fastText 分类器kenhktsui/llm-data-textbook-quality-fasttext-classifier-v2,输出 0 到 2 之间的连续分数。它计算 Low(0) / Mid(1) / High(2) 三档的概率加权期望,因此可用于排序、分位切分或阈值过滤。分数越高,表示该分类器越倾向把文本判为教科书或科普式内容;0.5 可作为初始观察点,但不是通用质量标准。该模型使用 CPU 推理。
PerplexityCalculator走的是 CCNet 的路线:先用 SentencePiece 分词,再用在维基百科上训练的 KenLM n-gram 模型逐行打分,把各行的对数概率和词数累加后折算成整篇文档的一个困惑度值。困惑度越低,文本越接近维基百科那种规范书面语。它支持中文和英文两种语言。
fromdaft.emr.functionsimportEnTextQualityScorer,PerplexityCalculator,emr_udf scored=(deduped.with_column("quality_score",emr_udf(EnTextQualityScorer,construct_args={"model_path":"oss://my-bucket/models",# 首次自动下载到本地缓存"batch_size":64,},num_cpus=1,concurrency=32,batch_size=512,)(texts=col("text")),).with_column("ppl",emr_udf(PerplexityCalculator,construct_args={"lang":"en","model_path":"/opt/emr/models"},num_cpus=1,concurrency=32,batch_size=512,)(texts=col("text")),))high_quality=scored.where(col("quality_score").not_null()&col("ppl").not_null()&(col("quality_score")>0.5)&(col("ppl")<1000))两个分数看的是不同的信号:困惑度衡量文本对当前 KenLM 参考分布的接近程度,质量分反映 fastText 分类器学到的教育内容偏好。将两者一起保留便于分析,但不建议在未评估误杀样本前直接将两个阈值都视为强过滤条件。例如,公式、代码和专业术语密集文档可能产生较高困惑度,却仍然对特定训练任务有价值。
这里有个容易被忽略的部署细节:EnTextQualityScorer的model_path支持oss://前缀,初始化时会把指定模型文件下载到本地缓存并复用。PerplexityCalculator依赖的 KenLM 和 SentencePiece 资产从本地模型目录加载,默认基础路径为/opt/emr/models。
需要提醒的是,EnTextQualityScorer是英文分类器,多语言场景应该先做语种识别,只把英文文本送入该算子。EMR Serverless Daft 也提供基于 fastTextlid.176模型的LanguageRecognizer,可输出语言代码和置信度,用于在管线前段分流。
管线第六段:语义加工与向量化入库
到这一步语料已经干净了,剩下的是加工成下游能直接消费的形态。
ai_query:在 DataFrame 中批量做 LLM 分类、打标和抽取
ai_query接收每行的文本 prompt,也支持可选的图像或视频数据列。在语料管线中,常见用法包括领域分类、模板文本判定、安全标签、摘要和结构化信息抽取。模型可以通过 EMR Serverless 中的 AI 函数默认配置、service_name或显式model路由;实际可用模型取决于当前工作空间和模型服务配置。
领域打标用ai_query,把大模型调用当成一个普通的 DataFrame 表达式:
fromdaftimportlitfromdaft.emr.functionsimportai_query tagged=high_quality.with_column("ai_result",ai_query(lit("Read the passage and reply with JSON only, fields: ""domain (science / code / news / forum / commerce / other), ""is_boilerplate (true or false).\n\n")+col("text"),model="qwen3.6-plus",concurrency=8,batch_size=64,options={"temperature":0,"max_tokens":128},),).with_column("domain_json",col("ai_result").get("content"))ai_query返回结构体,其中包括content、reasoning_content、finish_reason、prompt_tokens、completion_tokens、total_tokens、cached_tokens、reasoning_tokens、model、id和error等字段。这些字段方便统计 token 与模型响应元数据。上例的domain_json仍是字符串,下游应执行 JSON Schema 校验,不要只凭 prompt 就假设输出永远合法。
当前服务调用在重试后仍失败时会使作业失败,不应假设每个失败都会被静默转换成行内error。生产环境应配合幂等输出路径、分区重跑和运行监控。
[!tip] 先分块,再调用模型 网页正文可能超过模型上下文窗口或 embedding 单次输入上限。在
ai_query和ai_embedding之前,应按标题、段落或 token 数分块,并保留document_id、chunk_id和来源 URL。
ai_embedding:把文本列批量转换为向量
ai_embedding使用 OpenAI-compatible embeddings 接口语义调用配置的 embedding 服务,并将多行打包为批量请求。它返回包含embedding、prompt_tokens、completion_tokens、total_tokens、model和error的结构体。空字符串和空值不会发送到服务,结果为空向量;非空请求在重试后仍失败时会抛错。
fromdaft.emr.functionsimportai_embedding final=(tagged.with_column("emb",ai_embedding(col("text"),concurrency=8,batch_size=256,embedding_batch_size=8,),).with_column("vector",col("emb").get("embedding")))final.write_parquet("oss://my-bucket/corpus/en-clean/")落盘形态由下游消费方式决定。预训练数据通常保留文本、来源、语种、质量分和过滤版本,不必预先生成向量;RAG 和语义检索场景则应保留 chunk 文本、向量、主键和可回溯元数据。Parquet 适合作为 OSS 上的中间交付格式,再根据目标向量数据库的导入规范生成索引。
内置算子的优势
写管线时,用户主要做业务层决策:用哪个抽取后端、计算哪些质量信号、阈值如何校准、调用哪个模型以及保留哪些可回溯字段。算子则统一承担批处理接口、资源声明、模型加载和服务并发等工程语义。
这种封装的重点,是把复杂的逻辑收敛进一个个简单的算子调用里。正文抽取算子内部处理 WARC 遍历、编码兜底和异常记录降级;质量打分算子内部处理模型加载、缓存分发和批量推理;管线代码里看到的只是一个表达式和几个业务参数。写代码的人面对的是业务——抽哪些、留哪些、怎么打分——而不是每一环节的实现细节。
错误边界会按算子类型区分:例如 WARC 内单条正文抽取失败可降级为空内容,多个文本评分算子会对行内无效数据返回空值;模型服务请求在重试耗尽后则会使作业失败。这种区分避免了一概吞错,也便于把数据质量问题与系统性问题分开处理。
远程模型算子还将行批次、单次请求打包数、作业内并发和服务请求窗口分开配置。因此,用户可以在不改业务表达式的前提下,根据模型服务限流和单条文本大小调整吞吐参数。
写在最后
语料清洗的每一个环节都有成熟的参考方法,真正的难点是让抽取、质量评估、去重、模型调用和向量化共享同一套可观测、可回溯、可迭代的数据流程。
阿里云 EMR Serverless Daft 的核心价值,是将批量非结构化数据处理与ai_query、ai_embedding这类模型服务调用组合在同一个 DataFrame 作业中。对于已将原始数据存放在 OSS,并希望以批处理方式构建大模型训练语料或 RAG 索引的团队,它提供了一条可组合的工程路径。