跳到主要内容
Vane Data / 教程

构建多模态训练数据发布集

即使源资产包含文档、图像、音频和纯文本,训练数据发布也需要一个一致且可查询的 schema。本教程沿着 multimodal-training-data 用例,把四种模态转换成类型明确的特征表、发布数据集,以及可审查的拒绝记录集。

用例源码:AstroVela/demo-scene/multimodal-training-data

默认资产清单包含五个固定版本的公开来源资产,其中四个通过发布策略,一个低分辨率 SVG 被拒绝。这个示例关注的是数据处理流程,而不是模型训练。Vane 让每种模态使用合适的 Python 批处理函数,再把所有分支合并回同一个类型明确的 Relation,统一执行发布策略、汇总和写出。

处理流程会:

  1. 把资产清单读取为一个源 Relation。
  2. 按模态把这个 Relation 过滤成多个分支。
  3. 使用类型明确的批量 UDF 解析每个分支。
  4. union 把兼容的特征行合并成一张表。
  5. 使用 Relation 的筛选和聚合操作生成发布、拒绝与汇总产物。

默认资产与输入契约

training_assets.csv 中每个源资产一行。每行标识模态、源 URI、许可证、数据集划分、MIME 类型、本地文件路径、预期内容哈希、可选文本和模态元数据。

记录模态与来源用途默认结果
arrow-project-readmeApache Arrow Markdown 文档文档解码和行数统计通过
arrow-python-readmeApache Arrow Markdown 文本空白规范化和 token 统计通过
wikimedia-generic-fileWikimedia 512×512 SVG通过分辨率策略的图像通过
wikimedia-download-iconWikimedia 136×168 SVG低分辨率检查拒绝
wikimedia-audioWikimedia 2.4 秒 WAV时长和采样率抽取通过

这些文件及其来源和许可证元数据都固定在用例中,因此正常流程完全离线。这个仓库内资产快照是一组具体的发布示例数据,并不表示五个文件本身构成有价值的训练集。

1. 为所有模态定义一个特征契约

各个处理器可以产生不同的媒体指标和特征 JSON,但 build_feature_row 负责发布数据共享的列。它记录内容标识、大小、文本、token 数、质量、决策、风险标记、指标以及模态专属特征。

发布策略在这个边界上清晰可见:缺少许可证会成为风险标记;只有没有任何标记且质量分至少为 0.8 的记录才能通过发布策略。

example.py
def build_feature_row(
    row: dict[str, Any],
    *,
    payload: bytes,
    content_text: str,
    flags: list[str],
    quality_score: float,
    metrics: dict[str, int | float | None],
    features: dict[str, Any],
) -> dict[str, Any]:
    if not str(row.get("license_id") or "").strip():
        flags.append("missing_license")
    quality_score = round(max(0.0, quality_score - 0.25 * ("missing_license" in flags)), 3)
    decision = "accepted" if not flags and quality_score >= 0.8 else "rejected"
    return {
        "record_id": row["record_id"],
        "modality": row["modality"],
        "source_uri": row["source_uri"],
        "license_id": row.get("license_id") or "",
        "split": row["split"],
        "mime_type": row["mime_type"],
        "content_text": content_text,
        "content_sha256": hashlib.sha256(payload).hexdigest(),
        "byte_size": len(payload),
        "token_count": len(tokenize(content_text)),
        "quality_score": quality_score,
        "decision": decision,
        "risk_flags": flags,
        "media_metrics": metrics,
        "feature_json": json.dumps(features, sort_keys=True),
    }

在每个特征行中保留 decision 很有价值:完整特征表可以解释发布结果,而发布数据本身仍然只是一个简单的筛选 Relation。

2. 抽取模态专属特征

每个函数都通过同一个构建函数返回,因此媒体专属逻辑不会意外改变表格契约。

模态当前处理方式主要拒绝信号
document解码 UTF-8 并统计行数无效 UTF-8 或空内容
text解码 UTF-8、折叠空白并统计 token无效 UTF-8 或少于四个 token
image解析 SVG XML,并读取尺寸或 viewBox无效图像、缺少尺寸或低于 512×512
audio读取 PCM WAV 采样率和帧数无效音频或短于 0.005 秒

图像分支完整展示了模态边界:它在本地生成媒体指标和风险标记,再通过统一发布契约返回。

example.py
def process_image(row: dict[str, Any]) -> dict[str, Any]:
    payload = decode_payload(row)
    flags: list[str] = []
    metrics = empty_metrics()
    image_format = "unknown"
    if row["mime_type"] == "image/svg+xml":
        try:
            width, height = svg_dimensions(payload)
            image_format = "svg"
            metrics["width"] = width
            metrics["height"] = height
            if not width or not height:
                flags.append("missing_dimensions")
            elif width < 512 or height < 512:
                flags.append("low_resolution")
        except ET.ParseError:
            flags.append("invalid_image")
    else:
        flags.append("invalid_image")


    return build_feature_row(
        row,
        payload=payload,
        content_text=str(row.get("text") or ""),
        flags=flags,
        quality_score=1.0 - 0.5 * len(flags),
        metrics=metrics,
        features={
            "metadata": json.loads(row["metadata_json"] or "{}"),
            "format": image_format,
        },
    )

