跳到主要内容
Vane Data / 核心概念

UDF

一个 Vane Data UDF 可以从四个相关维度来描述。这些维度回答不同的问题, 但 API 会约束哪些组合有效:

  1. API 入口: SQL Expression、Python Expression 或 Relation。
  2. 调用形态: scalar value、单行或 Arrow batch;Relation 中分别对应 mapflat_mapmap_batches
  3. 基数: 该阶段保持还是改变行数。
  4. 生命周期: 无状态函数或有状态 callable class。

1. API 入口

根据输出契约以及该阶段所需的执行控制来选择 API 入口。

API 入口调用方式输出契约适用场景
SQL Expression UDF使用 vane.attach_function 挂载 callable,再在 SELECT 中调用其 alias一个有类型的 projection 列;v1 中保持行数每个输入行旁需要一个结果,并且周围的 pipeline 使用 SQL
Python Expression UDF使用 vane.funcvane.func.batchvane.clsvane.cls.batch 构建 Expression一个有类型、保持行数的 projection 列相同的 projection 更适合用 Python 对象构建
Relation UDF调用 rel.maprel.flat_maprel.map_batches返回 Relation;map 追加 valueflat_mapmap_batches 定义完整输出 schemacallable 需要控制表形或行数,或者该阶段需要显式 backend 或 actor pool 等 Relation 专属控制项

Expression UDF 都会保持行数。callable 需要改变行数时,应使用 Relation flat_mapmap_batches。Relation UDF 目前是 Python 方法;Vane 不会把 它们暴露成 SQL table function。

案例:SQL Expression、Python Expression 与 Relation

下面的程序分别通过 SQL 和 Python Expression API 执行相同的规范化逻辑, 再用 Relation UDF 生成多列结果:

example.py
import pyarrow as pa
import vane


con = vane.connect()
source = con.values(
    (
        vane.lit(1).alias("document_id"),
        vane.lit(" Refund requested ").alias("text"),
    ),
    (vane.lit(2), vane.lit(" Shipment delayed ")),
)




def normalize_value(value):
    return str(value).strip().lower()




@vane.func(return_dtype="VARCHAR")
def normalize_text(value):
    return normalize_value(value)




# SQL Expression UDF。
vane.attach_function(
    normalize_text,
    alias="normalize_text_sql",
    connection=con,
    parameters=["VARCHAR"],
)
try:
    sql_result = con.sql("""
        SELECT
            document_id,
            text,
            normalize_text_sql(text) AS normalized_text
        FROM source
        ORDER BY document_id
    """)
    sql_result.show()
finally:
    vane.detach_function("normalize_text_sql", connection=con)




# Python Expression UDF。
python_result = source.select(
    vane.col("document_id"),
    vane.col("text"),
    normalize_text(vane.col("text")).alias("normalized_text"),
).order("document_id")
python_result.show()




# Relation UDF:返回的 Arrow 表控制完整输出形态。
def enrich_batch(table):
    document_ids = table.column("document_id").to_pylist()
    values = table.column("text").to_pylist()
    normalized = [normalize_value(value) for value in values]
    return pa.table({
        "document_id": document_ids,
        "normalized_text": normalized,
        "text_length": [len(value) for value in normalized],
    })




relation_result = source.map_batches(
    enrich_batch,
    schema={
        "document_id": "INTEGER",
        "normalized_text": "VARCHAR",
        "text_length": "INTEGER",
    },
).order("document_id")
relation_result.show()

show() 或其他终端操作物化 Relation 之前,应保持 SQL alias 处于挂载状态。 Python Expression 调用会构建惰性的 projection 对象,不需要 alias。Relation callable 必须返回该阶段之后仍需保留的每一列。

2. 调用形态:mapflat_mapmap_batches

这些是 Relation 方法的实际名称。最接近的 Expression 形式具有相似的 callable 形态,但仍遵守单列 Expression 契约。

操作调用形态基数Relation 输出Expression 对应形式
map输入 scalar 列值,输出一个有类型的值保持行数保留输入列并追加 valuePython 中使用 vane.func,SQL 中使用挂载的 scalar alias
flat_map输入一个 row dict,输出零个或多个 row dict可以改变行数仅返回声明的 schemav1 中没有 Expression 对应形式;SQL 也没有 Relation table-function 对应形式
map_batches输入 pyarrow.Table,输出 pyarrow.Table可以保持或改变行数仅返回声明的 schemaPython 中使用 vane.func.batch,SQL 中使用挂载的 batch alias;Expression schema 必须恰好包含一列

案例:组合三种 Relation 操作

下面的 pipeline 依次使用 mapflat_mapmap_batches

example.py
import pyarrow as pa
import vane


con = vane.connect()
source = con.values(
    vane.lit(1).alias("document_id"),
    vane.lit("Refund requested").alias("text"),
)




# map:为每个输入行追加一个 scalar value。
def text_length(document_id, text):
    return len(text)




mapped = source.map(
    text_length,
    return_type=vane.sqltypes.INTEGER,
).project("document_id, text, value AS text_length")




# flat_map:把一篇文档转换成零个或多个单词行。
def split_words(row):
    for word in row["text"].split():
        yield {
            "document_id": row["document_id"],
            "text_length": row["text_length"],
            "word": word.lower(),
        }




words = mapped.flat_map(
    split_words,
    schema={
        "document_id": "INTEGER",
        "text_length": "INTEGER",
        "word": "VARCHAR",
    },
)




# map_batches:控制每个 Arrow batch 的完整输出表。
def add_word_length(table):
    document_ids = table.column("document_id").to_pylist()
    text_lengths = table.column("text_length").to_pylist()
    word_values = table.column("word").to_pylist()
    return pa.table({
        "document_id": document_ids,
        "text_length": text_lengths,
        "word": word_values,
        "word_length": [len(word) for word in word_values],
    })




