跳到主要内容
Vane Data / API 参考

vane.func.batch

使用 vane.func.batch,可以让 Python 函数在 select() 中一次处理一批数据。Expression 输入和 Python 字面量都会转换成 Arrow 列,函数需要返回一个等长的 Arrow 列。

签名

text
vane.func.batch(
    *,
    return_dtype: Any,
    name: str | None = None,
    batch_size: int | None = None,
    unnest: bool = False,
    gpus: float | None = None,
) -> Callable[[_PythonFunction], VaneBatchFunction]

参数

参数类型说明默认值
return_dtypeSQL 类型字符串、Vane DuckDBPyType 或受支持的 pyarrow.DataType返回列的类型必填
name非空 strNoneUDF 名称;省略时使用函数的 __qualname__None
batch_size正整数或 None函数每次最多处理的行数None
unnestbool返回 Struct 时,将字段展开成多列False
gpus非负有限数或 None每个 Task 使用的 GPU 资源;正值需要 RayNone

返回值与错误

select() 中调用时,输出行数与输入保持一致。unnest=False 时,查询产生一个 return_dtype 类型的结果列;unnest=True 时,Struct 返回类型的各个字段会展开成独立列。

直接调用只接受 pyarrow.Arraypyarrow.ChunkedArray,所有输入必须等长。规范化后的结果只有一个 chunk 时返回 Array,否则返回 ChunkedArray。

函数运行失败,或输入、返回值不是受支持的 Arrow 列,行数与输入不一致,类型无法转换为 return_dtype 时,会在直接调用或获取查询结果时报错。分布式后端可能重试整个批次,因此外部副作用必须具备幂等性。

示例

example.py
import vane




@vane.func.batch(return_dtype="BIGINT")
def add(a, b):
    import pyarrow.compute as pc


    return pc.add(a, b)




source = vane.sql("SELECT * FROM (VALUES (1, 4), (2, 5), (3, 6)) AS t(a, b)")
result = source.select(add(vane.col("a"), vane.col("b")).alias("total"))


print(result.order("total").fetchall())
vane.close()

输出:

text
[(5,), (7,), (9,)]

源码与相关页面