跳到主要内容

Vane Data:从多模态文件到可查询数据

· 阅读需 16 分钟

AI 应用往往需要处理 PDF、图片、音频和视频。进入数据管道前,需要对 PDF 拆页、对图片解码、对音频重采样,并沿视频时间轴展开,同时保留文件名、页码和帧号等定位信息。本文以这四类文件为例,介绍 Vane Data 如何以 Relation(多模态数据集)为核心,串联文件展开、批量处理、模型推理、查询和写出,把原始文件转化为可查询的数据。

阅读说明.

本文用概念性伪代码说明 Vane Data 如何组织多模态处理流程,省略 PyMuPDF、图片解码、Whisper 和 YOLO 等工具的具体实现。文末列出了相关的可运行示例。

Vane Data 将 PDF、图片、音频和视频处理为可查询数据

SQL 与 Python:两种 API,一条数据管道

在 Vane Data 中,文件内容、业务字段和处理结果都保存在 Relation 里。SQL 适合表达字段计算、筛选和 AI Function 调用;Python API 适合接入现有处理库、执行一对多展开,以及配置 Actor 和 GPU 等运行资源。

解析、解码和格式转换可以接入无状态 UDF;Whisper、YOLO 等不宜反复加载的模型则封装成有状态的 UDF,由 Actor 初始化一次,再连续处理多个批次。Prompt 和 Embedding 交给 AI Function。

下面四个示例分别展示 PDF、图片、音频和视频如何从原始文件变成可查询的 Relation,并概览各自的处理流程、主要产出和代码组织方式。

示例典型处理流程主要产出代码呈现
PDFPDF → 文本块 → Embedding / 语义字段每行一个文本块,保留来源、页码、块序号和文本,并增加 embeddingtopicschunk_summarySQL 读取文件并生成向量与语义字段;
Python flat_map 展开页面和文本块
图片图片 BLOB → 批量解码与检查 → 可用图片 → 视觉 Prompt → STRUCT 字段每行一张通过检查的图片,包含文件信息、宽高、中文摘要和模型自报置信度SQL 调用 inspect_image 完成检查与筛选,再用 AI_PROMPT 生成结构化描述;
Python 注册 UDF
音频音频 bytes → 解码与 16 kHz 重采样 → Whisper 输入特征 → 模型推理 → 转写文本每行一条音频,保留路径、语言和业务字段,并增加中文 transcriptionSQL CTE 串联批处理 UDF;
Python 注册无状态 UDF 和复用 Whisper 模型的 Actor
视频视频文件 → 帧 → 每帧检测结果 → 目标 → 裁剪图片每行一个检测目标,包含视频与帧定位信息、类别、置信度、边界框和裁剪后的 PNG BLOBPython 使用 VideoFrameSourcemap_batches 和 GPU Actor 完成帧读取、检测、目标展开与裁剪

PDF:拆成文本块,生成检索与语义字段

典型流程:PDF → 页面 → 文本块 → Embedding / 语义字段

PDF 案例处理的是一批文本型文档,目标是得到可检索的文本块。每个文本块都带有文件来源、页码和块序号,检索结果可以直接回到原文;向量用于相似度检索,主题和摘要则用于筛选和结果展示。

示例通过 read_blob 读入文件,PyMuPDF 按页提取文本,两次 Python flat_map 分别展开页面和文本块。切块大小与 Embedding 模型的 tokenizer 对齐并保留适度重叠,再通过 SQL AI_EMBED 生成向量字段。

Embedding 和 Prompt 可以处理同一份文本块 Relation,embedding 用于相似度检索,topicschunk_summary 用于筛选和结果展示。

example.py
import vane


# 以下占位函数分别代表 PyMuPDF 和基于 tokenizer 的辅助逻辑。
# 完整的 PDF 处理流程请参考文末链接的可运行示例。
EMBEDDING_MODEL = "sentence-transformers/all-MiniLM-L6-v2"
MAX_CHUNK_TOKENS = 240
CHUNK_TOKEN_OVERLAP = 32




