跳到主要内容

Vane Data:如何让 DuckDB 成为 AI 多模态数据引擎

· 阅读需 11 分钟

Vane Data 是基于 DuckDB 构建的多模态原生数据引擎,将数据处理、模型推理和异构资源调度纳入同一条 Relation 流水线。本文介绍它如何组合 AI Functions、Python UDF 与 vLLM,如何通过动态批处理、流水线并行、背压和容错让任务稳定运行,并以车险理赔为例展示从图片预处理到审核分流的完整流程。

AI 工作负载需要新的数据引擎

AI 工作负载正在从“查表”走向“理解和行动”。AI Agent 的检索、判断和执行依赖多种输入,包括文本、图片、音频、视频和文档。这些多模态数据需要先经过解析、筛选和组织,转化为可查询的结构化形式,才能稳定支撑 Agent 工作流。

  • 传统数据栈难以承载多模态工作流

多模态文件对象大小差异很大,单行可能从几字节到数百 MB;模型推理又常常需要 CPU、GPU、网络和存储,它们的吞吐并不一致。传统数据栈通常把数仓、Python 脚本、对象存储和模型服务拼在一起:表引擎负责结构化数据,脚本负责媒体预处理,另一个系统负责推理。数据在系统之间来回搬运,批大小、并发、重试和内存上限只能分别配置,流水线容易出现排队、空转或内存峰值。

  • Vane Data:把多模态处理带回数据引擎

Vane Data 是基于 DuckDB 构建的多模态原生数据引擎,它延续 DuckDB 的 Relation API,把文本、图片和音频等数据作为列,将 AI Functions 与 Python UDF 组合进同一条关系查询,并由 vLLM 等 Provider 执行模型推理。处理后的结果仍可继续筛选、连接、聚合和写出。

从 DuckDB 出发,构建 AI 数据引擎

DuckDB 是一款面向分析场景的嵌入式分析数据库,能够直接在应用或 Python 进程中查询本地文件、内存数据和对象存储中的数据,常用于数据探索、ETL、Notebook 与轻量分析,也因此在开发者社区中得到广泛采用。

  • 轻量嵌入:从一个进程开始。 DuckDB 以进程内、无服务依赖的形态运行,通过一条安装命令,即可在本地进行部署,开始分析;对本地文件和内存数据尤其顺手,适合把数据处理嵌入应用、Notebook 与任务脚本。可以说是 OLAP 的 SQLite。
  • 极致性能:为分析负载而生。 DuckDB 面向分析负载融合列式存储、向量化执行、SIMD 指令级并行与多线程调度等技术,形成高效数据流水线,单机即可实现每秒数亿行的分析吞吐。公开的 ClickBench 结果 表明了其优异的分析性能。
  • 丰富的生态:按需加载可扩展。 DuckDB 提供了灵活的 扩展机制,使得用户可以定义新的数据类型、函数、文件格式以及新的 SQL 语法。数据处理的过程中,经常会涉及到多种数据源。DuckDB 可以查询 Parquet、Iceberg、CSV、JSON 和 Arrow 等,读取 Pandas、Polars 等 Python 对象,并通过该扩展机制连接对象存储、湖仓格式和外部数据库。业务记录、文件和内存中的表可以在同一个查询里组合,少一些中间文件和格式转换。
  • SQL 与 Python:表达能力强。 对使用者来说,SQL 和 Python SDK 都可以作为入口。可以用 SQL 写筛选、关联和聚合,也可以通过 Python 的 Relation API 逐步组合查询。

因此我们认为,DuckDB 在单机分析领域性能卓越,其设计理念令人称道。基于它构建 AI 多模态原生数据引擎,具备得天独厚的先天优势。于是在 DuckDB 的基础上,构建了 AI 多模态原生数据引擎 Vane Data,旨在帮助用户轻松打造多模态 AI 流水线。

Vane Data 多模态原生引擎:让数据、模型与计算资源协同工作

Vane Data 多模态原生数据引擎架构

