Apache Airflow 资产触发的分区 Dag Run 中 partition_date 的解析与设置机制
2026/9/10 0:16:15 网站建设 项目流程

Apache Airflow 资产触发的分区 Dag Run 中 partition_date 的解析与设置机制

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

导读

在 Apache Airflow 的资产(Asset)驱动调度体系中,分区映射器(Partition Mapper)负责把上游资产事件中的分区键(partition key)映射到下游 Dag Run 的分区。本文围绕仓库中 airflow-core/newsfragments/68266.bugfix.rst 记录的缺陷修复展开:当消费者的分区映射器为时间型(temporal)映射器,或将其包裹在RollupMapper/FanOutMapper/ChainMapper等复合映射器中时,资产触发的分区 Dag Run 现在会正确设置partition_date;而非时间型映射器则保持partition_date为空。读完本文,你将理解partition_date的来源、各类映射器的解析委托链、调度器(Scheduler)的冲突裁决逻辑,以及对应的数据库迁移与测试验证。

背景:资产触发的分区 Dag Run 与 partition_date

Airflow 3.x 引入了资产(Asset)驱动的调度模型,配合分区映射器可以实现细粒度的"事件驱动 + 分区聚合"工作流。当一个下游 Dag 使用PartitionedAssetTimetable消费上游资产时,调度器会为每个满足条件的分区创建一条asset_partition_dag_run(APDR)记录,待分区满足后再据此创建实际的 Dag Run。

partition_date是该模型中的关键时间语义字段:它表示一个分区的时间锚点(period-start datetime),即该分区所代表的时间窗口的起始时刻。它被冻结在 APDR 创建时刻,从而保证消费者 Dag Run 的partition_date与产生其partition_key的分区映射器保持一致——这一点在该字段对应的数据库迁移文档注释中明确说明(见 0123_3_3_0_add_partition_date_to_asset_partition_dag_run.py)。

本次修复(PR 68266)的核心内容为:

Asset-triggered partitioned Dag runs now setpartition_datewhen the consumer's partition mapper is temporal (directly, or wrapped inRollupMapper/FanOutMapper/ChainMapper). Non-temporal mappers leavepartition_dateunset.

即:只有当消费者的分区映射器属于时间型(直接使用,或被RollupMapper/FanOutMapper/ChainMapper包裹)时,partition_date才会被设置;非时间型映射器(如IdentityMapperAllowedKeyMapperFixedKeyMapper)则保持partition_date为空。

partition_date 的解析入口:to_partition_date

partition_date的解析入口是PartitionMapper基类中定义的to_partition_date方法(见 base.py):

def to_partition_date(self, downstream_key: str) -> datetime | None: """ Return the temporal anchor (period-start datetime) for *downstream_key*. The scheduler stamps this on the asset-triggered Dag run as its ``partition_date``. The base implementation returns ``None`` — a plain partition key carries no temporal meaning. Temporal mappers override to decode the key into its window anchor; composite mappers (:class:`RollupMapper`, FanOutMapper) delegate to whichever child owns the downstream key's identity. """ return None

基类默认返回None——普通的(非时间型)分区键不携带时间语义,因此无法也不应推导出partition_date。这正是本修复中"非时间型映射器保持partition_date为空"的实现基础。

时间型映射器的实现

时间型映射器统一继承自_BaseTemporalMapper(见 temporal.py),其expected_decoded_type声明为datetime,并覆写了to_partition_date

def to_partition_date(self, downstream_key: str) -> datetime: anchor = self.normalize(self.decode_downstream(downstream_key)) # decode_downstream returns a naive datetime; localise it with the mapper's # own timezone, mirroring to_downstream, so the stored instant is correct. if anchor.tzinfo is None: anchor = make_aware(anchor, self._timezone) return anchor

解析流程分三步:先用decode_downstream从格式化后的下游分区键还原出分区起始时间(naive datetime),再经normalize归一到分区起点,最后用映射器自身的timezone本地化,得到带时区信息的锚点时间。

时间型映射器内置了六种常用粒度(同文件temporal.py):

映射器归一化语义默认 output_format示例
StartOfHourMapper取小时起点%Y-%m-%dT%H2024-03-13T10:42:152024-03-13T10
StartOfDayMapper取天起点%Y-%m-%d2024-03-13T10:42:152024-03-13
StartOfWeekMapper取 ISO 周周一%Y-%m-%d (W%V)2024-03-11 (W11)
StartOfMonthMapper取月首日%Y-%m2024-03
StartOfQuarterMapper取季度首日%Y-Q{quarter}2024-Q1
StartOfYearMapper取 1 月 1 日%Y2024

以 example_asset_partition.py 中的示例为例,StartOfDayMapper会把上游每小时的%Y-%m-%dT%H:%M:%S时间戳归一化为天起点%Y-%m-%d,再配合DayWindow声明"需等待全部 24 个小时分区到齐"。

复合映射器的委托链