# flat_map 回调:接收一条 PDF 记录,按页生成多条记录。
def extract_pages(row):
    for page_number, text in parse_pdf(bytes(row["content"])):
        yield {
            "source": row["source"],
            "size": row["size"],
            "page_number": page_number,
            "text": text,
        }




# flat_map 回调:接收一条页面记录,将其切分为按 tokenizer 对齐且相互重叠的文本块。
def split_chunks(row):
    for chunk_index, text in enumerate(
        split_text_by_tokens(
            row["text"],
            model=EMBEDDING_MODEL,
            max_tokens=MAX_CHUNK_TOKENS,
            overlap_tokens=CHUNK_TOKEN_OVERLAP,
        )
    ):
        yield {
            "source": row["source"],
            "size": row["size"],
            "page_number": row["page_number"],
            "chunk_index": chunk_index,
            "text": text,
        }




# flat_map 会改变行数,因此每个阶段都需要声明输出 schema。
PAGE_SCHEMA = {
    "source": vane.sqltypes.VARCHAR,
    "size": vane.sqltypes.BIGINT,
    "page_number": vane.sqltypes.INTEGER,
    "text": vane.sqltypes.VARCHAR,
}
CHUNK_SCHEMA = {
    **PAGE_SCHEMA,
    "chunk_index": vane.sqltypes.INTEGER,
}


# 同一个连接既执行 SQL 查询,也承接后续的 Python Relation 操作。
con = vane.connect()


# read_blob 将文件元数据和二进制内容读入同一个 Relation。
pdfs = con.sql("""
    SELECT filename AS source, size, content
    FROM read_blob('/data/pdfs/*.pdf')
""")


# 第一次 flat_map 将一条 PDF 记录展开为多条页面记录。
pages = pdfs.flat_map(extract_pages, schema=PAGE_SCHEMA)


# 第二次 flat_map 将一条页面记录展开为多条文本块记录。
chunks = pages.flat_map(split_chunks, schema=CHUNK_SCHEMA)


# 文本块形成 Relation 后切回 SQL,通过 AI_EMBED 增加向量字段。
embedded = chunks.query(
    "chunks",
    """
    SELECT
        *,
        -- 为每个文本块生成向量,同时保留定位字段。
        AI_EMBED(
            text,
            provider := 'transformers',
            model := 'sentence-transformers/all-MiniLM-L6-v2',
            options := struct_pack(
                device := 'cpu',
                batch_size := 32
            )
        ) AS embedding
    FROM chunks
    """,
)


# 在同一个 Relation 上继续调用 AI_PROMPT,生成可查询的语义字段。
result = embedded.query(
    "embedded",
    """
    WITH enriched AS (
        SELECT
            *,
            -- AI_PROMPT 返回 STRUCT(topics, summary),失败时返回 NULL。
            AI_PROMPT(
                text,
                return_format := json '{
                    "type": "object",
                    "properties": {
                        "topics": {
                            "type": "array",
                            "items": {"type": "string"}
                        },
                        "summary": {"type": "string"}
                    },
                    "required": ["topics", "summary"],
                    "additionalProperties": false
                }',
                system_message :=
                    'Extract the main topics from the text chunk and summarize it in one sentence.',
                provider := 'vllm',
                model := 'Qwen/Qwen2.5-7B-Instruct',
                on_error := 'ignore',
                options := struct_pack(
                    max_tokens := 128,
                    temperature := 0.0
                )
            ) AS metadata
        FROM embedded
    )
    SELECT
        source,
        size,
        page_number,
        chunk_index,
        text,
        embedding,
        -- 直接展开 STRUCT,供后续 SQL 筛选和展示。
        metadata.topics AS topics,
        metadata.summary AS chunk_summary
    FROM enriched
    """,
)


# write_parquet 触发整条 Relation 流水线执行。
result.write_parquet("/tmp/pdf_chunks_enriched.parquet")

write_parquet 触发整条 Relation 执行。示例仅覆盖文本型 PDF;扫描件、密码保护或损坏文件需要另行分流。

输出中每行对应一个文本块,字段如下:

