用 DataFusion 窗口函数排名与填补空值

原文:Window Functions。作者归属:Apache DataFusion Python 文档贡献者 / Apache Software Foundation。中文翻译与核验:未完纪 / Codex。原页未明确注明首发日期;本文核验日期为 2026-10-05。

窗口帧、排名与空值回填:排名:Attack 升序 → rank,相同攻击力可以并列;rank < 3 不保证恰好返回两行;ROWS(2, 0):当前行 + 前两行,15 → 20 → 20:第三行窗口均值为 55 / 3 ≈ 18.33;最后有效值:Grass → NULL → Poison,IGNORE_NULLS 得到 Grass → Grass → Poison;原始列仍保留 NULL;稳定性:顺序必须有确定含义,lag 与回填应指定排序;同值时按业务增加唯一次级键
图:窗口帧、排名与空值回填。根据原文与静态核验内容自主绘制,非软件界面截图。

窗口函数:每一行都有一个结果

窗口函数会使用一行或多行中的值,为每一行分别计算结果;聚合函数则通常把多行汇总为一个值。DataFusion 的窗口函数位于 datafusion.functions 模块。下面沿用原文来自 Ritchie Vink 的 Pokémon 数据集。

开始前,将官方构建脚本指定的 固定版本 pokemon.csv 保存到当前工作目录。该样例有 163 条数据,包含 Name、Type 1、Type 2、Attack、Speed 等 13 列,其中 Type 2 有 86 个空字段;这些是对源 CSV 的静态检查,不是执行下列 DataFusion 查询所得。

from datafusion import SessionContext
from datafusion import col, lit
from datafusion import functions as f

ctx = SessionContext()
df = ctx.read_csv("pokemon.csv")

下面的示例把每只 Pokémon 的速度与 DataFrame 中上一行的速度放到一起比较。

df.select(
    col('"Name"'),
    col('"Speed"'),
    f.lag(col('"Speed"')).alias("Previous Speed")
)
DataFrame()
+---------------------------+-------+----------------+
| Name                      | Speed | Previous Speed |
+---------------------------+-------+----------------+
| Bulbasaur                 | 45    |                |
| Ivysaur                   | 60    | 45             |
| Venusaur                  | 80    | 60             |
| VenusaurMega Venusaur     | 80    | 80             |
| Charmander                | 65    | 80             |
| Charmeleon                | 80    | 65             |
| Charizard                 | 100   | 80             |
| CharizardMega Charizard X | 100   | 100            |
| CharizardMega Charizard Y | 100   | 100            |
| Squirtle                  | 43    | 100            |
+---------------------------+-------+----------------+
Data truncated.

设置参数

排序

为 order_by 参数传入排序表达式列表,就可以控制窗口函数处理行的顺序。下例在每个第一属性(Type 1)分区内,按攻击力升序排名。

df.select(
    col('"Name"'),
    col('"Attack"'),
    col('"Type 1"'),
    f.rank(
        partition_by=[col('"Type 1"')],
        order_by=[col('"Attack"').sort(ascending=True)],
    ).alias("rank"),
).sort(col('"Type 1"'), col('"Attack"'))
DataFrame()
+------------+--------+--------+------+
| Name       | Attack | Type 1 | rank |
+------------+--------+--------+------+
| Metapod    | 20     | Bug    | 1    |
| Kakuna     | 25     | Bug    | 2    |
| Caterpie   | 30     | Bug    | 3    |
| Weedle     | 35     | Bug    | 4    |
| Butterfree | 45     | Bug    | 5    |
| Venonat    | 55     | Bug    | 6    |
| Venomoth   | 65     | Bug    | 7    |
| Paras      | 70     | Bug    | 8    |
| Beedrill   | 90     | Bug    | 9    |
| Parasect   | 95     | Bug    | 10   |
+------------+--------+--------+------+
Data truncated.

分区

与聚合函数类似,窗口函数可以接受一个 partition_by 列表。这样,窗口计算会在每个分区内部独立进行。上例按 Type 1 为每只 Pokémon 排名;下面筛选出各分区靠前的名次。

df.select(
    col('"Name"'),
    col('"Attack"'),
    col('"Type 1"'),
    f.rank(
        partition_by=[col('"Type 1"')],
        order_by=[col('"Attack"').sort(ascending=True)],
    ).alias("rank"),
).filter(col("rank") < lit(3)).sort(col('"Type 1"'), col("rank"))
DataFrame()
+-----------+--------+----------+------+
| Name      | Attack | Type 1   | rank |
+-----------+--------+----------+------+
| Metapod   | 20     | Bug      | 1    |
| Kakuna    | 25     | Bug      | 2    |
| Dratini   | 64     | Dragon   | 1    |
| Dragonair | 84     | Dragon   | 2    |
| Voltorb   | 30     | Electric | 1    |
| Magnemite | 35     | Electric | 2    |
| Clefairy  | 45     | Fairy    | 1    |
| Clefable  | 70     | Fairy    | 2    |
| Mankey    | 80     | Fighting | 1    |
| Machop    | 80     | Fighting | 1    |
+-----------+--------+----------+------+
Data truncated.