本修复的关键点在于:当时间型映射器被复合映射器包裹时,partition_date依然能被正确解析。三种复合映射器各自委托给"拥有下游分区键身份"的子映射器:

  • RollupMapper(N→1 聚合):下游键由upstream_mapper.to_downstream产生,因此锚点由upstream_mapper解析(见 base.py):
def to_partition_date(self, downstream_key: str) -> datetime | None: # The downstream key is in upstream_mapper's format (to_downstream delegates # to it), so the anchor is the upstream_mapper's to resolve. return self.upstream_mapper.to_partition_date(downstream_key)
  • FanOutMapper(1→N 扇出):下游键由downstream_mapper格式化,因此锚点由downstream_mapper解析(见 temporal.py)。

  • ChainMapper(顺序链式):链中最后一个映射器负责格式化最终下游键,因此锚点由self.mappers[-1]解析(见 chain.py)。

此外,基类还提供了carry_partition_date(见 base.py),用于在事件入队时把生产者的partition_date携带到消费者的 APDR 上。基类默认返回None,只有IdentityMapper覆写为透传——因为它的下游键与上游键相同,而该键本身无法解码出时间含义。对应测试见 test_base.py。

调度器侧的解析与冲突裁决

调度器在创建资产触发的分区 Dag Run 时,会调用_resolve_partition_date(见 scheduler_job_runner.py)来解析partition_date,其裁决逻辑分为四档:

  1. 时间型映射器产出锚点:遍历为该分区键贡献的所有上游资产(name, uri),通过timetable.get_partition_mapper(name, uri)取得映射器并调用to_partition_date。由于分区消费者只有一个分区身份,所有时间型映射器必须把同一分区键解析到同一时刻(按时间戳比较,时区感知的等价时刻会收敛为一个)。
  2. 无任何时间型映射器:当anchors集合为空(例如全部由IdentityMapper供数)时,回退使用 APDR 在入队时携带的生产者日期carried_partition_date;一旦存在时间型映射器解析出的锚点,则以键为权威来源优先采用。
  3. 锚点冲突:若多个时间型映射器对同一键解析出不同的时间(例如不同资产配置了不同时区的映射器,属于配置错误),调度器记录警告日志并返回None故意不用携带日期替代,以免掩盖被抑制的错误。
  4. 异常兜底:任一映射器抛出异常时,记录日志并返回None,保证故障映射器不会导致整个调度器 tick 崩溃。

同时,scheduler_job_runner.py 中对 rollup 状态求值的异常也会被捕获,并提示"这通常表明分区映射器配置有误"。

数据模型与迁移

partition_date作为asset_partition_dag_run表的新增列,由迁移脚本 0123_3_3_0_add_partition_date_to_asset_partition_dag_run.py 引入,目标版本为 Airflow 3.3.0:

def upgrade(): """Add partition_date column to asset_partition_dag_run.""" with op.batch_alter_table("asset_partition_dag_run", schema=None) as batch_op: batch_op.add_column(sa.Column("partition_date", UtcDateTime, nullable=True)) def downgrade(): """Remove partition_date column from asset_partition_dag_run.""" with op.batch_alter_table("asset_partition_dag_run", schema=None) as batch_op: batch_op.drop_column("partition_date")

该列使用UtcDateTime类型且允许为空(nullable)——空值正是"非时间型映射器下partition_date未设置"这一语义在数据层的体现。partition_datepartition_key共同描述一个分区的身份:partition_key是字符串形式的分区标识,partition_date是其可被日历/UI 理解的时间锚点。

测试验证

本次修复行为在单元测试中有多处直接覆盖:

  • test_chain.py 的test_to_partition_date_delegates_to_last_mapper验证ChainMapperto_partition_date委托给链中最后一个映射器;
  • test_product.py 验证ProductMapper(笛卡尔积组合映射器)同样能正确解析to_partition_date
  • test_temporal.py 的test_to_partition_date_uses_mapper_timezone验证锚点使用映射器自身的_timezone本地化,而非全局默认时区;同文件的复合映射器测试(L427-L449)验证各复合映射器的委托行为;
  • test_scheduler_job.py 则在调度器层面验证时间型与复合映射器经由to_partition_date被正确解析。

实践要点总结

  1. 判断一个映射器是否为时间型:看它是否覆写了to_partition_date(或expected_decoded_type是否为datetime)。内置的StartOfHour/Day/Week/Month/Quarter/YearMapper均为时间型;IdentityMapperAllowedKeyMapperFixedKeyMapper等为非时间型。
  2. 复合映射器不会丢失时间语义RollupMapper委托upstream_mapperFanOutMapper委托downstream_mapperChainMapper委托最后一个子映射器,因此包裹在其中的时间型映射器依然能为 Dag Run 提供partition_date
  3. 多资产供数需保持时间一致:同一消费者的多个上游资产如果配置了不同时区或粒度的映射器,可能导致锚点冲突,调度器会以日志警告并置空partition_date的方式暴露配置错误。
  4. 查看数据asset_partition_dag_run.partition_date为可空的UtcDateTime列(Airflow 3.3.0 起),非时间型映射器触发的运行该字段为NULL

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询