跳到主要内容
Vane Data / 教程

处理 Common Crawl 数据

Common Crawl 很适合搜索、检索和模型训练流程,但 WET 文件仍需经过解码、过滤,并被切分成适合模型处理的单元。本教程沿着 examples/common_crawl.py 的主线,从原始 WARC 结构记录一直走到页面、文本块与嵌入产物。

默认输入刻意保持小巧,并且适合离线运行。它使用与真实 WET 输入相同列结构的内置记录,因此以后切换到本地文件或 URL 时,转换阶段无需改变。

处理流程分为五个阶段:

  1. 加载样例记录、本地 WET/WARC 文件或 WET URL。
  2. 保留 conversion 记录并解码正文。
  3. 从 WARC 头部读取识别出的语言,只保留目标语言。
  4. 把每个页面切成有长度上限的句子块。
  5. 生成嵌入并写出便于检查的结果。

1. 选择数据源

数据源选择集中在一个函数中。真实数据分支会检查必需参数、解析 WET 字节,再构建与样例分支相同的 Relation schema。

example.py
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 表。

example.py
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 循环之外,可以让各阶段的数据契约更容易检查。

example.py
    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 保留源记录标识,并在每个页面内分配单调递增的块编号。

example.py
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 可以直接写出。

example.py
    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 解析实现请查看完整源码