字段类型用途
sourceVARCHAR定位原始 PDF
sizeBIGINT保留原始文件大小
page_numberINTEGER定位原始页面
chunk_indexINTEGER表示文本块在页面内的顺序
textVARCHAR保存送入模型的文本
embeddingFLOAT[384]供下游向量索引和相似度查询使用
topicsVARCHAR[]用于筛选文本块的主题列表
chunk_summaryVARCHAR用于结果展示的文本块摘要

图片:筛选可用图片,生成结构化描述

典型流程:图片 BLOB → 批量解码与检查 → 可用图片 → 视觉 Prompt → STRUCT 字段

图片部分关注结构化描述。图片检查逻辑由 Python 实现,会预先注册为 SQL UDF。SQL 调用该函数得到包含宽度、高度和 is_usable 标记的 STRUCT,再通过 WHERE 分流图片。检查失败的记录可以形成单独的 Relation;通过检查的图片继续交给视觉模型。AI_PROMPT 按 schema 生成 STRUCT(summary, model_confidence),SQL 负责约束字段和类型。

示例处理五张图片,并通过 OpenAI Provider 调用 gpt-4o-mini。代码中的 inspect_image 定义批量 UDF 的输入输出约定,底层图片解码逻辑由 inspect_image_blobs 实现。

example.py
import pyarrow as pa
import vane


# UDF 为每张图片返回一个 STRUCT,三个字段都可直接在 SQL 中访问。
INSPECTION_TYPE = pa.struct([
    pa.field("width", pa.int32()),
    pa.field("height", pa.int32()),
    pa.field("is_usable", pa.bool_()),
])


# 定义无状态批量 UDF;Vane 每批传入两个 Arrow Array。
# 函数返回等长的 StructArray,batch_size 控制单批大小。
@vane.func.batch(return_dtype=INSPECTION_TYPE, batch_size=32)
def inspect_image(image, minimum_side):
    return inspect_image_blobs(image, minimum_side)


con = vane.connect()


# 将 Python UDF 注册到当前连接,SQL 函数名为 inspect_image。
# parameters 声明函数签名为 (BLOB, INTEGER)。
vane.attach_function(
    inspect_image,
    alias="inspect_image",
    connection=con,
    parameters=["BLOB", "INTEGER"],
)


# 将图片文件读入包含元数据和 BLOB 值的 Relation。
images = con.sql("""
    SELECT filename, size, content AS image
    FROM read_blob('/data/images/*')
""")


# 在 SQL 投影中调用已注册的 Python UDF;64 表示最短边阈值。
inspected = images.query(
    "images",
    """
    SELECT
        *,
        -- 每行得到一个 STRUCT(width, height, is_usable)。
        inspect_image(image, 64) AS inspection
    FROM images
    """,
)


# 在 SQL 中按 is_usable 分流,只把可用图片交给模型。
result = inspected.query(
    "inspected",
    """
    WITH described AS (
        SELECT
            filename,
            size,
            inspection.width AS width,
            inspection.height AS height,
            -- 对通过检查的图片调用视觉 Prompt,并要求返回固定结构的 STRUCT。
            AI_PROMPT(
                'Describe the main subject in the image and provide a confidence score from 0 to 1.',
                image,
                return_format := json '{
                    "type": "object",
                    "properties": {
                        "summary": {"type": "string"},
                        "model_confidence": {"type": "number"}
                    },
                    "required": ["summary", "model_confidence"],
                    "additionalProperties": false
                }',
                system_message := 'Return only a Chinese-language result that matches the requested structure.',
                provider := 'openai',
                model := 'gpt-4o-mini',
                on_error := 'ignore',
                options := struct_pack(
                    use_chat_completions := true,
                    max_output_tokens := 128,
                    temperature := 0.0
                )
            ) AS answer
        FROM inspected
        WHERE inspection.is_usable
    )
    SELECT
        filename,
        size,
        width,
        height,
        -- 将模型返回的 STRUCT 展开为普通的可查询列。
        answer.summary AS summary,
        answer.model_confidence AS model_confidence
    FROM described
    """,
)