在 DuckDB 的基础上,Vane Data 将多模态算子、AI 调用和资源调度接入同一条基于 Relation 的多模态数据处理流水线,DuckDB 中熟悉的读数据、筛选、聚合与导出逻辑在 Vane Data 中可以继续使用。下面先介绍如何在 Relation 中组织 Prompt、Embedding、UDF 与 vLLM,再说明这些阶段如何在单机与分布式环境中协调 CPU、GPU 和 I/O。

AI Functions、UDF 与 vLLM:基于 Relation 构建多模态数据流水线

多模态处理的关键,是让图片或音频等多模态数据和文本、表格字段一样成为 Relation 中的列,沿着可组合、可延迟执行的执行计划流动。Vane Data 可以构建基于 Relation 的多模态数据处理管道,具体特性包括:

  • AI Functions:让模型调用像列计算一样自然。 Vane Data 支持 ai_promptai_embed,既能在 SQL 中调用,也能通过 Python API 使用。ai_prompt 可接收文本或图片,支持 STRUCT 结构化输出;ai_embed 把文本转换为固定维度的向量。它们的结果仍是可过滤、连接、聚合的列;支持按需配置 on_error、重试以及其他 Provider 选项进行更多控制。
  • Python UDF:让自定义逻辑进入执行计划。 @vane.func 可以将无状态的处理逻辑定义为可在 SQL 中直接调用的函数,@vane.cls 将有状态的处理逻辑定义为函数,它可以在 Actor 内复用模型、客户端或解码器等资源。同时,这两种函数也都提供 batch 模式来批量处理,从而提升性能。
  • vLLM Provider:使用前缀路由复用 KV Cache。 Vane 的原生 vLLM Provider 在推理之外增加了面向数据流的前缀感知路由。运行时把共享前缀的请求放入有界桶,并尽量送到同一 vLLM Actor,以复用 KV/prefix cache;当该 Actor 的在途负载明显高于其他 Actor 时再做负载均衡。这样既提高重复 Prompt 的缓存命中机会,也避免为了追求亲和性而让单个 Actor 堵塞。

在这条链路中,Prompt 调用、Embedding 生成和结果过滤可以通过 Relation API 连续组合,并在最终取数或写出时作为一条完整的数据处理流程执行。

从单机到分布式:统一调度 CPU、GPU 与 I/O

前面的能力解决了“如何把多模态操作写进数据流程”,但要让这条流程稳定跑完,还需要回答“在哪里运行、如何分配资源、遇到波动怎么办”。Vane Data 将这些问题交给统一运行时处理,从单机开发到分布式执行,围绕资源、吞吐和可靠性提供以下机制:

  • 两个 Runner,一套业务逻辑。 创建连接前使用 vane.configure(runner="local") 选择单机 Local Runner,适合开发和小规模任务;改为 runner="ray" 即可交给 Ray Runner 调度多进程或多节点 CPU/GPU。Relation 和业务逻辑保持一致,变化集中在运行时配置与资源声明。
  • 动态批处理:适应数据大小与计算成本。 运行时同时参考行数和字节数,不把大小悬殊的多模数据对象硬塞进固定行数;超大输入会被拆分,输出缓冲区达到行数或字节阈值就刷新,分区分配器还会根据已观察到的 split 大小逐步调整目标。这样能在模型吞吐、内存峰值和单任务开销之间取得更稳的平衡。
  • 流水线并行:让异构资源同时工作。 读取/解码、CPU 预处理、GPU 推理和 I/O 写出是不同阶段,各自声明资源、并发和批大小。异步执行图允许相邻阶段重叠:GPU 处理当前批时,CPU 准备下一批,上一批结果同时写出,从而减少某一种资源等待另一种资源的时间。
  • 背压:用有界队列换取稳定吞吐。 下游 GPU、模型服务或存储变慢时,运行时用有界的在途任务、输出窗口和资源准入限制上游继续提交,并跟踪排队字节与待处理任务。队列不会无限增长,内存和对象存储峰值更可控。
  • 容错:让失败可重试、可隔离、可恢复。 瞬时 Provider 错误、任务或 Worker 故障可以按策略重试,Actor 也能在重建后重新加载模型或客户端;AI Function 调用可用 on_error="ignore" 转为 NULL 并保留错误信息,或用 raise 让任务明确失败。