这些处理器刻意保持轻量。示例中可扩展的是边界:你可以用更强的 PDF 解析器、图像解码器、语音模型或语言检测器替换某个分支,同时继续返回相同的特征 schema。

3. 为每种模态运行一个带类型的分支

编排逻辑按模态过滤源 Relation,并分配对应的可导入批处理函数。map_batches 对每个分支应用共享 schema,union 再重建一个多模态特征 Relation。

example.py
def build_feature_relations(
    conn: Any,
    raw_assets: Any,
    args: argparse.Namespace,
) -> tuple[Any, dict[str, dict[str, Any]]]:
    udf_options = batch_udf_options(args.execution_backend)
    stage_functions = {
        "document": "process_document_batch",
        "image": "process_image_batch",
        "audio": "process_audio_batch",
        "text": "process_text_batch",
    }
    relations: list[Any] = []
    backend_metadata: dict[str, dict[str, Any]] = {}
    for modality in SUPPORTED_MODALITIES:
        source = raw_assets.filter(f"modality = '{modality}'").order("record_id")
        relations.append(
            source.map_batches(
                importable_batch_function(stage_functions[modality]),
                schema=TRAINING_FEATURE_SCHEMA,
                batch_size=args.batch_size,
                **udf_options,
            )
        )
        backend_metadata[f"process_{modality}"] = backend_metadata_entry(args.execution_backend)


    features = relations[0]
    for relation in relations[1:]:
        features = features.union(relation)
    return features, backend_metadata

执行后端会改变批处理函数的运行方式,但不会改变四个分支之后的筛选、合并或查询逻辑。

4. 将发布数据与审查集合分开

所有模态共享 schema 以后,发布就变成常规 Relation 操作。发布记录按数据集划分和模态排序;拒绝记录按审查顺序排列;汇总表按模态统计质量和决策数量。

example.py
    training_release_rel = feature_records.filter("decision = 'accepted'").order(
        "split, modality, record_id"
    )
    conn.sql("drop table if exists training_release")
    training_release_rel.to_table("training_release")
    training_release = conn.sql("select * from training_release")


    rejected_records_rel = feature_records.filter("decision = 'rejected'").order(
        "quality_score, modality, record_id"
    )
    conn.sql("drop table if exists rejected_records")
    rejected_records_rel.to_table("rejected_records")
    rejected_records = conn.sql("select * from rejected_records")


    modality_summary_rel = feature_records.aggregate(
        """
        modality,
        count(*) as records,
        sum(byte_size) as total_bytes,
        round(avg(quality_score), 3) as avg_quality_score,
        sum(case when decision = 'accepted' then 1 else 0 end) as accepted,
        sum(case when decision = 'rejected' then 1 else 0 end) as rejected
        """
    ).order("modality")

默认汇总会为每个分支提供直接反馈:

模态输入平均质量发布拒绝
audio11.00010
document11.00010
image20.75011
text11.00010

被拒绝的记录是 wikimedia-download-icon:实际尺寸 136×168 产生 low_resolution,质量分为 0.500

运行 VANE_RUNNER=local-fast .venv/bin/python src/multimodal_training_data.py 即可执行离线处理流程。它会报告五个原始资产、四条发布记录和一条拒绝记录。这里使用 local-fast 是有意的:项目启动检查要求 Vane 使用进程内 Relation 路径,因为命名表位于客户端连接中。它是这个离线 fixture 的实现要求,不是公共 local runner 的另一个名称。

输出用途
feature_records.parquet所有通过和未通过发布策略的特征行,包括风险标记和媒体指标
training_release.parquet四条通过示例发布策略的记录
rejected_records.csv被拒绝记录及其具体原因
modality_summary.csv按模态统计的数量、字节、平均质量和决策
manifest.json源模式、发布策略、schema、计数和执行后端

完整特征表与发布表同样重要:它让发布决策可以复现,也让被拒绝的数据能回到数据整理和修订流程,而不是在没有记录的情况下被丢弃。

适配真实训练资产

  • --input 指向另一个保留资产清单字段和逐资产许可证元数据的 manifest。
  • 可以用 PDF 解析、栅格图像解码、ASR 或更丰富的质量分析替换某个模态处理器,但仍返回 TRAINING_FEATURE_SCHEMA
  • 继续用共享列表示发布策略,这样新增媒体类型不需要另一条发布路径。
  • 将示例阈值视为需要版本管理的策略输入;它们是演示用检查条件,不是校准后的质量基准。

输入投影、公开来源资产清单、Arrow schema、写出逻辑和快照元数据请查看完整用例