result = words.map_batches(
    add_word_length,
    schema={
        "document_id": "INTEGER",
        "text_length": "INTEGER",
        "word": "VARCHAR",
        "word_length": "INTEGER",
    },
)
result.show()

map 会自动保留上游列。flat_mapmap_batches 只返回 schema 中声明的 列,因此应把后续仍需使用的上游列包含在 schema 中。

3. 基数

基数与调用形态是两个不同的问题。所有 Expression UDF 和 Relation map 都会为每个输入行保留一个输出行。Relation flat_map 明确为每个输入行生成 零个或多个输出行,Relation map_batches 则可以返回相同或不同数量的行。 callable 需要改变行数时,应使用这两个 Relation API。

4. 无状态与有状态

状态描述 callable 的生命周期,而不是输出形态。Expression 和 Relation UDF 都可以是无状态或有状态的。

生命周期Expression 形式Relation 形式运行时行为
无状态vane.funcvane.func.batch;wrapper 也可以挂载为 SQL alias将函数传给 task backend不保证在多次调用之间保留 callable instance;不同 batch 可以独立处理
有状态实例化 vane.clsvane.cls.batch;在 Python 中调用,或把实例挂载为 SQL alias将 callable class 传给 subprocess_actorray_actor每个 actor 有一个 instance;各 actor 拥有独立、仅在本次执行中有效的状态

对于不需要复用初始化结果的确定性工作,使用无状态函数。需要复用模型、 tokenizer、client、lookup table 或 cache 时,使用有状态 class。每次执行 有状态 Expression 时,都会按 actor_number 创建相互独立的 Actor 实例; Relation 的 Actor 后端也使用这个配置值作为池大小。

案例:无状态函数与有状态 Actor

下面的示例通过 Python 和 SQL 对比无状态与有状态的 Expression UDF, 然后展示有状态 Relation UDF 的 actor 复用:

example.py
import pyarrow as pa
import vane


con = vane.connect()
source = con.values(
    (
        vane.lit(1).alias("document_id"),
        vane.lit(" Refund requested ").alias("text"),
    ),
    (vane.lit(2), vane.lit(" Shipment delayed ")),
    (vane.lit(3), vane.lit(" Refund requested ")),
)




def normalize_value(value):
    return str(value).strip().lower()




@vane.func(return_dtype="VARCHAR")
def stateless_normalize(value):
    return normalize_value(value)




@vane.cls(actor_number=1, return_dtype="VARCHAR")
class CachedNormalizer:
    def __init__(self):
        self.cache = {}


    def __call__(self, value):
        if value not in self.cache:
            self.cache[value] = normalize_value(value)
        return self.cache[value]




cached_normalize = CachedNormalizer()


# Python Expression:一个 task-backed 函数与一个有状态 actor 并列使用。
python_result = source.select(
    vane.col("document_id"),
    stateless_normalize(vane.col("text")).alias("stateless_result"),
    cached_normalize(vane.col("text")).alias("stateful_result"),
).order("document_id")
python_result.show()


# SQL Expression:挂载无状态 wrapper 和已经实例化的有状态 wrapper,
# 而不是被装饰的 class 本身。
vane.attach_function(
    stateless_normalize,
    alias="stateless_normalize_sql",
    connection=con,
    parameters=["VARCHAR"],
)
vane.attach_function(
    cached_normalize,
    alias="cached_normalize_sql",
    connection=con,
    parameters=["VARCHAR"],
)
try:
    sql_result = con.sql("""
        SELECT
            document_id,
            stateless_normalize_sql(text) AS stateless_result,
            cached_normalize_sql(text) AS stateful_result
        FROM source
        ORDER BY document_id
    """)
    sql_result.show()
finally:
    vane.detach_function("cached_normalize_sql", connection=con)
    vane.detach_function("stateless_normalize_sql", connection=con)




# Relation actor 复用:每个 actor 构造并持有一个 class instance。
class CachedNormalizerBatch:
    def __init__(self):
        self.cache = {}


    def __call__(self, table):
        document_ids = table.column("document_id").to_pylist()
        values = table.column("text").to_pylist()
        normalized = []
        for value in values:
            if value not in self.cache:
                self.cache[value] = normalize_value(value)
            normalized.append(self.cache[value])
        return pa.table({
            "document_id": document_ids,
            "normalized_text": normalized,
        })




actor_result = source.map_batches(
    CachedNormalizerBatch,
    schema={"document_id": "INTEGER", "normalized_text": "VARCHAR"},
    execution_backend="ray_actor",
    actor_number=1,
    batch_size=1,
    gpus=0.0,
).order("document_id")
actor_result.show()

batch_size=1 会让这个小型 source 经历多次 actor 调用,从而展示 cache 复用; 这不是吞吐量调优建议。注册有状态 SQL alias 不会让其状态延续到多个 query。

状态不是 durable、checkpointed、global、keyed、transactional 或 exactly-once。 多个 actor 不共享内存,分布式输入没有稳定的全局行顺序,retry 还可能重复外部 副作用。应尽量避免让状态影响结果语义。

组合这些维度

在遵守上述有效组合的前提下,可以通过回答每个问题来概括一个 UDF:

示例API 入口调用形态基数生命周期
SQL 文本规范化 aliasSQL ExpressionScalar value保持行数无状态
Python 缓存规范化Python ExpressionScalar value保持行数有状态
Token 展开Relation单行(flat_map改变行数无状态
Actor-backed 模型推理RelationArrow batch(map_batches由返回的表声明有状态