处理 Common Crawl 数据
Common Crawl 很适合搜索、检索和模型训练流程,但 WET 文件仍需经过解码、过滤,并被切分成适合模型处理的单元。本教程沿着 examples/common_crawl.py 的主线,从原始 WARC 结构记录一直走到页面、文本块与嵌入产物。
默认输入刻意保持小巧,并且适合离线运行。它使用与真实 WET 输入相同列结构的内置记录,因此以后切换到本地文件或 URL 时,转换阶段无需改变。
处理流程分为五个阶段:
- 加载样例记录、本地 WET/WARC 文件或 WET URL。
- 保留 conversion 记录并解码正文。
- 从 WARC 头部读取识别出的语言,只保留目标语言。
- 把每个页面切成有长度上限的句子块。
- 生成嵌入并写出便于检查的结果。
1. 选择数据源
数据源选择集中在一个函数中。真实数据分支会检查必需参数、解析 WET 字节,再构建与样例分支相同的 Relation schema。
def load_source_relation(conn: Any, args: argparse.Namespace) -> Any: if args.source == "sample": return sample_relation(conn, args.limit) if args.source == "wet-file": if not args.wet_path: raise SystemExit("--wet-path is required when --source wet-file.") return wet_relation(conn, read_wet_file(args.wet_path, args.limit)) if args.source == "wet-url": if not args.wet_url: raise SystemExit("--wet-url is required when --source wet-url.") return wet_relation(conn, read_wet_url(args.wet_url, args.limit)) raise ValueError(f"Unsupported source: {args.source}")
解析器同时支持压缩和未压缩的 WET 输入。它会分离 WARC 头部与正文、提取标准记录元数据,并在达到请求条数后停止。
2. 解码并过滤页面
DecodeWarcBatch 是二进制 WARC 数据与类型化页面行之间的边界。它将正文按 UTF-8 解码,解析 JSON 形式的头部,并返回结构精简且明确的 Arrow 表。
class DecodeWarcBatch: """Decode WARC content bytes and parse WARC headers.""" def __call__(self, batch: pa.Table) -> pa.Table: record_ids = batch["WARC-Record-ID"].to_pylist() target_uris = batch["WARC-Target-URI"].to_pylist() dates = batch["WARC-Date"].to_pylist() lengths = batch["Content-Length"].to_pylist() content_values = batch["warc_content"].to_pylist() header_values = batch["warc_headers"].to_pylist() texts = [] languages = [] for content, raw_headers in zip(content_values, header_values, strict=True): try: text = bytes(content or b"").decode("utf-8") except UnicodeDecodeError: text = None try: headers = json.loads(raw_headers or "{}") except json.JSONDecodeError: headers = {} languages.append(str(headers.get("WARC-Identified-Content-Language") or "")) texts.append(text) return pa.table( { "record_id": pa.array(record_ids, type=pa.string()), "target_uri": pa.array(target_uris, type=pa.string()), "warc_date": pa.array(dates, type=pa.string()), "content_length": pa.array(lengths, type=pa.int64()), "language": pa.array(languages, type=pa.string()), "text": pa.array(texts, type=pa.string()), } )
编排逻辑先用 SQL 删除非 conversion 记录,再应用 UDF 并执行语言过滤。把这些过滤留在 Python 循环之外,可以让各阶段的数据契约更容易检查。
conn = vane.connect() rel = load_source_relation(conn, args) filtered = rel.query( "cc", """ select * from cc where "WARC-Type" = 'conversion' """, ) decoder = DecodeWarcBatch() pages = filtered.map_batches( decoder.__call__, schema={ "record_id": vane.sqltypes.VARCHAR, "target_uri": vane.sqltypes.VARCHAR, "warc_date": vane.sqltypes.VARCHAR, "content_length": vane.sqltypes.BIGINT, "language": vane.sqltypes.VARCHAR, "text": vane.sqltypes.VARCHAR, }, batch_size=args.batch_size, ).query( "pages", f""" select * from pages where text is not null and language = {sql_literal(args.language)} """, )
3. 生成适合嵌入的文本块
句子切分会先规范化空白,并优先使用标点边界。第二个辅助函数会在附近的空格处继续拆分过长句子。批量 UDF 保留源记录标识,并在每个页面内分配单调递增的块编号。
def regex_sentences(text: str) -> list[str]: normalized = re.sub(r"\s+", " ", text).strip() if not normalized: return [] pieces = re.split(r"(?<=[.!?])\s+", normalized) return [piece.strip() for piece in pieces if piece.strip()] def split_long_text(text: str, max_chars: int) -> list[str]: if len(text) <= max_chars: return [text] chunks = [] start = 0 while start < len(text): end = min(len(text), start + max_chars) if end < len(text): split_at = text.rfind(" ", start, end) if split_at > start + max_chars // 2: end = split_at chunk = text[start:end].strip() if chunk: chunks.append(chunk) start = end return chunks class ChunkTextBatch: """Split decoded web page text into embedding-sized chunks.""" def __init__(self, *, max_doc_chars: int, max_chunk_chars: int): self.max_doc_chars = max_doc_chars self.max_chunk_chars = max_chunk_chars def __call__(self, batch: pa.Table) -> pa.Table: output = { "record_id": [], "target_uri": [], "warc_date": [], "language": [], "chunk_id": [], "text": [], } rows = batch.to_pylist() for row in rows: text = str(row["text"] or "") if self.max_doc_chars and len(text) > self.max_doc_chars: text = text[: self.max_doc_chars] sentence_id = 0 for sentence in regex_sentences(text): for chunk in split_long_text(sentence, self.max_chunk_chars): output["record_id"].append(row["record_id"]) output["target_uri"].append(row["target_uri"]) output["warc_date"].append(row["warc_date"]) output["language"].append(row["language"]) output["chunk_id"].append(sentence_id) output["text"].append(chunk) sentence_id += 1 return pa.table( { "record_id": pa.array(output["record_id"], type=pa.string()), "target_uri": pa.array(output["target_uri"], type=pa.string()), "warc_date": pa.array(output["warc_date"], type=pa.string()), "language": pa.array(output["language"], type=pa.string()), "chunk_id": pa.array(output["chunk_id"], type=pa.int64()), "text": pa.array(output["text"], type=pa.string()), } )
默认文档上限为 1,000 个字符,默认文本块上限为 1,024 个字符。这些是适合教程的值,并非通用生产配置;应根据数据质量和嵌入模型进行调整。
4. 生成嵌入
示例使用 Vane 的 Transformers provider。Relation 形式的 embed 会保留源列并新增指定的 embedding 列,因此物化后的 Relation 可以直接写出。
embedded_table = None if not args.skip_embeddings: embedded = embed( chunks, vane.col("text"), provider="transformers", model=args.embedding_model_id, output_column="embedding", max_chunk_chars=args.max_chunk_chars, batch_size=args.embedding_batch_size, ) embedded_table = collect_relation(embedded)
模型默认为 Sentence Transformers 的 all-MiniLM-L6-v2。如果模型已经缓存,并且执行节点不应访问 Hugging Face,可以启用仅本地文件选项。也可以跳过此阶段,在不安装模型依赖的情况下检查 WARC 解码和切分结果。
5. 检查输出
默认输出目录包含:
- filtered_pages.csv:页面元数据与文本预览;
- chunks.csv:每个文本块一行;
- chunk_embeddings.parquet:启用嵌入时保存完整向量;
- chunk_embeddings.csv:用嵌入维度代替完整向量,便于快速检查。
脚本会打印源记录数、解码页面数、文本块数和嵌入维度,再展示有行数上限的 Relation 预览。如果没有页面匹配目标语言,或没有产生任何文本块,脚本会在写出结果前报错。
扩展同一条处理流程
脚本通过 Vane 配置的 runner 物化每个 Relation。可以设置 VANE_RUNNER=local 在本地验证 schema;使用 Ray 时则保持 runner 未设置并配置 RAY_ADDRESS。数据加载与输出契约保持不变。
所有参数与 WET 解析实现请查看完整源码。