UDF
一个 Vane Data UDF 可以从四个相关维度来描述。这些维度回答不同的问题, 但 API 会约束哪些组合有效:
- API 入口: SQL Expression、Python Expression 或 Relation。
- 调用形态: scalar value、单行或 Arrow batch;Relation 中分别对应 map、flat_map 和 map_batches。
- 基数: 该阶段保持还是改变行数。
- 生命周期: 无状态函数或有状态 callable class。
1. API 入口
根据输出契约以及该阶段所需的执行控制来选择 API 入口。
| API 入口 | 调用方式 | 输出契约 | 适用场景 |
|---|---|---|---|
| SQL Expression UDF | 使用 vane.attach_function 挂载 callable,再在 SELECT 中调用其 alias | 一个有类型的 projection 列;v1 中保持行数 | 每个输入行旁需要一个结果,并且周围的 pipeline 使用 SQL |
| Python Expression UDF | 使用 vane.func、vane.func.batch、vane.cls 或 vane.cls.batch 构建 Expression | 一个有类型、保持行数的 projection 列 | 相同的 projection 更适合用 Python 对象构建 |
| Relation UDF | 调用 rel.map、rel.flat_map 或 rel.map_batches | 返回 Relation;map 追加 value,flat_map 和 map_batches 定义完整输出 schema | callable 需要控制表形或行数,或者该阶段需要显式 backend 或 actor pool 等 Relation 专属控制项 |
Expression UDF 都会保持行数。callable 需要改变行数时,应使用 Relation flat_map 或 map_batches。Relation UDF 目前是 Python 方法;Vane 不会把 它们暴露成 SQL table function。
案例:SQL Expression、Python Expression 与 Relation
下面的程序分别通过 SQL 和 Python Expression API 执行相同的规范化逻辑, 再用 Relation UDF 生成多列结果:
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. 调用形态:map、flat_map 与 map_batches
这些是 Relation 方法的实际名称。最接近的 Expression 形式具有相似的 callable 形态,但仍遵守单列 Expression 契约。
| 操作 | 调用形态 | 基数 | Relation 输出 | Expression 对应形式 |
|---|---|---|---|---|
| map | 输入 scalar 列值,输出一个有类型的值 | 保持行数 | 保留输入列并追加 value | Python 中使用 vane.func,SQL 中使用挂载的 scalar alias |
| flat_map | 输入一个 row dict,输出零个或多个 row dict | 可以改变行数 | 仅返回声明的 schema | v1 中没有 Expression 对应形式;SQL 也没有 Relation table-function 对应形式 |
| map_batches | 输入 pyarrow.Table,输出 pyarrow.Table | 可以保持或改变行数 | 仅返回声明的 schema | Python 中使用 vane.func.batch,SQL 中使用挂载的 batch alias;Expression schema 必须恰好包含一列 |
案例:组合三种 Relation 操作
下面的 pipeline 依次使用 map、flat_map 和 map_batches:
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_map 和 map_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.func 或 vane.func.batch;wrapper 也可以挂载为 SQL alias | 将函数传给 task backend | 不保证在多次调用之间保留 callable instance;不同 batch 可以独立处理 |
| 有状态 | 实例化 vane.cls 或 vane.cls.batch;在 Python 中调用,或把实例挂载为 SQL alias | 将 callable class 传给 subprocess_actor 或 ray_actor | 每个 actor 有一个 instance;各 actor 拥有独立、仅在本次执行中有效的状态 |
对于不需要复用初始化结果的确定性工作,使用无状态函数。需要复用模型、 tokenizer、client、lookup table 或 cache 时,使用有状态 class。每次执行 有状态 Expression 时,都会按 actor_number 创建相互独立的 Actor 实例; Relation 的 Actor 后端也使用这个配置值作为池大小。
案例:无状态函数与有状态 Actor
下面的示例通过 Python 和 SQL 对比无状态与有状态的 Expression UDF, 然后展示有状态 Relation UDF 的 actor 复用:
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 文本规范化 alias | SQL Expression | Scalar value | 保持行数 | 无状态 |
| Python 缓存规范化 | Python Expression | Scalar value | 保持行数 | 有状态 |
| Token 展开 | Relation | 单行(flat_map) | 改变行数 | 无状态 |
| Actor-backed 模型推理 | Relation | Arrow batch(map_batches) | 由返回的表声明 | 有状态 |