Vane Data / API 参考
Relation.map_batches
Relation.map_batches 把一批输入行作为 pyarrow.Table 交给同步 Python 可调用对象。函数可以改变列结构和行数,完整输出由 schema 声明。
签名
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.Table、pyarrow.RecordBatch、列 dict、产生这些类型的同步可迭代对象或 None | 必填 |
| schema | 非空 dict[str, DuckDBPyType] | 完整输出列名、顺序和类型。当前运行时要求显式传入 | 必填 |
| batch_size | 正整数或 None | 一次调用最多接收的行数 | None |
| output_batch_size | 正整数或 None | 输出 Arrow 数据块的目标行数 | None |
| min_task_batch_size | 正整数或 None | Task 输入合并的软下限;要求设置 batch_size,且不能小于它 | None |
| preserve_compute_batch_boundaries | bool 或 None | 为 True 时,在每个计算批次结束后刷新输出 | None |
| cpus | 非负有限数或 None | 每个 Task 或 Actor 的 CPU 资源 | None |
| gpus | 非负有限数或 None | 每个 Task 或 Actor 的 GPU 资源;正值只支持 Ray | None |
| memory_bytes | 正整数或 None | 每个 Task 或 Actor 使用的内存资源;仅 Ray 后端可用 | None |
| execution_backend | subprocess_task、subprocess_actor、ray_task、ray_actor 或 None | 执行后端;省略时根据 runner 和可调用对象形态选择 | None |
| actor_number | 正整数或 None | Actor 实例数量;Actor 后端必填,Task 后端不能设置 | None |
| ray_actor_thread_policy | ray_native、managed 或 None | Ray Actor 线程策略;仅适用于 ray_actor,当前默认解析为 ray_native | None |
| target_max_batch_bytes | 正整数或 None | Task 输入和输出数据块的共同字节目标 | 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 重建会清空本地状态。
示例
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()
输出:
[(2,), (3,)]可调用对象把三行输入筛成两行,直接体现 N → M。结果只包含 schema 声明的完整输出。