构建多模态训练数据发布集
即使源资产包含文档、图像、音频和纯文本,训练数据发布也需要一个一致且可查询的 schema。本教程沿着 multimodal-training-data 用例,把四种模态转换成类型明确的特征表、发布数据集,以及可审查的拒绝记录集。
用例源码:AstroVela/demo-scene/multimodal-training-data。
默认资产清单包含五个固定版本的公开来源资产,其中四个通过发布策略,一个低分辨率 SVG 被拒绝。这个示例关注的是数据处理流程,而不是模型训练。Vane 让每种模态使用合适的 Python 批处理函数,再把所有分支合并回同一个类型明确的 Relation,统一执行发布策略、汇总和写出。
处理流程会:
- 把资产清单读取为一个源 Relation。
- 按模态把这个 Relation 过滤成多个分支。
- 使用类型明确的批量 UDF 解析每个分支。
- 用 union 把兼容的特征行合并成一张表。
- 使用 Relation 的筛选和聚合操作生成发布、拒绝与汇总产物。
默认资产与输入契约
training_assets.csv 中每个源资产一行。每行标识模态、源 URI、许可证、数据集划分、MIME 类型、本地文件路径、预期内容哈希、可选文本和模态元数据。
| 记录 | 模态与来源 | 用途 | 默认结果 |
|---|---|---|---|
| arrow-project-readme | Apache Arrow Markdown 文档 | 文档解码和行数统计 | 通过 |
| arrow-python-readme | Apache Arrow Markdown 文本 | 空白规范化和 token 统计 | 通过 |
| wikimedia-generic-file | Wikimedia 512×512 SVG | 通过分辨率策略的图像 | 通过 |
| wikimedia-download-icon | Wikimedia 136×168 SVG | 低分辨率检查 | 拒绝 |
| wikimedia-audio | Wikimedia 2.4 秒 WAV | 时长和采样率抽取 | 通过 |
这些文件及其来源和许可证元数据都固定在用例中,因此正常流程完全离线。这个仓库内资产快照是一组具体的发布示例数据,并不表示五个文件本身构成有价值的训练集。
1. 为所有模态定义一个特征契约
各个处理器可以产生不同的媒体指标和特征 JSON,但 build_feature_row 负责发布数据共享的列。它记录内容标识、大小、文本、token 数、质量、决策、风险标记、指标以及模态专属特征。
发布策略在这个边界上清晰可见:缺少许可证会成为风险标记;只有没有任何标记且质量分至少为 0.8 的记录才能通过发布策略。
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 秒 |
图像分支完整展示了模态边界:它在本地生成媒体指标和风险标记,再通过统一发布契约返回。
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。
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 操作。发布记录按数据集划分和模态排序;拒绝记录按审查顺序排列;汇总表按模态统计质量和决策数量。
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")
默认汇总会为每个分支提供直接反馈:
| 模态 | 输入 | 平均质量 | 发布 | 拒绝 |
|---|---|---|---|---|
| audio | 1 | 1.000 | 1 | 0 |
| document | 1 | 1.000 | 1 | 0 |
| image | 2 | 0.750 | 1 | 1 |
| text | 1 | 1.000 | 1 | 0 |
被拒绝的记录是 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、写出逻辑和快照元数据请查看完整用例。