Daft Common Crawl 实战教程:从网页语料到句子级文本嵌入的完整流水线
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
Common Crawl 是规模最大的开放网络数据集之一,包含超过 2500 亿个网页、横跨 18 年的抓取历史,自 2020 年以来已成为生成式 AI 训练数据的核心来源——GPT-3 等模型所用训练数据的绝大部分都来自 Common Crawl。本文以 Daft 仓库中的教程 docs/examples/common-crawl-daft-tutorial.md 为主体,完整演示如何用daft.datasets.common_crawl加载抓取数据、理解其 WARC/WET 数据模式、用 UDF 做语言过滤与句子切分,并用本地小模型批量生成文本嵌入。读完并照做后,你将掌握一条从原始网页语料到可入库(向量库/Parquet)嵌入数据的端到端 Daft 流水线。
一、环境准备:安装 Daft 与依赖
教程使用的所有依赖可以通过一条命令安装。daft[transformers]额外带上 Transformers 相关依赖,spacy用于后续的句子切分:
!pip install uv && uv pip install "daft[transformers]" spacy以下是教程用到的完整导入。注意PYTORCH_CUDA_ALLOC_CONF必须在导入 torch之前设置,用于缓解 GPU 显存碎片问题:
import os # We need to do this _before_ importing torch os.environ["PYTORCH_CUDA_ALLOC_CONF"] = "expandable_segments:True" import json from collections.abc import Iterator from datetime import datetime from typing import TypedDict import spacy import spacy.cli import torch from transformers import AutoConfig import daft from daft import col二、用daft.datasets.common_crawl加载抓取数据
Daft 通过内置函数daft.datasets.common_crawl提供对 Common Crawl 的访问,参数与底层行为在 daft/datasets/common_crawl.py 中完整定义。教程中的基本用法:
IN_AWS = False # Set this to `True` if you're running in the us-east-1 AWS region. # For Google Colab, this must be set to `False`. df = daft.datasets.common_crawl( crawl="CC-MAIN-2025-33", # The specific crawl to get. Crawls are listed on the # Common Crawl website. This one was crawled ~ Spring 2025. content="warc", # Options are "raw", "text", "metadata", "warc", "wet", "wat". # These change the names, types, and number of columns you'll get # and, importantly, what kind of data you'll be able to work with! num_files=1, # Optional number of files to fetch. None means get all files. # Fetching a single file is useful for rapid prototyping. in_aws=IN_AWS, # Let's Daft select the optimal download location! You **MUST** # set this to False if running outside of us-east-1 in AWS. )2.1 参数详解
结合 daft/datasets/common_crawl.py 的函数签名与文档字符串,各参数含义如下:
| 参数 | 类型/默认值 | 说明 |
|---|---|---|
crawl | str | 抓取批次标识,如"CC-MAIN-2025-33"(2025 年春季抓取) |
segment | str \| None,默认None | 指定 crawl 内的某个 segment(每个 crawl 被切分为 100 个 segment);None表示全部 segment |
content | 默认"raw" | 可选"raw"/"text"/"metadata"/"warc"/"wet"/"wat",分别对应 WARC 原始记录、WET 抽取文本、WAT 元数据 |
num_files | int \| None,默认None | 限制要处理的文件数量;None表示处理全部匹配文件,设为1非常适合快速原型验证 |
io_config | IOConfig \| None | 存储访问配置(如 S3 凭证、HuggingFace token) |
in_aws | bool,默认False | 已弃用,源码中标注将在 v0.9.0 移除,请改用source="s3" |
source | "s3" / "hf" / "http" \| None | 数据源选择;None时自动尝试 HuggingFace,找不到则回退 HTTP |
几个值得注意的实现细节:
num_files有正数校验:num_files <= 0会直接抛ValueError(见 daft/datasets/common_crawl.py)。in_aws弃用迁移:源码中当in_aws=True时会发出warnings.warn并自动转换为source="s3";若同时传了source,in_aws优先并告警。因此新代码建议直接写source="s3"(在 AWS us-east-1 内)或source="hf"(AWS 之外,最省成本)。
2.2 数据源解析:manifest 文件如何被定位
从源码结构看,daft.datasets.common_crawl并非直接硬编码文件列表,而是先读取各数据源的manifest 清单文件再解析出实际文件路径。_get_mainfest_path函数(daft/datasets/common_crawl.py)按数据源拼出清单地址:
| 数据源 | 清单路径 | 文件前缀 |
|---|---|---|
| S3 | s3://commoncrawl/crawl-data/{crawl}/{file_type}.paths.gz | s3://commoncrawl/ |
| HuggingFace(默认首选) | hf://buckets/commoncrawl/commoncrawl/crawl-data/{crawl}/{file_type}.paths.gz | hf://buckets/commoncrawl/commoncrawl/ |
| HTTPS | https://data.commoncrawl.org/crawl-data/{crawl}/{file_type}.paths.gz | https://data.commoncrawl.org/ |
_get_common_crawl_paths(daft/datasets/common_crawl.py)用daft.read_text读取清单,拼接前缀得到完整 URL,再按需做segment过滤(contains)和num_files截断(limit)。如果首选的 HuggingFace 源抛FileNotFoundError(比如该 crawl 尚未同步到 HF bucket),会自动回退到source="http"。
content参数通过一张映射表落到具体的文件类型(daft/datasets/common_crawl.py):
content_type_map = {"raw": "warc", "text": "wet", "metadata": "wat", "warc": "warc", "wet": "wet", "wat": "wat"}最后一步是调用read_warc(warc_paths, io_config=io_config)生成 DataFrame。
三、理解数据模式:WARC 全记录 vs 抽取文本
先看看content="warc"(即 raw WARC)里有什么:
# Show is a materializing operation in Daft. It will cause the Daft query to run. # It only computes the first 8 records and then prints them to STDOUT. df.show()WARC 文件信息量非常大,每一行对应一条 WARC 记录:
WARC-Record-ID:uuid,每条记录的唯一标识;WARC-Type:记录类型——是元数据?是发往网站的原始请求?还是网站的完整响应?(取值如warcinfo、request、response、metadata、conversion);WARC-Date:Timestamp,这条记录何时被抓取;Content-Length:int,响应体的长度;WARC-Identified-Payload-Type:已识别出的响应负载类型(MIME);warc_content:网站内容(或响应/元数据等)的原始字节;warc_headers:一个 JSON 对象字符串,包含Content-Type、WARC-Block-Digest、WARC-Refers-To、WARC-Target-URI(即被爬取的网站 URL)等。
这个模式与 daft/io/_warc.py 中read_warc声明的 schema 一致:WARC-Record-ID为DataType.uuid(),WARC-Target-URI/WARC-Type/WARC-Identified-Payload-Type/warc_headers为字符串,WARC-Date为纳秒级 UTC 时间戳,Content-Length为Int64,warc_content为Binary。
如果需要非常精细的 WARC 级细节(请求头、响应状态、块摘要等),就用content="warc"。但本教程只需要网页的文本,不需要 HTML,因此改用content="text"(WET 数据)——这也意味着要处理和下载的数据量更小。
3.1 只取文本:content="text"与try_decode
下面是content="text"数据的处理与预览。这里用col("X")构建"引用某列"的表达式:
df_sample = ( daft.datasets.common_crawl( crawl="CC-MAIN-2025-33", content="text", num_files=1, in_aws=IN_AWS, ) # we only care about website responses that were converted into # this `"text"` content we selected for in the Common Crawl data .where(col("WARC-Type") == "conversion") # try to decode the byte content as UTF-8 encoded text .with_column("warc_content", col("warc_content").try_decode("utf-8")) # failed decodes result in a `None` value -- we remove these records .drop_null(col("warc_content")) .select("WARC-Record-ID", "WARC-Target-URI", "WARC-Date", "Content-Length", "warc_content", "warc_headers") ) df_sample.show()两个要点:
WARC-Type == "conversion"过滤:WET 文件中,网页响应经过抽取转换后的记录类型为conversion,其他类型(warcinfo、request等)不是正文文本;try_decode("utf-8")容错解码:warc_content是 Binary 列,try_decode尝试按 UTF-8 解码为字符串;解码失败的行得到None而不是抛异常,随后用drop_null清掉这些记录。
预览结果可以看到,Common Crawl 包含多种语言的网页。为了聚焦,教程只处理英文页面(扩展成多语言留给读者练习)。
四、UDF 解析warc_headers:过滤英文页面
语言信息藏在warc_headers这个 JSON 字符串列的WARC-Identified-Content-Language字段里。用一个@daft.func标量 UDF 解析 JSON,再取值过滤:
WarcHeaders = TypedDict( "WarcHeaders", { "Content-Type": str, "WARC-Block-Digest": str, "WARC-Identified-Content-Language": str, "WARC-Refers-To": str, "WARC-Target-URI": str, }, ) @daft.func def json_load_warc_headers(x: str) -> WarcHeaders: return json.loads(x) df_lang = df_sample.with_column("warc_headers", json_load_warc_headers(col("warc_headers"))).with_column( "language", col("warc_headers").get("WARC-Identified-Content-Language") ) df_lang.select("warc_content", "language").show()这里能观察到两个重要现象:
eng表示英文(English);- 同一条记录可能带多个语言标记(
language是列表)。
为简化处理,只保留全部文本均为英文的记录:
df_lang = df_lang.where(col("language") == "eng") df_lang.select("warc_content", "language").show()五、嵌入流水线:为下游任务准备 Common Crawl
熟悉数据集之后,进入核心目标——生成文本嵌入(embedding)。嵌入是对数据(文本、图像、音频等)的数值向量表示,编码了语义信息,可用于语义检索、去重、多语言应用等场景;实践中常把嵌入连同标识性元数据一起存入向量数据库。
要生成嵌入,首先要把网页文本拆成有意义的片段(chunk)。文本是分层的:
Document → Sections → Paragraphs → Sentences → Words → Characters切分策略取决于用途:
- 句子级:大多数场景通用,尤其当文档结构不明确或不一致时(网页正是如此,所以本教程选它);
- 段落级:适合 RAG(检索增强生成)这类需要跨句保持上下文的场景;
- 章节级:适合结构清晰划分的长文档;
- 固定长度:实现简单,但可能在任意边界切断语义。
5.1 全局配置变量
教程把所有可调参数集中定义,并做正数/非空校验,方便复现和调参:
######## CONFIGURATION: Options ######## MAX_SEQ_LEN_SPACY: int = 1_000 # Maximum text length for sentence splitting. NLP_MODEL_NAME: str = "en_core_web_sm" # spaCy model for sentence detection CHUNKING_PARALLELISM: int = 4 # Parallel chunking processes MAX_SEQ_LEN_SENTENCE_TRANSFORMER: int = 1024 * 1 # Maximum text length for any individual embedding. EMBEDDING_MODEL_NAME: str = "Qwen/Qwen3-Embedding-0.6B" # Text embedding model EMBEDDING_BATCH_SIZE: int = 16 # Batch size for embeddings EMBEDDING_SIZE: int = AutoConfig.from_pretrained(EMBEDDING_MODEL_NAME).hidden_size ######## CONFIGURATION: Validation ######## if MAX_SEQ_LEN_SPACY <= 0: raise ValueError(f"MAX_SEQ_LEN_SPACY must be positive! {MAX_SEQ_LEN_SPACY=}") if len(NLP_MODEL_NAME) == 0: raise ValueError("NLP_MODEL_NAME must be specified!") if CHUNKING_PARALLELISM <= 0: raise ValueError(f"CHUNKING_PARALLELISM must be positive! {CHUNKING_PARALLELISM=}") if MAX_SEQ_LEN_SENTENCE_TRANSFORMER <= 0: raise ValueError(f"MAX_SEQ_LEN_SENTENCE_TRANSFORMER must be positive! {MAX_SEQ_LEN_SENTENCE_TRANSFORMER=}") if len(EMBEDDING_MODEL_NAME) == 0: raise ValueError("EMBEDDING_MODEL_NAME must be specified!") if EMBEDDING_BATCH_SIZE <= 0: raise ValueError(f"EMBEDDING_BATCH_SIZE must be positive! {EMBEDDING_BATCH_SIZE=}")各参数职责:MAX_SEQ_LEN_SPACY(1000 字符)限制送入 spaCy 做句子检测的文本长度;MAX_SEQ_LEN_SENTENCE_TRANSFORMER(1024)限制单个送入嵌入模型的句子长度;EMBEDDING_BATCH_SIZE(16)是 GPU 推理批大小;EMBEDDING_SIZE则直接从Qwen3-Embedding-0.6B的AutoConfig.hidden_size读取,保证下游声明的嵌入维度与模型真实输出一致。
5.2 下载 spaCy 模型
在切分之前先下载句子检测模型。注意要在主流程之外一次性下载,而不是在 UDF 内部下载——因为 Daft 可能对同一个 UDF 创建多个实例,多进程并发下载会产生竞态条件:
try: spacy.cli.download(NLP_MODEL_NAME) except: print(f"ERROR: Invalid spacy model name: {NLP_MODEL_NAME=}") raise5.3 句子切分 UDF(@daft.cls+ spaCy)
Daft 的类式 UDF 由@daft.cls装饰器定义(其参数与语义见 daft/udf/init.py)。教程使用max_concurrency=1, use_process=True:
use_process=True:让每个类实例在独立进程中运行,隔离 spaCy 的 GIL 竞争与 C 扩展;max_concurrency=1:限制该类 UDF 的并发实例数,模型加载昂贵,用类式 UDF 可以在多行数据间复用一次初始化。
class TextChunk(TypedDict): text: str chunk_id: int @daft.cls(max_concurrency=1, use_process=True) class ChunkingUDF: """Chunks text into sentences using Spacy.""" def __init__(self) -> None: # ensure model is already present via: # python -m spacy download {NLP_MODEL_NAME} # Or via Python: # spacy.cli.download(NLP_MODEL_NAME) # We **DON'T** download it here otherwise we could have a race # condition as Daft _can_ make multiple copies of our UDF. self.nlp = spacy.load(NLP_MODEL_NAME) @daft.method def __call__(self, text: str) -> Iterator[TextChunk]: n_truncated_spacy = 0 n_truncated_sentence_transformer = 0 if len(text) > MAX_SEQ_LEN_SPACY: n_truncated_spacy += 1 text = text[:MAX_SEQ_LEN_SPACY] doc = self.nlp(text) for i, sentence in enumerate(doc.sents): if len(sentence.text) > MAX_SEQ_LEN_SENTENCE_TRANSFORMER: s_text = sentence.text[:MAX_SEQ_LEN_SENTENCE_TRANSFORMER] n_truncated_sentence_transformer += 1 else: s_text = sentence.text chunked_text = TextChunk(text=s_text, chunk_id=i) yield chunked_text if n_truncated_spacy > 0: print(f"Truncated {n_truncated_spacy} sentences that were longer than {MAX_SEQ_LEN_SPACY} characters.") if n_truncated_sentence_transformer > 0: print( f"Truncated {n_truncated_sentence_transformer} sentences that were longer than {MAX_SEQ_LEN_SENTENCE_TRANSFORMER} characters." )这个方法的关键点:
- 用
Iterator[TextChunk]一行产出多行:一个网页会被展开为多个句子,每条带自增chunk_id; - 两级截断保护:先截到
MAX_SEQ_LEN_SPACY再进 spaCy,切出的每个句子若超过MAX_SEQ_LEN_SENTENCE_TRANSFORMER也截断,并对截断数量打印统计,方便监控数据质量。
在小样本上看效果:
chunker = ChunkingUDF() df_chunk = ( df_lang.with_column("chunks", chunker(col("warc_content"))) # and we want to see each object's fields as their own columns .with_column("text_chunk", col("chunks").get("text")) .with_column("text_index", col("chunks").get("chunk_id")) .select("WARC-Record-ID", "WARC-Target-URI", "WARC-Date", "text_chunk", "text_index") ) df_chunk.show()切分后,chunks是一个 Struct 列(text+chunk_id),用.get("text")/.get("chunk_id")拆成独立列。
5.4 批量文本嵌入 UDF(sentence-transformers + Qwen3)
有了句子级文本,就可以生成嵌入了。Daft 让"在数据上跑模型"变得非常直接:再写一个类式 UDF,用本地运行的sentence-transformers模型计算嵌入。嵌入列的类型显式声明为daft.DataType.embedding(float32, EMBEDDING_SIZE),这让 Daft 在计划阶段就知道输出是固定维度的向量:
from sentence_transformers import SentenceTransformer @daft.cls(max_concurrency=1, use_process=True) class EmbedderUDF: def __init__(self): self.device = "cuda" if torch.cuda.is_available() else "cpu" self.model = SentenceTransformer(EMBEDDING_MODEL_NAME).to(self.device) self.model = self.model.eval() self.model.compile() @daft.method.batch( return_dtype=daft.DataType.embedding(daft.DataType.float32(), EMBEDDING_SIZE), batch_size=EMBEDDING_BATCH_SIZE, ) def embed_text(self, texts): with torch.inference_mode(): embeddings = self.model.encode( texts, batch_size=EMBEDDING_BATCH_SIZE, output_value="sentence_embedding", precision="float32", show_progress_bar=False, convert_to_numpy=True, ) return embeddings要点:
@daft.method.batch(batch_size=EMBEDDING_BATCH_SIZE):批式方法,Daft 会把batch_size行文本攒成一批交给embed_text,模型按批推理,GPU 利用率远高于逐行调用;torch.inference_mode()关闭梯度计算;self.model.compile()启用 torch.compile 加速;设备自动选择 CUDA/CPU;return_dtype声明嵌入向量维度,与AutoConfig...hidden_size动态读取的EMBEDDING_SIZE呼应,无需手填魔法数字。
在小样本上预览嵌入列:
( df_chunk.with_column("embedding", EmbedderUDF().embed_text(col("text_chunk"))) .select("WARC-Record-ID", "WARC-Target-URI", "text_chunk", "text_index", "embedding") .show() )六、组装完整流水线
把各部分串起来,就是一条从 Common Crawl 原始数据到句子嵌入的完整查询计划:
chunker = ChunkingUDF() embedder = EmbedderUDF() df = ( daft.datasets.common_crawl( crawl="CC-MAIN-2025-33", segment=None, content="text", num_files=10, # INCREASE THIS NUMBER TO RUN ON MORE CRAWL FILES # OR REMOVE IT / SET IT TO `None` TO RUN ON ALL FILES! in_aws=IN_AWS, ) # only run on actual website text .where(col("WARC-Type") == "conversion") # UTF-8 decode the text .with_column("text", col("warc_content").try_decode("utf-8")) .drop_null(col("text")) # extract the language & filter english pages only .with_column("warc_headers", json_load_warc_headers(col("warc_headers"))) .with_column("language", col("warc_headers").get("WARC-Identified-Content-Language")) .where(col("language") == "eng") # chunk text into sentences .into_batches(batch_size=EMBEDDING_BATCH_SIZE * 10) .with_column("sentences", chunker(col("text"))) .with_column("text", col("sentences").get("text")) .with_column("chunk_id", col("sentences").get("chunk_id")) .exclude("sentences") # perform text embedding using the GPU .into_batches(batch_size=EMBEDDING_BATCH_SIZE) .with_column("embedding", embedder.embed_text(col("text"))) # our final columns .select("WARC-Record-ID", "WARC-Target-URI", "WARC-Date", "chunk_id", "text", "embedding") )其中两处into_batches是性能关键。into_batches(实现见 daft/dataframe/dataframe.py)按目标行数重新切分分区,其启发式是"攒够batch_size * 0.8行就发出一个批次",以处理效率优先、批次大小近似为准。教程的编排逻辑是:
- 切分前:
into_batches(EMBEDDING_BATCH_SIZE * 10)——给 spaCy 切分提供更大批次,减少 UDF 调用开销; - 嵌入前:
into_batches(EMBEDDING_BATCH_SIZE)——由于切分把一行网页膨胀成多行句子,这里重新按嵌入批大小组批,正好对齐@daft.method.batch的batch_size; select收尾,只保留WARC-Record-ID、WARC-Target-URI、WARC-Date、chunk_id、text、embedding六列,形成"句子文本 + 嵌入向量 + 溯源元数据"的最终形态。
预览整体结果:
df.show()七、把结果落地:write_parquet
真实场景中必须把产出保存下来。daft.DataFrame.write_parquet接受本地路径或 S3 key,Daft 会自动写分区文件:
start = datetime.now() output = df.write_parquet("./local_chunked_cc_text_and_embeddings") end = datetime.now() print(f"Complete! Took {end-start} -- Wrote output partitions:\n{output}")输出目录中的每个分区 Parquet 都带有嵌入列与溯源列(记录 ID、URL、抓取时间),可直接加载到向量数据库或后续的去重、检索流程中。
八、小结与延伸阅读
本教程串起了四个核心技术点,均可在 Daft 仓库中溯源验证:
| 环节 | API | 源码位置 |
|---|---|---|
| 数据接入 | daft.datasets.common_crawl | daft/datasets/common_crawl.py |
| WARC 解析 | daft.read_warc | daft/io/_warc.py |
| 标量/类式 UDF | @daft.func、@daft.cls、@daft.method.batch | daft/udf/init.py |
| 批次重分区 | DataFrame.into_batches | daft/dataframe/dataframe.py |
注意事项与适用前提:
- 该 API 处于 beta 阶段,可能随 Common Crawl 数据集演进而变化;
- 在 AWS us-east-1 内运行必须传 S3 源(
in_aws=True,新写法source="s3")以获得最优路径;在 AWS 之外运行 S3 会产生出口流量费,教程中IN_AWS必须设为False; in_aws参数已弃用、将在 v0.9.0 移除,新代码建议直接使用source="s3" / "hf" / "http"。
进一步可阅读 Common Crawl 数据源参考文档(三种数据源的访问方式、认证配置与 WARC/WET/WAT 加载差异)以及教程对应的 Jupyter 版本。
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考