恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
Apache Airflow 分区窗口方向控制:`Window.Direction.FORWARD` 与 `BACKWARD` 详解
首页
资讯中心
/
Apache Airflow 分区窗口方向控制:`Window.Direction.FORWARD` 与 `BACKWARD` 详解
Apache Airflow 分区窗口方向控制:`Window.Direction.FORWARD` 与 `BACKWARD` 详解
发布时间:2026/9/10 7:05:25
Apache Airflow 分区窗口方向控制Window.Direction.FORWARD与BACKWARD详解【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读在 Apache Airflow 的资产分区Asset Partitioning体系中Window窗口类负责描述一个下游分区由哪些上游分区组成而本仓库的新特性见 67475.feature.rst为所有Window子类引入了一个direction关键字参数默认的Window.Direction.FORWARD从上游 key 所在时刻向前展开窗口周期而Window.Direction.BACKWARD则改为向上游 key 所在时刻收拢其之前紧邻的周期。本文将从该特性的 API 语义出发结合airflow-core的窗口实现、FanOutMapper扇形展开逻辑、单元测试与官方示例 DAG完整讲解direction的用法、底层实现与边界条件帮助你在编写PartitionedAssetTimetable时准确选择窗口展开方向。一、特性概述Window子类新增direction关键字按变更说明原文该特性核心是Window子类接受一个direction关键字 ——Window.Direction.FORWARD默认将周期从上游 key 开始向前展开时间上向未来传入directionWindow.Direction.BACKWARD例如WeekWindow(directionWindow.Direction.BACKWARD)则改为向上游 key 收拢其之前紧邻的周期时间上向过去。两个方向的含义可以精确定义为Window.Direction.FORWARD默认生成从上游 key 开始的那一个周期。例如上游 key 是某周一WeekWindow()展开为该周一及其后 6 天Window.Direction.BACKWARD生成以上游 key 为终点收拢的尾随周期即 FORWARD 的镜像。例如同一个周一 keyWeekWindow(directionWindow.Direction.BACKWARD)展开为以该周一为最后一天的 7 天。在airflow-core的源码中Direction是定义在抽象基类Window内部的枚举window.pyclass Direction(str, Enum): Direction of a Window fan-out relative to the upstream key. BACKWARD backward Yield the trailing period ending at the upstream key (the mirror of FORWARD). FORWARD forward Default; yield the period starting at the upstream key (forward in time).而Window.__init__通过self.direction self.Direction(direction)window.py完成类型收拢既可以直接传枚举成员也可以传字符串forward/backward传入拼写错误或大小写错误的值如forwrd、FORWARD会在构造时立即抛出ValueError而不是留到调度器 tick 中才暴露。这一行为由单元测试 test_window.py 中的TestDirectionValidation专门覆盖。二、方向在窗口展开中的底层实现2.1_build_directional_steps统一的步进生成器除MonthWindow外的所有时间窗口HourWindow、DayWindow、WeekWindow、QuarterWindow、YearWindow都通过辅助函数_build_directional_steps生成成员window.pydef _build_directional_steps( period_start: datetime, count: int, step: Callable[[datetime, int], datetime], direction: Window.Direction, ) - Iterable[datetime]: base step(period_start, -(count - 1)) if direction is Window.Direction.BACKWARD else period_start return (step(base, i) for i in range(count))其逻辑很直观FORWARD直接以period_start为起点用step依次推进count个周期起点BACKWARD先把基准点回退count - 1个单位再同样向前推进count个周期起点——这样得到的序列恰好以period_start为最后一个成员即尾随周期。各窗口传入的step与count分别为窗口类成员数步进函数FORWARD 语义BACKWARD 语义HourWindow60timedelta(minutesi)从整点开始的 60 个分钟起点以该整点为终点的 60 个分钟起点DayWindow24timedelta(hoursi)从零点开始的 24 个小时起点以该零点为终点的 24 个小时起点WeekWindow7timedelta(daysi)从该日开始的 7 个日起点以该日为终点的 7 个日起点QuarterWindow3_shift_months按月进位该季度起始月开始的 3 个月起点以该月为终点的 3 个月起点YearWindow12_shift_months按月进位该年 1 月开始的 12 个月起点以该月为终点的 12 个月起点以单元测试 test_fan_out.py 中test_fan_out_with_directional_window的断言为例上游 key 为 2024-03-04周一StartOfWeekMapper将其归一化为自身WeekWindow()FORWARD展开为2024-03-04…2024-03-10周一至周日WeekWindow(directionWindow.Direction.BACKWARD)展开为2024-02-27…2024-03-04含 2024 闰年 2 月 29 日的尾随 7 天最后一个成员恰好是上游 key 本身。2.2MonthWindow的特例BACKWARD 不是简单镜像MonthWindow不走_build_directional_steps因为其成员数不固定2831 天且 BACKWARD 语义是开闭区间(prev_month_start, anchor]而非先回退再前进的镜像window.pyif self.direction is Window.Direction.BACKWARD: # 尾随周期 上个月 1 号之后到 anchor 为止的每一天不含上个月 1 号含 anchor prev _shift_months(period_start, -1) days (period_start - prev).days return (prev timedelta(daysi 1) for i in range(days))这意味着MonthWindow(directionWindow.Direction.BACKWARD)以2024-03-01为锚点会展开出 2024-02-02 … 2024-03-012024 闰年共 29 个成员并不对齐自然月。该行为由测试 test_window.py 显式固定结果必须包含锚点 3 月 1 日且必须不包含2 月 1 日。2.3 月以上窗口的 day-1 前置条件MonthWindow、QuarterWindow、YearWindow的周期起点必须是当月 1 号_require_day_onewindow.py否则会抛出带明确提示的ValueErrorexpects a period start on day 1 of the month。内置的时间上游 Mapper如StartOfMonthMapper本身会把 key 归一化到 1 号但自定义decode_downstream可能返回任意日期这一检查把调度器 tick 中的莫名崩溃提前转化为对上游 Mapper 契约的显式约束。注意该前置条件与方向无关——MonthWindow(directionWindow.Direction.BACKWARD)传入非 1 号的锚点同样抛错见 test_window.py而HourWindow/DayWindow/WeekWindow的 BACKWARD 不受此限制。三、方向与FanOutMapper1 → N 扇形展开direction参数在FanOutMapper上游一个事件展开为下游多个 Dag run中的语义最直观。FanOutMapper由三部分组成temporal.pyupstream_mapper解析粗粒度上游 key 并归一化为周期起点window枚举该周期内的全部成员downstream_mapper把每个成员格式化为下游 key 字符串可省略按窗口类自动查表选择默认值。其方向语义在类文档中写得很清楚默认 FORWARD 展开上游 key 所代表的周期BACKWARD 展开以上游 key 为终点的尾随周期temporal.py。仓库自带示例 DAG example_asset_partition.py 给出了一个非常贴合实际场景的对照daily_inferenceDAG 使用WeekWindow()FORWARD每周模型产物周一 key扇出 7 个每日推理 run覆盖该周trailing_week_inferenceDAG 使用WeekWindow(directionWindow.Direction.BACKWARD)同样的周一 key 扇出 7 个每日 run但覆盖截至该周一的前 7 天例如用于对模型发布前一周的数据打分。两个 DAG 唯一区别就是direction正如源码注释所述Direction is the only difference between the two Dags.使用方式来自官方文档 assets.rstfrom airflow.sdk import FanOutMapper, PartitionedAssetTimetable, StartOfWeekMapper, WeekWindow, Window PartitionedAssetTimetable( assetsweekly_model_artifact, default_partition_mapperFanOutMapper( upstream_mapperStartOfWeekMapper(), windowWeekWindow(directionWindow.Direction.BACKWARD), ), )注意FanOutMapper的默认downstream_mapper查找表temporal.py按窗口类名解析方向不会影响默认解析WeekWindow(direction...)仍是WeekWindow自动解析为StartOfDayMapper。这一保证由测试 test_fan_out.pytest_fan_out_with_directional_window_resolves_default_downstream_mapper验证。此外方向信息会随窗口一起序列化serialize输出{direction: ...}window.pyDAG 经调度器序列化/反序列化后方向保持不变见 test_fan_out.py 的方向往返测试。四、方向与RollupMapperN → 1 汇总同样适用方向不限于扇形展开RollupMapper下游一个 run 等待多个上游分区全部到达复用同一批Window类因此同样支持direction。官方文档明确指出DayWindow(directionWindow.Direction.BACKWARD)会一直持有下游 run直到下游 key 对应午夜之前的 24 小时分区全部到达。也就是说汇总语义从等待该 key 之后的周期翻转为等待该 key 之前的尾随周期。无论是RollupMapper还是FanOutMapper窗口都只在上游 Mapper 的解码形式时间窗口为datetime上工作不接触 key 字符串、时区与格式——这些职责归上游 Mapperwindow.py。每个窗口类还通过expected_decoded_type声明其接受的解码类型RollupMapper.__init__会在配对错误如上游decode_downstream返回str而窗口需要datetime时直接拒绝组合。五、搭配时间 Mapper 与示例 DAG 的完整拼图5.1 与内置时间 Mapper 的配合direction的锚点是上游 key 归一化后的周期起点因此通常与StartOf*系列 Mapper 配合使用temporal.pyStartOfHourMapper2024-03-13T10:42:15→2024-03-13T10默认输出格式%Y-%m-%dT%HStartOfDayMapper→2024-03-13StartOfWeekMapper→ 所在 ISO 周的周一2024-03-11 (W11)StartOfMonthMapper→2024-03归一化到 1 号满足月级窗口的 day-1 前置条件StartOfQuarterMapper→2024-Q1StartOfYearMapper→2024。例如FanOutMapper(upstream_mapperStartOfWeekMapper(), windowWeekWindow(directionWindow.Direction.BACKWARD))周中任意时刻的 key如周三2024-01-17T13:42:00先被归一化为周一再展开为以该周一为终点的 7 个日 key。测试 test_fan_out.py 验证了周三输入仍正确产出周一到周日的完整序列FORWARD 场景。5.2 扇形展开的数量上限FanOutMapper支持max_downstream_keys限制单个上游事件最多创建的下游 run 数量assets.rst超过时这些 run不会被排队而是记录一条 partition fan-out exceeded 审计日志省略时回退到全局配置[scheduler] partition_mapper_max_downstream_keys默认 1000。WeekWindow恰好 7 个成员因此示例 DAG 中把max_downstream_keys7设在了临界值上example_asset_partition.py。该参数与direction相互独立均可随 Mapper 序列化往返见 test_fan_out.py 的编码往返测试。六、注意边界DST 与本地时区使用direction时仍需留意DayWindow文档中标注的 DST 边界window.py春令时跳转时钟前拨本地日实际不足 24 小时DayWindow产出的某个朴素小时步进落在时间空洞中被编码为下一个本地小时该 key 上游永远不会产生导致窗口永远无法满足秋令时回拨时钟重复本地日实际有 25 小时但DayWindow只枚举 24 步多出的那一小时的上游事件不会被任何汇总包含。DayWindow的算术基于朴素datetime步进不感知 DST缓解方案是使用 UTC 的input_format如%Y-%m-%dT%H%z并让上游生产者输出 UTC 分区 key。方向参数不改变这一前提——DayWindow(directionWindow.Direction.BACKWARD)同样以 24 个小时步进收拢到锚点而非回退到本地日历上的前一天。上述限制在 test_window.py 中以xfail测试显式固定为已知接受项。七、小结如何选择方向场景方向说明上游 key 代表一个周期的开始下游处理该周期本身FORWARD默认周度模型产物 → 该周 7 天的日级推理上游 key 代表一个事件的锚点下游处理其之前的周期BACKWARD周度模型产物 → 发布前 7 天的日级评分下游汇总需要等待key 之前的周期全部到达BACKWARDRollupMapperDayWindow(direction...)等待前 24 小时direction是本仓库分区机制中一个语义清晰、实现轻量的增强它不改变窗口的成员生成算法只是为相对锚点向前还是向后引入显式表达并通过枚举类型校验、序列化往返与测试矩阵保证了 DAG 作者可以安全地在FanOutMapper与RollupMapper中混用。深入源码可继续阅读窗口实现airflow-core/src/airflow/partition_mappers/window.py扇形展开与时间 Mapperairflow-core/src/airflow/partition_mappers/temporal.py方向行为测试airflow-core/tests/unit/partition_mappers/test_window.py 与 test_fan_out.py官方分区文档airflow-core/docs/authoring-and-scheduling/assets.rst完整示例 DAGairflow-core/src/airflow/example_dags/example_asset_partition.py【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考