这些机制共同把一次模型调用扩展成可持续运行的数据流水线:上游准备数据,中间阶段完成推理,下游汇总结果,资源和错误在各阶段之间被显式管理。Vane Data 的价值在于让这套能力仍然可以用熟悉的 DuckDB 方式开始。

案例:一条 SQL 串起图片预处理、AI 识别与审核分流

以车险理赔为例,理赔记录和现场照片分别保存在业务表和图片表中。我们需要先筛选待处理案件并关联照片,再用 Python UDF 完成图片方向校正、尺寸检查和格式统一,随后调用 AI Function 输出结构化的损伤判断,最后结合赔付金额与业务规则决定后续流程。除 UDF 的定义和注册外,整个数据处理过程都由一条 SQL 表达。

结构化输出的 JSON Schema 直接写在示例代码中,将模型返回约束为受损部位、严重程度和置信度三个字段;后续 SQL 可以直接读取这些结构化结果。

example.py
from io import BytesIO


from PIL import Image, ImageOps


import vane


# Constrain model output to directly queryable structured fields.
damage_schema = """{
  "type": "object",
  "properties": {
    "damage_part": {"type": "string"},
    "severity": {"type": "string", "enum": ["low", "medium", "high"]},
    "confidence": {"type": "number"}
  },
  "required": ["damage_part", "severity", "confidence"],
  "additionalProperties": false
}"""


# Correct orientation, validate dimensions, and normalize the image to JPEG.
@vane.func(return_dtype="BLOB")
def prepare_image(raw: bytes) -> bytes | None:
    try:
        with Image.open(BytesIO(raw)) as image:
            image = ImageOps.exif_transpose(image).convert("RGB")
            if min(image.size) < 480:
                return None
            image.thumbnail((1024, 1024))
            output = BytesIO()
            image.save(output, "JPEG", quality=85)
            return output.getvalue()
    except (OSError, TypeError):
        return None


connection = vane.connect()
# Compose relational processing, the Python UDF, the AI Function, and business rules.
routes = connection.sql(
    """
    WITH prepared AS (
        -- Join, filter, and preprocess each image with the Python UDF.
        SELECT c.claim_id, c.claim_amount,
               prepare_image(i.content) AS image
        FROM pending_claims AS c
        LEFT JOIN claim_images AS i USING (claim_id)
        WHERE c.status = 'pending'
    )
    SELECT
        claim_id,
        assessment.*,
        -- Combine the model result and claim amount for business routing.
        CASE
            WHEN image IS NULL THEN 'Request more photos'
            WHEN assessment IS NULL OR assessment.confidence < 0.65 THEN 'Manual review'
            WHEN assessment.severity = 'high' OR claim_amount >= 50000 THEN 'Claims review'
            ELSE 'Automatic routing'
        END AS next_step
    FROM (
        SELECT *,
               -- Convert the preprocessed image into a structured assessment.
               CASE WHEN image IS NOT NULL THEN ai_prompt(
                   'Identify the damaged vehicle part and severity.',
                   image,
                   return_format := $damage_schema,
                   system_message := 'Use only visual evidence; lower confidence when uncertain.',
                   provider := 'openai',
                   model := 'gpt-4o-mini',
                   on_error := 'ignore'
               ) END AS assessment
        FROM prepared
    ) AS assessed
    """,
    params={"damage_schema": damage_schema},
)


# Materialize the complete processing plan.
routes.write_parquet("claim_routes.parquet")

在这个流程中,Python 负责图片解码和预处理,SQL 负责关联、筛选与业务分流,ai_prompt 则把图片转换为结构化判断。AI 返回的结果是可以继续参与条件判断和后续写出的列。最终调用 write_parquet 时,Vane Data 才执行这条完整的数据处理计划。这个案例展示了 Vane 如何将关系计算、Python 逻辑与 AI 推理放在同一个数据流水线中。

从 DuckDB 走向多模态 AI 数据处理

Vane Data 让熟悉 DuckDB 的开发者以更低学习成本进入多模态、智能化数据处理时代:用 SQL/Python 表达关系,用 AI Functions/UDF 处理多模数据,再由统一运行时协调 CPU、GPU 与分布式资源。

开始探索: