把 DataFusion 表达式分发给 Python 进程

有些数据任务很适合直接拆开:驱动进程事先决定分片,每个工作进程把同一类计算应用到自己负责的数据,最后返回结果。DataFusion Python 可以把 Expr 表达式序列化后送到另一个进程,让表达式在那里求值。这样分发的是计算表达式;数据怎样切分、谁负责哪一份,仍由调用方决定。

本文根据 Apache DataFusion Python 的 Distributing work 完整翻译整理,并核对官方 multiprocessing 伴随示例。文档及示例归 Apache Software Foundation 和相关贡献者,适用 Apache License 2.0;原页未标个人作者。示例与接口仅经静态审阅,没有运行结果保证。

驱动进程构造 Expr 并发送给两个工作进程,各自在独立 SessionContext 中处理自己的数据批次,结果再返回;内联 Python UDF 只能来自可信发送者
表达式级分发的数据流。未完纪自绘;数组是源文演示数据,不是运行截图。

表达式怎样经过进程边界

把表达式作为参数传给 multiprocessing.Pool、Ray 远程调用或类似设施时,Python 的标准 pickle 机制负责搬运参数,DataFusion 的表达式序列化会随之发生。通常不必手工先调用一次编码函数再解码。

表达式里的内建函数、算术和比较运算具有可移植表示。Python 编写的标量、聚合、窗口 UDF 则可以随表达式内联传输,接收端自动重建可调用对象、签名与捕获状态。这种便利也意味着必须认真对待代码来源:恢复内联 Python UDF 可能执行任意 Python 代码。

一个完整的基本进程池例子

下面把原文分开的工作函数和驱动片段拼成一个文件。与原文相比,增加了主模块守卫,并把启动方式显式改成 spawn,以便 Windows 以及支持它的其他平台使用;这不是执行验证。原文的 forkserver 适用于提供该方式的 POSIX 环境,不能原样用于 Windows。

# Adapted from Apache DataFusion Python documentation.
# Licensed under the Apache License, Version 2.0.
import multiprocessing as mp

import pyarrow as pa
from datafusion import SessionContext, col, udf


def evaluate(expr, batch):
    ctx = SessionContext()
    df = ctx.from_pydict({"a": batch})
    return df.with_column("result", expr).select("result").to_pydict()["result"]


def main():
    double = udf(
        lambda arr: pa.array([(v.as_py() or 0) * 2 for v in arr]),
        [pa.int64()],
        pa.int64(),
        volatility="immutable",
        name="double",
    )
    expr = double(col("a"))
    with mp.get_context("spawn").Pool(processes=4) as pool:
        results = pool.starmap(
            evaluate,
            [(expr, [1, 2, 3]), (expr, [10, 20, 30])],
        )
    print(results)


if __name__ == "__main__":
    main()

驱动先构造 double(col("a")),然后把它与两个数据批次一起交给进程池。每个工作调用创建自己的 SessionContext,把批次变成 DataFrame,增加结果列并只取回该列。源文给出的预期值为 [[2, 4, 6], [20, 40, 60]];这里没有实际执行,不能据此推断启动时间、加速比或进程调度方式。

if __name__ == "__main__": 对使用 spawn 或 forkserver 的脚本很重要。子进程重新导入模块时,不能再次执行创建进程池的驱动逻辑。工作函数还应定义在能够被导入的模块层级。UDF 能内联序列化,并不意味着外围所有进程池回调都能随意写成不可导入的局部对象。

这段 UDF 保留原文的 v.as_py() or 0 处理:空值被转换为零,然后翻倍。真实业务应明确 null 是传播、忽略、报错还是替换为零。volatility="immutable" 是对函数性质的声明;只有当同样输入始终给出相同输出、没有依赖外部可变状态时,才应这样声明。示例逐元素进入 Python,重点在传输机制,不是 Arrow 向量化性能优化。

哪些东西随表达式走,哪些需要预先注册

表达式内容 传输方式 接收端条件
内建函数、算术、比较 DataFusion 可移植表示 无需为这些函数额外注册
Python 标量 UDF、UDAF、窗口 UDF 默认可内联函数、签名及闭包 匹配的 Python 小版本,以及可导入的依赖
由 FFI capsule 协议导入的 UDF 仅传名称 工作上下文必须已有匹配注册,否则求值报错

内联 Python UDF 依赖 cloudpickle。Python 字节码不是跨小版本的稳定协议,因此驱动和工作进程必须使用相同的 Python 主/小版本,例如 3.12 不能与 3.11 或 3.13 混用。文档说明传输格式带有发送方版本,版本不匹配时会报出双方版本。

第二个条件是导入依赖可用。cloudpickle 可以按值捕获函数体和闭包,但通过 import 引用的模块通常按模块路径引用。如果 UDF 使用 from mylib import transform,工作进程也必须安装可兼容的 mylib;从导入类取得的绑定方法也有类似问题。自包含的小函数最容易搬运,但仍须确保 PyArrow 和 DataFusion 本身的版本兼容。

截至核对日期,官方主分支的 pyproject.toml 声明 Python 至少 3.10、cloudpickle 至少 2.0,并按 Python 版本分别要求 PyArrow:3.14 之前至少 16.0.0,3.14 及之后至少 22.0.0。这是所读主分支的依赖证据,不是所有已发布 wheel 的统一保证。部署时要记录实际安装版本,并检查本地是否提供本文所用的 IPC API。

为工作进程安装注册表

FFI UDF 只传名称,因而需要在每个工作进程启动时创建上下文、注册函数,再用 set_worker_ctx() 安装它。以下是注册模式;my_ffi_aggregate 必须由使用者实际安装的扩展导入,因此这个片段不是独立可运行文件。

from datafusion import SessionContext
from datafusion.ipc import set_worker_ctx

def init_worker():
    ctx = SessionContext()
    ctx.register_udaf(my_ffi_aggregate)  # 来自已安装且可信的 FFI 扩展
    set_worker_ctx(ctx)

# 在 main() 的主模块守卫之下创建池:
# with mp.get_context("spawn").Pool(
#     processes=4, initializer=init_worker
# ) as pool:
#     ...

接收表达式时,按名称引用的函数会到工作上下文中解析。如果没有安装工作上下文,则回退到全局 SessionContext。仅包含内建函数和内联 Python UDF 的表达式通常不需要额外注册;FFI 函数则必须在实际负责解析的上下文里注册,不能只在驱动进程的局部变量里准备好就期待子进程知道。

Python 3.14 在 Linux 上把默认启动方式从 fork 改为 forkserver;macOS 自 Python 3.8 起默认 spawn,Windows 使用 spawn。原文因此建议使用初始化函数,而不依赖 fork 的写时复制恰好继承了父进程状态。官方伴随例子还提醒 PyArrow/Tokio 等多线程运行时与不恰当的 fork 组合存在风险。

传输成本和闭包状态

原文用数量级描述编码体积:只含内建函数的表达式可能是几十字节,携带 Python UDF 后可能达到数百字节。这些数字不是固定上限,更不是性能基准。序列化大小随函数和捕获对象增加;一个闭包如果抓住大字典、可变对象、文件路径甚至访问令牌,实际负担和风险可能远超示例。

频繁发送相同函数时,可以考虑在每个工作进程预先注册等价 FFI UDF,再按名称引用,以降低重复携带函数体的开销。反过来,如果不同任务确实依赖不同的闭包值,就不能把它们错误压成同一个按名称函数。官方完整伴随文件为阈值 10、30、60 分别创建闭包 UDF,统计数据中大于各阈值的元素,并另加一个自定义求和 UDAF;返回结果还带工作进程 PID,用于观察分发。本文已静态阅读该文件的 accumulator、任务构建、求值和主模块入口,但没有运行它。

伴随示例的静态问题:所读版本的 _SumAccumulator.merge() 对每个状态数组只取 s[0]。当一个状态字段的数组包含多个分区的部分和时,这会遗漏后续元素。官方 Accumulator API 示例则对 states[0] 整个数组求和。若复用伴随文件里的单字段求和 accumulator,可把合并方法改为下面的形式;这是与原文不同的静态修订,未执行多分区验证:

def merge(self, states: list[pa.Array]) -> None:
    for subtotal in states[0]:
        self._total += subtotal.as_py() or 0

这里的 states[0] 表示唯一状态字段的整列部分和,不是只取该列第一个值。不能因为小演示只产生一个部分状态,就认定同一实现适用于任意分区数。

关闭内联能缩小风险面,却不能修复 pickle

调用 SessionContext.with_python_udf_inlining(enabled=False),可让该上下文产生或接收的表达式不再内联 Python 函数。Python UDF 此时也仅传名称,接收方必须有兼容注册。这个模式适合非 Python 接收者,例如 Java、C++ 或另一个 Rust 程序;它们无法重建 cloudpickle 负载。

严格接收上下文还有安全作用:Expr.from_bytes() 在这种配置下不会为输入里的内联 UDF 调用 cloudpickle.loads,违规内联负载会报错,而不是悄悄降级执行。若希望 pickle.dumps(expr) 间接调用无上下文参数的 Expr.to_bytes() 时也采用严格编码,可设置发送上下文:

from datafusion import SessionContext
from datafusion.ipc import set_sender_ctx, set_worker_ctx

# 驱动线程:安装严格发送上下文。
sender = SessionContext().with_python_udf_inlining(enabled=False)
set_sender_ctx(sender)

# 每个工作进程的初始化函数内:安装严格接收上下文。
worker = SessionContext().with_python_udf_inlining(enabled=False)
set_worker_ctx(worker)

这只是上下文设置示意,实际发送端和工作端通常处于不同进程,且还需要注册允许的函数。显式调用 Expr.to_bytes(ctx) 或 Expr.from_bytes(blob, ctx=ctx) 时,以传入上下文为准。

安全边界必须分两层看:关闭内联仅缩小 Expr.from_bytes() 的攻击面,外部字节经过 pickle.loads() 仍然不安全,不受这个开关保护。不要把一个“严格 Expr”封装进来自不可信发送者的 pickle,然后认为反序列化已安全。对不可信来源工作流,应禁用 Python UDF 内联、只允许内建函数和预先注册的 Rust 侧 UDF,并完全避免对外部字节调用 pickle.loads()。即使如此,函数白名单、输入规模、授权和计算资源限额仍由系统设计负责。

同一个 SessionContext 类型的四种位置

位置 生命周期 作用
用户持有 局部变量或属性 构造和执行查询,例如 ctx = SessionContext()
全局 进程单例,按需初始化 模块级 read_parquet()、read_csv() 等读取函数,以及表达式解码的最终回退;可由 SessionContext.global_ctx() 访问
发送槽位 驱动端线程局部 决定没有显式上下文的 pickle.dumps()/Expr.to_bytes() 编码设置,由 set_sender_ctx() 安装
工作槽位 工作端线程局部 提供没有显式上下文的解码函数注册表,由 set_worker_ctx() 安装

这些位置没有创建四种不同上下文类。同一个对象可以占据多个位置,安装操作保存引用,并不复制对象。解码解析顺序是显式参数、工作上下文、全局上下文;发送槽位不参与解码,工作槽位不参与编码。

线程局部也是常见陷阱:后台线程在自行安装前看不到主线程的发送或工作槽位。fork 可能把父进程当前线程局部值和全局上下文通过写时复制带入子进程;spawn 与 forkserver 则从新的状态起步。因此每个工作初始化函数都应显式安装或清除对应槽位。内联开关属于上下文本身,也不是整个进程唯一的总开关。

表达式级和查询级分发的界线

本文模式适合驱动能预先决定分区、各任务彼此独立的工作。自动把一个逻辑计划或物理计划拆成阶段,在节点之间 shuffle,再合并结果,则属于查询级分发。

截至所核对页面,datafusion-distributed 的 DataFusion Python 集成仍在开发,Apache Ballista 的调度器/执行器式查询集成也仍列在路线图中;页面明确标为尚不能从 datafusion-python 使用。不能因为 Python 已能发送 Expr,就把这两条路线写成开箱即用的集群执行功能。进一步阅读可从 datafusion.ipc 和仓库中的 Ray 伴随示例入手。

Apache Arrow DataFusion、Apache 及相关标志属于 Apache Software Foundation 的商标。本文中文翻译、编排、spawn 修订和技术示意图由未完纪制作。Apache License 2.0 包含按现状提供及无保证条款,完整许可及上游示例署名见文末。核对日期:2026-10-05。

上游示例署名与许可

保留 Apache DataFusion Python 官方伴随示例的原始署名通知;本文的中文翻译、示例改写与修订已在正文标明。Apache License 2.0 官方文本。

# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements.  See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership.  The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License.  You may obtain a copy of the License at
#
#   http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied.  See the License for the
# specific language governing permissions and limitations
# under the License.

Apache License 2.0


                                 Apache License
                           Version 2.0, January 2004
                        http://www.apache.org/licenses/

   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION

   1. Definitions.

      "License" shall mean the terms and conditions for use, reproduction,
      and distribution as defined by Sections 1 through 9 of this document.

      "Licensor" shall mean the copyright owner or entity authorized by
      the copyright owner that is granting the License.

      "Legal Entity" shall mean the union of the acting entity and all
      other entities that control, are controlled by, or are under common
      control with that entity. For the purposes of this definition,
      "control" means (i) the power, direct or indirect, to cause the
      direction or management of such entity, whether by contract or
      otherwise, or (ii) ownership of fifty percent (50%) or more of the
      outstanding shares, or (iii) beneficial ownership of such entity.

      "You" (or "Your") shall mean an individual or Legal Entity
      exercising permissions granted by this License.

      "Source" form shall mean the preferred form for making modifications,
      including but not limited to software source code, documentation
      source, and configuration files.

      "Object" form shall mean any form resulting from mechanical
      transformation or translation of a Source form, including but
      not limited to compiled object code, generated documentation,
      and conversions to other media types.

      "Work" shall mean the work of authorship, whether in Source or
      Object form, made available under the License, as indicated by a
      copyright notice that is included in or attached to the work
      (an example is provided in the Appendix below).

      "Derivative Works" shall mean any work, whether in Source or Object
      form, that is based on (or derived from) the Work and for which the
      editorial revisions, annotations, elaborations, or other modifications
      represent, as a whole, an original work of authorship. For the purposes
      of this License, Derivative Works shall not include works that remain
      separable from, or merely link (or bind by name) to the interfaces of,
      the Work and Derivative Works thereof.

      "Contribution" shall mean any work of authorship, including
      the original version of the Work and any modifications or additions
      to that Work or Derivative Works thereof, that is intentionally
      submitted to Licensor for inclusion in the Work by the copyright owner
      or by an individual or Legal Entity authorized to submit on behalf of
      the copyright owner. For the purposes of this definition, "submitted"
      means any form of electronic, verbal, or written communication sent
      to the Licensor or its representatives, including but not limited to
      communication on electronic mailing lists, source code control systems,
      and issue tracking systems that are managed by, or on behalf of, the
      Licensor for the purpose of discussing and improving the Work, but
      excluding communication that is conspicuously marked or otherwise
      designated in writing by the copyright owner as "Not a Contribution."

      "Contributor" shall mean Licensor and any individual or Legal Entity
      on behalf of whom a Contribution has been received by Licensor and
      subsequently incorporated within the Work.

   2. Grant of Copyright License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      copyright license to reproduce, prepare Derivative Works of,
      publicly display, publicly perform, sublicense, and distribute the
      Work and such Derivative Works in Source or Object form.

   3. Grant of Patent License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      (except as stated in this section) patent license to make, have made,
      use, offer to sell, sell, import, and otherwise transfer the Work,
      where such license applies only to those patent claims licensable
      by such Contributor that are necessarily infringed by their
      Contribution(s) alone or by combination of their Contribution(s)
      with the Work to which such Contribution(s) was submitted. If You
      institute patent litigation against any entity (including a
      cross-claim or counterclaim in a lawsuit) alleging that the Work
      or a Contribution incorporated within the Work constitutes direct
      or contributory patent infringement, then any patent licenses
      granted to You under this License for that Work shall terminate
      as of the date such litigation is filed.

   4. Redistribution. You may reproduce and distribute copies of the
      Work or Derivative Works thereof in any medium, with or without
      modifications, and in Source or Object form, provided that You
      meet the following conditions:

      (a) You must give any other recipients of the Work or
          Derivative Works a copy of this License; and

      (b) You must cause any modified files to carry prominent notices
          stating that You changed the files; and

      (c) You must retain, in the Source form of any Derivative Works
          that You distribute, all copyright, patent, trademark, and
          attribution notices from the Source form of the Work,
          excluding those notices that do not pertain to any part of
          the Derivative Works; and

      (d) If the Work includes a "NOTICE" text file as part of its
          distribution, then any Derivative Works that You distribute must
          include a readable copy of the attribution notices contained
          within such NOTICE file, excluding those notices that do not
          pertain to any part of the Derivative Works, in at least one
          of the following places: within a NOTICE text file distributed
          as part of the Derivative Works; within the Source form or
          documentation, if provided along with the Derivative Works; or,
          within a display generated by the Derivative Works, if and
          wherever such third-party notices normally appear. The contents
          of the NOTICE file are for informational purposes only and
          do not modify the License. You may add Your own attribution
          notices within Derivative Works that You distribute, alongside
          or as an addendum to the NOTICE text from the Work, provided
          that such additional attribution notices cannot be construed
          as modifying the License.

      You may add Your own copyright statement to Your modifications and
      may provide additional or different license terms and conditions
      for use, reproduction, or distribution of Your modifications, or
      for any such Derivative Works as a whole, provided Your use,
      reproduction, and distribution of the Work otherwise complies with
      the conditions stated in this License.

   5. Submission of Contributions. Unless You explicitly state otherwise,
      any Contribution intentionally submitted for inclusion in the Work
      by You to the Licensor shall be under the terms and conditions of
      this License, without any additional terms or conditions.
      Notwithstanding the above, nothing herein shall supersede or modify
      the terms of any separate license agreement you may have executed
      with Licensor regarding such Contributions.

   6. Trademarks. This License does not grant permission to use the trade
      names, trademarks, service marks, or product names of the Licensor,
      except as required for reasonable and customary use in describing the
      origin of the Work and reproducing the content of the NOTICE file.

   7. Disclaimer of Warranty. Unless required by applicable law or
      agreed to in writing, Licensor provides the Work (and each
      Contributor provides its Contributions) on an "AS IS" BASIS,
      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
      implied, including, without limitation, any warranties or conditions
      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
      PARTICULAR PURPOSE. You are solely responsible for determining the
      appropriateness of using or redistributing the Work and assume any
      risks associated with Your exercise of permissions under this License.

   8. Limitation of Liability. In no event and under no legal theory,
      whether in tort (including negligence), contract, or otherwise,
      unless required by applicable law (such as deliberate and grossly
      negligent acts) or agreed to in writing, shall any Contributor be
      liable to You for damages, including any direct, indirect, special,
      incidental, or consequential damages of any character arising as a
      result of this License or out of the use or inability to use the
      Work (including but not limited to damages for loss of goodwill,
      work stoppage, computer failure or malfunction, or any and all
      other commercial damages or losses), even if such Contributor
      has been advised of the possibility of such damages.

   9. Accepting Warranty or Additional Liability. While redistributing
      the Work or Derivative Works thereof, You may choose to offer,
      and charge a fee for, acceptance of support, warranty, indemnity,
      or other liability obligations and/or rights consistent with this
      License. However, in accepting such obligations, You may act only
      on Your own behalf and on Your sole responsibility, not on behalf
      of any other Contributor, and only if You agree to indemnify,
      defend, and hold each Contributor harmless for any liability
      incurred by, or claims asserted against, such Contributor by reason
      of your accepting any such warranty or additional liability.

   END OF TERMS AND CONDITIONS

   APPENDIX: How to apply the Apache License to your work.

      To apply the Apache License to your work, attach the following
      boilerplate notice, with the fields enclosed by brackets "[]"
      replaced with your own identifying information. (Don't include
      the brackets!)  The text should be enclosed in the appropriate
      comment syntax for the file format. We also recommend that a
      file or class name and description of purpose be included on the
      same "printed page" as the copyright notice for easier
      identification within third-party archives.

   Copyright [yyyy] [name of copyright owner]

   Licensed under the Apache License, Version 2.0 (the "License");
   you may not use this file except in compliance with the License.
   You may obtain a copy of the License at

       http://www.apache.org/licenses/LICENSE-2.0

   Unless required by applicable law or agreed to in writing, software
   distributed under the License is distributed on an "AS IS" BASIS,
   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
   See the License for the specific language governing permissions and
   limitations under the License.
© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容