窗口帧

聚合函数作为窗口函数使用时,窗口帧定义参与计算的行。如果不显式指定窗口帧,原文给出的默认规则是:

  • 设置了 order_by:从分区开始(无界前向)到当前行。
  • 未设置 order_by:从无界前向到无界后向,也就是整个分区。

窗口帧由三个参数定义:单位类型、起始边界、结束边界。支持以下单位:

  • Rows:边界由相对于当前行的行数定义。
  • Range:原文要求 order_by 恰好包含一项;边界由各行排序表达式的值与当前值之间的距离定义。
  • Groups:一个组由在全部 order_by 表达式上取值相同的行组成。

下例计算当前 Pokémon 及其前两行的速度滚动平均值。

from datafusion.expr import Window, WindowFrame

df.select(
    col('"Name"'),
    col('"Speed"'),
    f.avg(col('"Speed"'))
    .over(Window(window_frame=WindowFrame("rows", 2, 0), order_by=[col('"Speed"')]))
    .alias("Previous Speed"),
)
DataFrame()
+------------+-------+--------------------+
| Name       | Speed | Previous Speed     |
+------------+-------+--------------------+
| Slowpoke   | 15    | 15.0               |
| Jigglypuff | 20    | 17.5               |
| Geodude    | 20    | 18.333333333333332 |
| Paras      | 25    | 21.666666666666668 |
| Grimer     | 25    | 23.333333333333332 |
| Rhyhorn    | 25    | 25.0               |
| Snorlax    | 30    | 26.666666666666668 |
| Metapod    | 30    | 28.333333333333332 |
| Oddish     | 30    | 30.0               |
| Parasect   | 30    | 30.0               |
+------------+-------+--------------------+
Data truncated.

空值处理

把聚合函数用作窗口函数时,常常需要指定如何处理空值。在线原文说明可使用构建接口,并预计未来会简化接口。下面的实际代码通过 Window 传入 null_treatment。

一种常见需求是查找截至当前行为止最后出现的有效值。把空值处理方式设为忽略空值后,就会使用最近一条非空记录的值填补结果。还需要把窗口终点设为当前行,以免读到后面的值。

示例先筛选第一属性为 Bug 的 Pokémon,这些记录的第二属性 Type 2 中有一些空值;随后同时计算忽略空值和保留空值的结果。

from datafusion.common import NullTreatment

df.filter(col('"Type 1"') == lit("Bug")).select(
    '"Name"',
    '"Type 2"',
    f.last_value(col('"Type 2"'))
    .over(
        Window(
            window_frame=WindowFrame("rows", None, 0),
            order_by=[col('"Speed"')],
            null_treatment=NullTreatment.IGNORE_NULLS,
        )
    )
    .alias("last_wo_null"),
    f.last_value(col('"Type 2"'))
    .over(
        Window(
            window_frame=WindowFrame("rows", None, 0),
            order_by=[col('"Speed"')],
            null_treatment=NullTreatment.RESPECT_NULLS,
        )
    )
    .alias("last_with_null"),
)
DataFrame()
+------------+--------+--------------+----------------+
| Name       | Type 2 | last_wo_null | last_with_null |
+------------+--------+--------------+----------------+
| Paras      | Grass  | Grass        | Grass          |
| Metapod    |        | Grass        |                |
| Parasect   | Grass  | Grass        | Grass          |
| Kakuna     | Poison | Poison       | Poison         |
| Caterpie   |        | Poison       |                |
| Venonat    | Poison | Poison       | Poison         |
| Weedle     | Poison | Poison       | Poison         |
| Butterfree | Flying | Flying       | Flying         |
| Beedrill   | Poison | Poison       | Poison         |
| Pinsir     |        | Poison       |                |
+------------+--------+--------------+----------------+
Data truncated.

把聚合函数用作窗口函数

任意聚合函数都可以用作窗口函数。下例用 datafusion.functions.avg(),把每只 Pokémon 的攻击力与其第一属性分组的平均攻击力放在同一行。窗口帧显式覆盖整个分区。

