Vane Data:从多模态文件到可查询数据
AI 应用往往需要处理 PDF、图片、音频和视频。进入数据管道前,需要对 PDF 拆页、对图片解码、对音频重采样,并沿视频时间轴展开,同时保留文件名、页码和帧号等定位信息。本文以这四类文件为例,介绍 Vane Data 如何以 Relation(多模态数据集)为核心,串联文件展开、批量处理、模型推理、查询和写出,把原始文件转化为可查询的数据。
本文用概念性伪代码说明 Vane Data 如何组织多模态处理流程,省略 PyMuPDF、图片解码、Whisper 和 YOLO 等工具的具体实现。文末列出了相关的可运行示例。
SQL 与 Python:两种 API,一条数据管道
在 Vane Data 中,文件内容、业务字段和处理结果都保存在 Relation 里。SQL 适合表达字段计算、筛选和 AI Function 调用;Python API 适合接入现有处理库、执行一对多展开,以及配置 Actor 和 GPU 等运行资源。
解析、解码和格式转换可以接入无状态 UDF;Whisper、YOLO 等不宜反复加载的模型则封装成有状态的 UDF,由 Actor 初始化一次,再连续处理多个批次。Prompt 和 Embedding 交给 AI Function。
下面四个示例分别展示 PDF、图片、音频和视频如何从原始文件变成可查询的 Relation,并概览各自的处理流程、主要产出和代码组织方式。
| 示例 | 典型处理流程 | 主要产出 | 代码呈现 |
|---|---|---|---|
| PDF → 文本块 → Embedding / 语义字段 | 每行一个文本块,保留来源、页码、块序号和文本,并增加 embedding、topics、chunk_summary | SQL 读取文件并生成向量与语义字段; Python flat_map 展开页面和文本块 | |
| 图片 | 图片 BLOB → 批量解码与检查 → 可用图片 → 视觉 Prompt → STRUCT 字段 | 每行一张通过检查的图片,包含文件信息、宽高、中文摘要和模型自报置信度 | SQL 调用 inspect_image 完成检查与筛选,再用 AI_PROMPT 生成结构化描述; Python 注册 UDF |
| 音频 | 音频 bytes → 解码与 16 kHz 重采样 → Whisper 输入特征 → 模型推理 → 转写文本 | 每行一条音频,保留路径、语言和业务字段,并增加中文 transcription | SQL CTE 串联批处理 UDF; Python 注册无状态 UDF 和复用 Whisper 模型的 Actor |
| 视频 | 视频文件 → 帧 → 每帧检测结果 → 目标 → 裁剪图片 | 每行一个检测目标,包含视频与帧定位信息、类别、置信度、边界框和裁剪后的 PNG BLOB | Python 使用 VideoFrameSource、map_batches 和 GPU Actor 完成帧读取、检测、目标展开与裁剪 |
PDF:拆成文本块,生成检索与语义字段
典型流程:PDF → 页面 → 文本块 → Embedding / 语义字段
PDF 案例处理的是一批文本型文档,目标是得到可检索的文本块。每个文本块都带有文件来源、页码和块序号,检索结果可以直接回到原文;向量用于相似度检索,主题和摘要则用于筛选和结果展示。
示例通过 read_blob 读入文件,PyMuPDF 按页提取文本,两次 Python flat_map 分别展开页面和文本块。切块大小与 Embedding 模型的 tokenizer 对齐并保留适度重叠,再通过 SQL AI_EMBED 生成向量字段。
Embedding 和 Prompt 可以处理同一份文本块 Relation,embedding 用于相似度检索,topics 和 chunk_summary 用于筛选和结果展示。
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;扫描件、密码保护或损坏文件需要另行分流。
输出中每行对应一个文本块,字段如下:
| 字段 | 类型 | 用途 |
|---|---|---|
| source | VARCHAR | 定位原始 PDF |
| size | BIGINT | 保留原始文件大小 |
| page_number | INTEGER | 定位原始页面 |
| chunk_index | INTEGER | 表示文本块在页面内的顺序 |
| text | VARCHAR | 保存送入模型的文本 |
| embedding | FLOAT[384] | 供下游向量索引和相似度查询使用 |
| topics | VARCHAR[] | 用于筛选文本块的主题列表 |
| chunk_summary | VARCHAR | 用于结果展示的文本块摘要 |
图片:筛选可用图片,生成结构化描述
典型流程:图片 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 实现。
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.JPEG | NULL | NULL |
| n01784675_1352.JPEG | 蜈蚣的特写画面,可以看到分节的身体和细长的足 | 0.95 |
| n02790996_10925.JPEG | 一名男子正在健身房进行卧推 | 0.95 |
| n02018207_15713.JPEG | NULL | NULL |
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 负责串联各个阶段。
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_id、video_path 和 frame_index 始终随记录保留,因此检测结果既可以按目标字段查询,也可以回到来源视频。
Detector 由 Actor 复用 YOLO 模型;crop_object_batch 展开检测列表,并按边界框裁出 PNG。
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_id、video_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 清晰、来源明确的数据,无需重新拼接散落在脚本和中间文件里的结果。多模态文件由此进入日常的数据处理,不再停留在一次模型调用里。
开始使用: