读取 DataFusion 查询的逐算子执行指标

来源:Apache DataFusion Python 官方指南《Execution Metrics》,维护方为 Apache Software Foundation。本文于 2026 年 10 月 5 日依全文译编;原页没有个人作者署名,文档路径为滚动更新版本。

查询变慢时,仅看最后返回了多少行,通常还不足以判断问题。DataFusion 会把逻辑计划编译成物理算子树,例如 FilterExec、ProjectionExec 或 HashAggregateExec。这些算子可以在执行期间记录统计数据,包括输出行数、计算时间以及落盘情况。把这些执行指标与算子树放在一起阅读,就能进一步判断过滤、聚合、扫描或分区分布是否值得调查。

本文从一个完全在内存中构建的销售表开始,保留原文的完整查询示例,再说明指标何时可用、怎样读取分区信息,以及为什么不能把汇总计算时间直接当成查询的端到端耗时。这里没有实测结果;文中的“应如何理解”是 API 语义说明,而不是本次运行报告。

DataFusion 指标读取流程:先执行查询并完整消费结果,再取得最近执行的物理计划,遍历算子并分别查看分区指标和汇总指标。
指标属于实际执行使用的物理计划。未完纪依据 Apache DataFusion 文档绘制;图中结构仅表示读取流程,不声称是某次运行生成的精确计划。

每个算子的指标,按分区采集

DataFusion 可能把同一个算子分配到多个分区并行执行。指标首先按分区记录,MetricsSet 的便捷属性再按指标名称汇总所有分区,给出该算子的整体统计。因此,一个 MetricsSet 对应的是一个算子的指标集合,不是整条查询的唯一计时器。

属性 含义 阅读时的注意点
output_rows 算子输出的行数,各分区求和 它是该算子的输出,不是源表输入行数,也不必等于最终结果行数
elapsed_compute 算子计算循环中的耗时,单位纳秒,排除 I/O 等待;各分区求和 并行分区和多个算子的时间不能简单视为查询的墙钟耗时
spill_count 因内存压力触发的落盘事件次数,各分区求和 这是无量纲事件数,不表示字节数或行数
spilled_bytes 落盘期间写出的总字节数,各分区求和 用来观察落盘数据量
spilled_rows 落盘期间写出的总行数,各分区求和 与事件数、字节数分别观察,三者不可互换

原文把 elapsed_compute 称为计算循环内部的“wall-clock CPU time”。这里保留其可操作的解释:它按纳秒记录算子计算循环的耗时、排除 I/O 等待,并在便捷属性中跨分区求和。不要因此把它等同于进程 CPU 记账、从提交到取完结果的总时间,或简单相加后的整个流水线耗时。

先执行并消费结果,再读取指标

某些算子,例如 DataSourceExec,会在建立物理计划时就创建 MetricsSet。因此,执行前调用 metrics() 可能已能拿到一个集合。但“存在集合”不代表“已经处理数据”:执行前的数值通常是 0 或 None,不能拿来说明查询做了多少工作。

有意义的数值来自终结执行操作:

  • collect():执行并收集结果。
  • collect_partitioned():执行并按分区收集结果。
  • execute_stream():必须把返回的流完整消费后,再读取最终指标。
  • execute_stream_partitioned():必须消费完所有分区的流,而不只是读完其中一个。

Notebook 中的 display(df) 或自动显示的 repr 不会填充这里要读取的指标。显示表格时,DataFusion 会做一次受限的内部执行来取得预览行,但不缓存那次物理计划,所以之后的 collect_metrics() 不代表这次预览操作。要取得这些指标,仍需显式调用上述终结操作。

同一个 DataFrame 每调用一次 collect(),或另一个终结操作,都会创建新的物理计划。df.execution_plan() 对应最近一次执行。如果要对比两次执行,应在每次执行完成后及时记录对应计划的指标,并明确记录输入、配置及采集时点,避免把不同执行的数据混在一起。

一个完整示例

下面保留源文的 API 调用和数据,只将注释译为中文。它创建三行内存数据,从中筛选第一列大于 1 的记录,然后遍历每个算子的指标。没有连接外部数据库、下载数据或写入生产目标。

from datafusion import SessionContext

ctx = SessionContext()
ctx.sql("CREATE TABLE sales AS VALUES (1, 100), (2, 200), (3, 50)")

df = ctx.sql("SELECT * FROM sales WHERE column1 > 1")

# 执行查询,填充指标
results = df.collect()

# 获取带有本次执行指标的物理计划
plan = df.execution_plan()

# 遍历全部算子并打印指标
for operator_name, ms in plan.collect_metrics():
    if ms.output_rows is not None:
        print(f"{operator_name}")
        print(f"  output_rows    = {ms.output_rows}")
        print(f"  elapsed_compute = {ms.elapsed_compute} ns")

# 查看原始的逐分区指标
for operator_name, ms in plan.collect_metrics():
    for metric in ms.metrics():
        print(
            f"  partition={metric.partition}  "
            f"{metric.name}={metric.value}  "
            f"labels={metric.labels()}"
        )

从 SQL 条件可以推导,逻辑上满足条件的是 (2, 200) 和 (3, 50) 两行。这一推导不意味着所有物理算子的 output_rows 都等于 2;每个算子有自己的输入输出语义,优化器还可能改变计划结构。本文没有构造或展示“实测”算子耗时、落盘次数和物理计划输出。

第一段遍历仅打印存在 output_rows 的节点,沿用原文的过滤条件。这适合快速浏览,但不表示被跳过的算子完全没有其他有用指标。第二段遍历输出原始指标,原文没有在每个节点前再次打印名称;如果日志需要独立保存,可以加上下面这一行,避免不同算子的同名指标混淆。

# 编辑补充:置于第二个外层 for 循环内部、内层循环之前
print(f"operator={operator_name}")

物理计划树与算子名称

execution_plan() 返回物理计划树的根 ExecutionPlan 节点。树形结构反映算子流水线:根节点通常是投影或合并节点,子节点可能是过滤、聚合、扫描等。

plan.collect_metrics() 遍历整棵树,返回 (operator_name, MetricsSet) 对。这里的 operator_name 是节点的显示名称,例如 FilterExec: column1@0 > 1,与 plan.display() 显示的节点字符串一致。它可能包含表达式,不只是一个短的算子类型名称。

如果只关心一个节点,调用该节点的 metrics(),而不是误以为它会自动遍历所有子节点。定位耗时或输出行数异常时,可先用整树遍历确定位置,再结合具体节点和分区深入检查。

用原始指标看分区差异

MetricsSet.output_rows 等属性适合看整体;要发现数据倾斜,需要保留分区身份。对一个已取得的 metrics_set,原文给出的读取方式是:

for metric in metrics_set.metrics():
    print(f"  partition={metric.partition}  {metric.name}={metric.value}")

partition 从 0 开始编号,表示处理该指标的并行分区。如果是与特定分区无关的全局指标,其值为 None。不能把 None 强行归入第 0 分区,也不能仅因存在多个分区就认为负载已经均衡。

检查倾斜时,应比较同一算子、同一指标、相同语义标签下的分区。某个分区处理远多于其他分区的行数,是值得继续调查的线索;仅凭一个汇总数值看不到这种差异。这里只说明观察方法,未对示例宣称存在或不存在倾斜。

标签可能改变指标的含义

单个 Metric 可以携带键值标签,提供算子特有的额外上下文。大多数指标的标签字典为空;某些算子会用它区分指标变体。例如,原文举出的 HashAggregateExec 可能分别记录中间输出和最终输出的 output_rows:

for metric in metrics_set.metrics():
    print(metric.name, metric.labels())
# 原文给出的标签示意,不是本文执行输出:
# output_rows  {'output_type': 'final'}
# output_rows  {'output_type': 'intermediate'}

output_rows 便捷属性以及 sum_by_name() 按名称求和时,不会按标签区分;同名指标都会加在一起。因此,若要只统计最终输出,就必须遍历原始 Metric 并按实际标签筛选。不要假设所有算子、所有版本都会提供示例中的 output_type 标签。

未列入便捷属性的其他指标,也可以通过 sum_by_name() 取得,或者遍历 metrics() 返回的原始对象。具体指标是否存在、单位为何、是否带标签,应以运行版本的 API 和对应算子为准。

版本与代码静态审查

本次核对的是 2026 年 10 月 5 日可读取的官方滚动文档,没有据此给任意旧版 Python 包承诺兼容。执行前应核对本地 datafusion 是否提供 ExecutionPlan.collect_metrics()、MetricsSet、Metric 等接口。物理计划 API 参考与本文原页中的 API 链接给出相应定义。

静态审查未发现原文示例中存在硬编码秘密或拼接外部输入的 SQL 注入入口:两条 SQL 都是固定字符串,数据是三个常量元组。将这个示例扩展成服务时,不能把未经验证的输入直接拼入 SQL;输出计划及表达式的日志也可能泄露业务表名、列名或查询常量,应按实际数据环境处理。

collect() 会把结果收集到内存,三行示例没有规模问题,但不能未经评估地照搬到大结果集。改用流式接口时,必须完整消费才能得到最终统计。本文没有安装或执行 DataFusion,没有生成性能数字,也没有验证某个生产查询;上述结论只覆盖已展示代码的静态检查,没有发现问题不等于不存在漏洞。

原文与项目:Execution Metrics、Apache DataFusion Python。原项目按 Apache License 2.0 提供,许可副本随交付文件保留;源文译编、中文注释及补充说明于 2026 年 10 月 5 日修改。Apache、DataFusion 及相关名称和标识的权利属于 Apache Software Foundation;译编不表示基金会背书。原有版权、归属与无担保声明继续保留。

保留的完整许可声明

以下为原项目适用许可文本,原作者与文档归属及本文改动说明见正文。

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 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容