df.select(
    col('"Name"'),
    col('"Attack"'),
    col('"Type 1"'),
    f.avg(col('"Attack"')).over(
        Window(
            window_frame=WindowFrame("rows", None, None),
            partition_by=[col('"Type 1"')],
        )
    ).alias("Average Attack"),
)
DataFrame()
+-------------------+--------+--------+----------------+
| Name              | Attack | Type 1 | Average Attack |
+-------------------+--------+--------+----------------+
| Dragonite         | 134    | Dragon | 94.0           |
| Dratini           | 64     | Dragon | 94.0           |
| Dragonair         | 84     | Dragon | 94.0           |
| Gastly            | 35     | Ghost  | 53.75          |
| Haunter           | 50     | Ghost  | 53.75          |
| Gengar            | 65     | Ghost  | 53.75          |
| GengarMega Gengar | 65     | Ghost  | 53.75          |
| Jynx              | 50     | Ice    | 67.5           |
| Articuno          | 85     | Ice    | 67.5           |
| Kabuto            | 80     | Rock   | 87.5           |
+-------------------+--------+--------+----------------+
Data truncated.

可用函数

  • 排名函数:rank()、dense_rank()、ntile()、row_number()。
  • 分析函数:cume_dist()、percent_rank()、lag()、lead()。
  • 聚合函数:全部聚合函数都可以用作窗口函数。

用户自定义窗口函数

继承 WindowEvaluator 并通过 udwf() 注册,就能把自定义窗口函数交给引擎执行。求值器接口及完整示例见 datafusion.user_defined。原文在这里提供扩展入口,没有给出完整自定义求值器实现。

序列化时,Python 窗口 UDF 会内嵌在经 pickle 或 to_bytes() 序列化的表达式中。求值器类通过 cloudpickle 按值捕获,因此工作进程无需预先注册 UDF;求值器通过 import 解析的名称则按引用捕获,接收端工作进程必须能够导入相应模块。完整 IPC 模型与安全注意事项见 datafusion.ipc。

补充核验:当前 main 分支的窗口选项规则

以下内容译自已固定提交的官方 windows.md,属于比在线页面更新的补充说明。将窗口函数的 partition_by、order_by 与 window_frame 集中在一处设置:使用函数关键字参数、一次 over(),或者以 build() 结束的一条构建链。若窗口函数已经有其中任意选项,再追加构建方法或 over() 会报错。

# 报错:lead already has window options (order_by)
f.lead(col("v"), order_by="t").over(Window(partition_by=[col("g")]))

# 应在同一处设置
f.lead(col("v")).over(Window(partition_by=[col("g")], order_by="t"))

构建好的窗口函数只保存具体窗口帧,不记录它是用户选择的还是由排序派生的,因此无法在不猜测意图的情况下合并选项。原文跟踪 apache/datafusion#25934,未来解决后或可支持合并。已有的 null_treatment 会保留;聚合函数作为窗口函数时,其 filter 和 distinct 也会保留。整个分区的帧被视作“没有帧”,所以添加排序时会派生累计帧,即使先前显式传入过整个分区帧。

whole = WindowFrame("rows", None, None)

# 两者都得到累计求和
f.sum(col("v")).over(Window()).order_by(col("v")).build()
f.sum(col("v")).over(Window(window_frame=whole)).order_by(col("v")).build()

# 若要保留整个分区帧,应在排序之后设置,或同时传给 Window
f.sum(col("v")).over(Window()).order_by(col("v")).window_frame(whole).build()
f.sum(col("v")).over(Window(order_by=col("v"), window_frame=whole))

over() 会保留聚合函数构建时指定的 filter、distinct 与 null_treatment。例如,下面对值 1.0 和 4.0 去重后求均值:

f.avg(col("v"), distinct=True).over(Window())

聚合调用自身的 order_by 与 OVER 组合会报错,与 SQL 中聚合调用内部的 ORDER BY 情况一致。窗口的排序只设置窗口帧和行顺序,不会把排序传给聚合函数。例外是 percentile_cont 这类 WITHIN GROUP 函数:作为窗口函数时按升序计算,接受升序 sort_expression,拒绝降序。

f.percentile_cont(col("v"), 0.25).over(Window())  # 接受:升序
f.percentile_cont(col("v").sort(ascending=False), 0.25).over(Window())  # 报错

来源、版本与许可说明

官方 windows.md 带 ASF Apache-2.0 许可头,原文注明 Pokémon 样例来自 Ritchie Vink。

本译稿保留原文的技术步骤、示例与限制,并以“译注”或“补充核验”区分修正和新增说明。示例只做静态审查,未安装、运行或测试。

伴随源码核验提交:apache/datafusion-python @ a2ddb5531ac8。在线文档可能与源码构建时间不同,文中已注明实质差异。

原文许可与版权声明

Apache DataFusion 贡献者的窗口函数文档及其中代码按 Apache License 2.0 提供。以下原许可头逐字保留。本文为中文翻译,编校补充已单独标示,修改日期为 2026-10-05。Pokémon 示例数据来源:Ritchie Vink。

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

请登录后发表评论

    暂无评论内容