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

Relation.map_batches

Relation.map_batches 把一批输入行作为 pyarrow.Table 交给同步 Python 可调用对象。函数可以改变列结构和行数,完整输出由 schema 声明。

签名

text
Relation.map_batches(
    function: Callable[..., typing.Any],
    schema: dict[str, sqltypes.DuckDBPyType] | None = None,
    *,
    batch_size: int | None = None,
    output_batch_size: int | None = None,
    min_task_batch_size: int | None = None,
    preserve_compute_batch_boundaries: bool | None = None,
    cpus: float | None = None,
    gpus: float | None = None,
    memory_bytes: int | None = None,
    execution_backend: typing.Literal["subprocess_task", "subprocess_actor", "ray_task", "ray_actor"] | None = None,
    actor_number: int | None = None,
    ray_actor_thread_policy: typing.Literal["managed", "ray_native"] | None = None,
    target_max_batch_bytes: int | None = None,
    task_input_max_bytes: int | None = None,
    output_target_max_bytes: int | None = None,
) -> DuckDBPyRelation

参数

参数类型说明默认值
function同步函数、绑定方法或可零参数构造的可调用类接收 pyarrow.Table,返回物化的 pyarrow.Tablepyarrow.RecordBatch、列 dict、产生这些类型的同步可迭代对象或 None必填
schema非空 dict[str, DuckDBPyType]完整输出列名、顺序和类型。当前运行时要求显式传入必填
batch_size正整数或 None一次调用最多接收的行数None
output_batch_size正整数或 None输出 Arrow 数据块的目标行数None
min_task_batch_size正整数或 NoneTask 输入合并的软下限;要求设置 batch_size,且不能小于它None
preserve_compute_batch_boundariesboolNoneTrue 时,在每个计算批次结束后刷新输出None
cpus非负有限数或 None每个 Task 或 Actor 的 CPU 资源None
gpus非负有限数或 None每个 Task 或 Actor 的 GPU 资源;正值只支持 RayNone
memory_bytes正整数或 None每个 Task 或 Actor 使用的内存资源;仅 Ray 后端可用None
execution_backendsubprocess_tasksubprocess_actorray_taskray_actorNone执行后端;省略时根据 runner 和可调用对象形态选择None
actor_number正整数或 NoneActor 实例数量;Actor 后端必填,Task 后端不能设置None
ray_actor_thread_policyray_nativemanagedNoneRay Actor 线程策略;仅适用于 ray_actor,当前默认解析为 ray_nativeNone
target_max_batch_bytes正整数或 NoneTask 输入和输出数据块的共同字节目标None
task_input_max_bytes正整数或 None单次 Task 或 Actor 调用的输入字节目标None
output_target_max_bytes正整数或 None输出数据块的字节目标None

返回值与错误

map_batches() 返回一个新的 Relation,原 Relation 不会改变。结果只包含 schema 中声明的列。

可调用对象必须返回已经物化的结果,不支持 pyarrow.RecordBatchReader

函数、schema 或执行参数不正确时,调用 map_batches() 就会报错。函数运行失败、返回值格式不正确或返回的列与 schema 不匹配时,会在 fetchall() 等取结果操作中报错。

Task 和 Actor 后端可能重试调用,因此外部副作用必须具备幂等性。可调用类会在相互独立、随时可能重建的 Actor 中运行;任务不保证固定分配给某个 Actor,也没有全局执行顺序,Actor 重建会清空本地状态。

示例

example.py
import vane




def keep_large_values(table):
    import pyarrow.compute as pc


    return table.filter(pc.greater(table["value"], 1))




source = vane.sql("SELECT * FROM (VALUES (1), (2), (3)) AS t(value)")
result = source.map_batches(
    keep_large_values,
    schema={"value": vane.sqltypes.BIGINT},
)


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

输出:

text
[(2,), (3,)]

可调用对象把三行输入筛成两行,直接体现 N → M。结果只包含 schema 声明的完整输出。

源码与相关页面