# 写出时执行前面的 UDF 检查和模型调用。
result.write_parquet("/tmp/image_descriptions.parquet")

五张图片的运行结果如下:

文件名(filename内容摘要(summary模型自报置信度(model_confidence
n02094114_4707.JPEG一只毛茸茸的小狗正在草地上奔跑0.95
n02398521_13903.JPEGNULLNULL
n01784675_1352.JPEG蜈蚣的特写画面,可以看到分节的身体和细长的足0.95
n02790996_10925.JPEG一名男子正在健身房进行卧推0.95
n02018207_15713.JPEGNULLNULL

NULL 表示图片通过了解码和尺寸检查,但 Prompt 阶段没有生成有效的结构化结果。原因可能是 Provider 调用失败,也可能是输出未通过校验;它不表示图片不可用。model_confidence 可用于辅助排序。

音频:分阶段转写并复用 Whisper 模型

典型流程:音频 bytes → 解码与 16 kHz 重采样 → Whisper 输入特征 → 模型推理 → 转写文本

这里以普通话音频为例。音频样本的编码格式、采样率和时长并不统一,送入 Whisper 前需要先解码并重采样为 16 kHz。查询通过一组 SQL CTE 和最终投影显式展开各个阶段:先提取音频 BLOB,再依次完成解码与重采样、输入特征生成、Whisper 推理,以及把输出 token 解码成中文文本;后四个阶段分别调用已经注册的批处理 UDF。模型推理时显式指定 language="zh"task="transcribe",避免短音频依赖自动语言识别。

解码、特征生成和 token 解码由无状态批量 UDF 完成;WhisperTranscriber Actor 为 whisper_transcribe_zh SQL UDF 提供推理能力,并在生命周期内复用同一模型。SQL 负责串联各个阶段。

example.py
import vane




MODEL_ID = "openai/whisper-tiny"
SAMPLING_RATE = 16000
BATCH_SIZE = 128
NUM_GPU_ACTORS = 1


# 声明三个中间结果的类型,同时用于 UDF 返回值和 SQL 参数。
FEATURE_MELS = 80
FEATURE_FRAMES = 3000
RESAMPLED_AUDIO_TYPE = vane.list_type(vane.sqltypes.FLOAT)
INPUT_FEATURES_TYPE = vane.tensor_type(
    vane.sqltypes.FLOAT,
    (FEATURE_MELS, FEATURE_FRAMES),
)
TOKEN_IDS_TYPE = vane.list_type(vane.sqltypes.INTEGER)




# 以下占位函数分别代表音频解码、特征构建和 token 解码。
# 完整的音频处理流程请参考文末链接的可运行示例。


# 第一个无状态批量 UDF:将不同格式的音频 BLOB 解码并重采样为 16 kHz 波形。
@vane.func.batch(return_dtype=RESAMPLED_AUDIO_TYPE, batch_size=BATCH_SIZE)
def decode_resample_16k(audio_bytes):
    return decode_and_resample(audio_bytes, sample_rate=SAMPLING_RATE)




# 第二个无状态批量 UDF:将波形转换为 Whisper 要求的定长特征。
@vane.func.batch(return_dtype=INPUT_FEATURES_TYPE, batch_size=BATCH_SIZE)
def prepare_whisper_features(waveform):
    return build_whisper_features(
        waveform,
        model=MODEL_ID,
        sample_rate=SAMPLING_RATE,
    )




# 定义有状态 Actor UDF;每个 Actor 使用一张 GPU,并连续处理多个批次。
@vane.cls.batch(
    actor_number=NUM_GPU_ACTORS,
    gpus=1.0,
    return_dtype=TOKEN_IDS_TYPE,
    batch_size=BATCH_SIZE,
)
class WhisperTranscriber:
    def __init__(self):
        # Actor 创建时只执行一次 __init__,避免每批重复加载模型。
        self.model = load_whisper(MODEL_ID, device="cuda")


    def __call__(self, input_features):
        # __call__ 处理一个特征批次,并显式指定中文转写。
        return transcribe_to_token_ids(
            input_features,
            model=self.model,
            language="zh",
            task="transcribe",
        )




# 第三个无状态批量 UDF:将模型输出的 token ID 解码为中文字符串。
@vane.func.batch(return_dtype=vane.sqltypes.VARCHAR, batch_size=BATCH_SIZE)
def decode_whisper_tokens(token_ids):
    return decode_token_ids(token_ids, model=MODEL_ID)




con = vane.connect()


# 注册解码与重采样 UDF,对应 SQL 签名 decode_resample_16k(BLOB)。
vane.attach_function(
    decode_resample_16k,
    alias="decode_resample_16k",
    connection=con,
    parameters=["BLOB"],
)


# 注册特征生成 UDF,其输入类型与上一阶段的返回类型一致。
vane.attach_function(
    prepare_whisper_features,
    alias="prepare_whisper_features",
    connection=con,
    parameters=[RESAMPLED_AUDIO_TYPE],
)


# 注册 Actor 实例;后续 SQL 调用都会复用已经加载的模型。
vane.attach_function(
    WhisperTranscriber(),
    alias="whisper_transcribe_zh",
    connection=con,
    parameters=[INPUT_FEATURES_TYPE],
)


# 注册 token 解码 UDF,将 INTEGER[] 转换为 VARCHAR。
vane.attach_function(
    decode_whisper_tokens,
    alias="decode_whisper_tokens",
    connection=con,
    parameters=[TOKEN_IDS_TYPE],
)


# 读取包含 audio STRUCT 和样本元数据的 Parquet Relation。
source = con.sql("""
    SELECT *
    FROM read_parquet('/data/chinese-speech/*.parquet')
""")


# 通过分阶段的 SQL CTE 和最终投影串联已注册的 UDF。
result = source.query(
    "source",
    """
    WITH audio AS (
        SELECT
            * EXCLUDE (audio),
            -- 从 audio STRUCT 中提取原始音频 BLOB。
            audio.bytes AS audio_bytes
        FROM source
    ),
    resampled AS (
        SELECT
            * EXCLUDE (audio_bytes),
            -- 阶段 1:调用 Python UDF,将音频解码并重采样为 16 kHz。
            decode_resample_16k(audio_bytes) AS waveform
        FROM audio
    ),
    featured AS (
        SELECT
            * EXCLUDE (waveform),
            -- 阶段 2:调用 Python UDF,生成 Whisper 输入张量。
            prepare_whisper_features(waveform) AS input_features
        FROM resampled
    ),
    tokens AS (
        SELECT
            * EXCLUDE (input_features),
            -- 阶段 3:调用 Actor UDF,在 GPU 上执行模型推理。
            whisper_transcribe_zh(input_features) AS token_ids
        FROM featured
    )
    SELECT
        * EXCLUDE (token_ids),
        -- 阶段 4:调用 Python UDF,将 token ID 解码为最终文本。
        decode_whisper_tokens(token_ids) AS transcription
    FROM tokens
    """,
)


# 只写出业务字段和转写文本;写出操作会触发整条 UDF 链执行。
result.write_parquet("/tmp/chinese_audio_transcriptions.parquet")

视频:逐帧检测并提取目标

典型流程:视频文件 → 帧 Relation → 每帧检测结果 → 目标级 Relation → 裁剪图片

视频处理涉及自定义数据源、连续帧张量、GPU Actor 配置和批量目标裁剪,这些步骤使用 Python API 表达更加直接,因此这一节采用 Python First。

视频处理包含两次粒度变化:VideoFrameSource 先把一个视频沿时间轴展开成多帧,目标检测再把一帧中的多个目标展开成多条记录。source_idvideo_pathframe_index 始终随记录保留,因此检测结果既可以按目标字段查询,也可以回到来源视频。

Detector 由 Actor 复用 YOLO 模型;crop_object_batch 展开检测列表,并按边界框裁出 PNG。

example.py
import vane
from vane.datasource import read_datasource
from vane.datasource.video_reader import VideoFrameSource




# 有状态批量回调:传给 map_batches 后,由 Vane 创建并管理 Actor。
class Detector:
    def __init__(self):
        # 每个 Actor 在生命周期内只加载一次 YOLO 模型。
        self.model = load_yolo("yolo11n.pt", device="cuda")


    def __call__(self, table):
        # __call__ 接收一个 Arrow Table 批次,对其中所有视频帧执行目标检测。
        return run_object_detection(
            table,
            model=self.model,
            frame_column="frame",
            fields={
                "label": "boxes.cls",
                "confidence": "boxes.conf",
                "bbox": "boxes.xyxy",
            },
            # 保留来源字段和帧序号,方便定位检测结果。
            pass_through=(
                "source_id",
                "video_path",
                "frame_index",
                "frame",
            ),
        )




# 无状态批量回调:将每帧的目标列表展开为多条裁剪结果。
def crop_object_batch(table):
    # 一帧可能包含多个目标,因此需要展开为多行,并分别裁剪。
    return crop_detected_objects(
        table,
        frame_column="frame",
        features_column="features",
        pass_through=("source_id", "video_path", "frame_index"),
        object_column="object",
    )




con = vane.connect()


# VideoFrameSource 负责视频解码,read_datasource 将输出组织成帧 Relation。
# 每行表示一帧,并携带来源、路径、帧序号和帧张量。
frames = read_datasource(
    VideoFrameSource(
        video_paths,
        height=640,
        width=640,
    ),
    con=con,
)


# 第一次 map_batches 调用 Detector;Vane 创建 GPU Actor 并复用其中的模型。
detected = frames.map_batches(
    Detector,
    # schema 描述检测阶段生成的 Relation 中的全部字段。
    schema={
        "source_id": vane.sqltypes.VARCHAR,
        "video_path": vane.sqltypes.VARCHAR,
        "frame_index": vane.sqltypes.BIGINT,
        "frame": FRAME_TYPE,
        "features": FEATURE_LIST_TYPE,
    },
    batch_size=16,
    actor_number=1,
    gpus=1.0,
)


# 第二次 map_batches 在 CPU 上调用无状态函数,展开并裁剪目标。
objects = detected.map_batches(
    crop_object_batch,
    # 展开后每行表示一个目标,并增加裁剪后的 PNG BLOB。
    schema={
        "source_id": vane.sqltypes.VARCHAR,
        "video_path": vane.sqltypes.VARCHAR,
        "frame_index": vane.sqltypes.BIGINT,
        "features": FEATURE_TYPE,
        "object": vane.sqltypes.BLOB,
    },
)


# 显式选择最终输出字段;裁剪阶段已经移除帧张量。
result = objects.project(
    "source_id, video_path, frame_index, features, object"
)


# write_parquet 触发前面的读取、检测和裁剪阶段执行。
result.write_parquet("/tmp/video_objects.parquet")

输出的每一行对应一个检测目标:

字段用途
source_idvideo_path标识目标来自哪个视频
frame_index定位目标所在的解码帧
features.label模型输出的数值类别
features.confidence目标检测置信度
features.bbox检测模型输入帧坐标系中的边界框
object按 bbox 裁出的 PNG BLOB

features.label 可按 YOLO 的 names 映射为类别名。bbox 基于 640×640 的模型输入帧;映射回原视频时,还需保留原始宽高和缩放参数。frame_index 只是解码帧序号,若要精确定位播放时间,还需记录时间戳,或 PTS 与 time base。

结束语

摘要、转写和检测结果要继续用于查询、复核和写出,就不能和原始文件脱节。文件名、页码、帧号及业务字段随结果保留,查到某条记录时,仍能知道它来自哪里。Vane Data 用 Relation 维持这种对应关系,解析和推理继续交给现有的 Python 工具与模型。以后更换模型或加入业务规则,下游拿到的依然是 schema 清晰、来源明确的数据,无需重新拼接散落在脚本和中间文件里的结果。多模态文件由此进入日常的数据处理,不再停留在一次模型调用里。

开始使用: