Apache Airflow 分区资产:将 producer 的 partition_date 透传到消费者 DAG 运行与任务模板
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本篇技术文章基于 Airflow 仓库中的变更说明 67285.feature.rst 展开。该变更让分区资产(partitioned assets)的 producer DagRun 上的partition_date能够透传到下游消费者 DagRun,使日期形状的分区值在消费者的任务模板上下文中可直接使用。读完本文,你将掌握:partition_date从产出侧到消费侧的完整数据流、时间型 mapper 与IdentityMapper两种日期解析路径的分工、冲突时的降级策略,以及如何在消费者 DAG 的 Jinja 模板与代码中实际读取该值。
功能背景:分区资产的 partition_key 与 partition_date
Airflow 的分区资产调度中,一个分区感知(partition-aware)的消费者 DAG 由PartitionedAssetTimetable驱动:上游资产事件携带partition_key,mapper 将其映射为消费者侧的分区键,调度器据此为每个分区创建独立的 DagRun。在partition_date透传引入之前,DagRun 上的两个分区字段语义并不对称:
partition_key:始终存在,是分区身份的字符串标识;partition_date:日期形状的分区锚点(period-start datetime),此前只有当消费者的 mapper 能从键本身解码出时间时才能被调度器推导出来。
对于IdentityMapper这类"键即原样透传"的 mapper,键本身不含时间语义,调度器无法反推出日期——即使 producer 侧明明知道这批数据属于哪个日期。本功能解决的就是这一缺口:把 producer DagRun 已经携带的partition_date一路带到消费者 DagRun,再暴露给任务模板。变更说明原文为:
Propagate
partition_datefrom producer DagRun to consumers of partitioned assets, so date-shaped partitions are available in consumer task templates.
完整数据流:从产出侧事件到消费者 DagRun
整条透传链路可以在源码中逐段印证:
1. 入口:register_asset_change 接收 partition_date
当产出任务完成并登记资产变更时,AssetManager.register_asset_change 的方法签名中新增了partition_date: datetime | None参数(L281),其取值即产出侧 DagRun 的partition_date。该方法随后把它连同partition_key一起交给内部队列逻辑:
cls._queue_dagruns( asset_id=asset_model.id, dags_to_queue=dags_to_queue, partition_key=partition_key, partition_date=partition_date, ... )见 manager.py 的 _queue_dagruns 调用处。
2. 分流:只有分区感知消费者参与透传
_queue_dagruns 将待排队 DAG 按timetable_partitioned属性切分为两组:分区 DAG 走 _queue_partitioned_dags,非分区 DAG 走既有的AssetDagRunQueue路径,partition_date对后者无意义。
3. 关键决策:mapper 的 carry_partition_date 是否"承运"日期
在_queue_partitioned_dags中,每个目标 DAG 的 mapper 会被调用一次:
target_partition_date = mapper.carry_partition_date(partition_date)见 manager.py L661-L680。这里体现了一条清晰的设计契约(源码注释原文明确了它):
IdentityMapper:覆写了 carry_partition_date,把source_partition_date原样返回——因为消费者键与生产者键相同且不含时间语义,调度器后续无法从键反推日期,只能靠"承运";- 时间型/复合型 mapper:基类 PartitionMapper.carry_partition_date 默认返回
None——消费者的日期会在建运行时由to_partition_date从键本身解码得出,无需承运; - 容错:自定义 mapper 的
carry_partition_date若抛出异常,管理器会记录日志并降级为None(消费者仍按partition_key排队,只是没有日期),不会中断整个写库流程。
4. 落库:APDR 行保存承运的 partition_date
承运得到的target_partition_date随 AssetPartitionDagRun(APDR) 行一起持久化,对应数据库迁移 0123_3_3_0_add_partition_date_to_asset_partition_dag_run.py 新增的partition_date列。
5. 建运行时:调度器最终裁决
调度器在为待处理的 APDR 创建消费者 DagRun 前,调用 _resolve_partition_date 做最终裁决,优先级如下:
- 时间型 mapper 优先:遍历贡献该分区的全部上游资产,各自用
mapper.to_partition_date(partition_key)解码锚点;所有时间型 mapper 解码出同一时刻(按 timezone-aware 的瞬时比较)时,该锚点即为 DagRun 的partition_date; - 回退到承运日期:若没有任何时间型 mapper 贡献锚点(典型即"全
IdentityMapper喂入"的场景),返回 APDR 上承运的carried_partition_date; - 冲突置空:时间型 mapper 解码出互不相同的锚点(例如不同资产对同一键使用了不同时间区配置的 mapper),记录警告并返回
None——此时刻意不用承运日期顶替,避免掩盖被记录的抑制事件。
DagRun 创建代码见 scheduler_job_runner.py L2429-L2445。
多来源并发下的冲突调和策略
分区消费者的多个上游资产可能对同一个 APDR先后贡献事件,_get_or_create_apdr 对已存在的 pending APDR 做了 best-effort 调和(L749-L779):
| 场景 | 处理 |
|---|---|
| APDR 上尚无日期,新事件携带日期 | 采纳新事件的日期(避免后续 identity 事件的日期被丢弃) |
| APDR 已有日期,新事件携带不同的日期 | 视为上游资产对分区时间不一致:将承运日期抑制为None并记录警告,让消费者 DagRun 得到None而非一个顺序依赖的、不稳定的值 |
| 两者一致,或新事件不携带日期 | 保留已有值 |
这一策略的核心思想是:与其让消费者拿到一个"看事件到达顺序"的日期,不如让partition_date明确为空、由用户从日志排查上游 mapper 配置。
在消费者侧使用 partition_date
透传的最终收益落在任务模板与任务代码上。Task SDK 的 Jinja2 模板上下文中声明了partition_date字段:
class Context(TypedDict, total=False): """Jinja2 template context for task rendering.""" ... partition_key: NotRequired[str | None] partition_date: NotRequired[DateTime | None]见 task-sdk/src/airflow/sdk/definitions/context.py L66-L67。因此消费者任务中可以直接写{{ partition_date }}模板变量,或按仓库示例 DAG 的方式读取dag_run.partition_date属性,example_asset_partition.py 展示了典型用法:
with DAG( dag_id="clean_and_combine_player_stats", schedule=PartitionedAssetTimetable( assets=team_a_player_stats & team_b_player_stats & team_c_player_stats, default_partition_mapper=StartOfHourMapper(), ), catchup=False, ): @task(outlets=[combined_player_stats]) def combine_player_stats(dag_run=None): if TYPE_CHECKING: assert dag_run print(dag_run.partition_key, dag_run.partition_date) combine_player_stats()该示例同时还演示了IdentityMapper的兜底行为:当PartitionedAssetTimetable未指定partition_mapper时回退到IdentityMapper(见 example_asset_partition.py L115-L123)——这正是本次功能最典型的受益场景:上一跳的时间型 mapper(如StartOfHourMapper)解析出的日期,现在会沿 identity 链继续透传到下一跳消费者。
两种日期来源路径对照
从源码结构看,同一个消费者 DagRun 的partition_date有两个来源,按优先级排列:
| 来源 | 生效条件 | 对应源码 |
|---|---|---|
键解码(to_partition_date) | 消费者至少有一个上游使用时间型 mapper(或复合 mapper 委托到时间型子 mapper) | base.py L113-L124、temporal.py L205 |
producer 透传(carry_partition_date) | 全部上游均为IdentityMapper(键不含时间语义,无法解码) | identity.py L33-L37 |
时间型路径的键是"权威来源",调度器可在任意时刻重新推导,因此始终优先;透传路径填补的是 identity 链上"解码不出来"的空档。自定义 mapper 开发者可以按需覆写这两个方法之一:需要解码语义就实现to_partition_date,需要承运语义就实现carry_partition_date,且注意基类对decode_downstream/encode_upstream成对覆写的强制校验(base.py L56-L67),避免 rollup 窗口永不满足导致运行被永久挂起。
测试与验证入口
该功能的实现事实可进一步通过仓库中的测试用例核对:
- test_identity.py 与 test_base.py:
carry_partition_date的默认行为(返回None)与 identity 透传行为; - test_manager.py:
register_asset_change/ APDR 的partition_date落库与冲突抑制逻辑; - test_scheduler_job.py:
_resolve_partition_date的时间型优先、identity 回退与冲突置空三类裁决。
小结
本次变更以最小侵入的方式打通了分区资产时间语义的"最后一公里":partition_date从 producer DagRun 出发,经register_asset_change入参、mapper 的carry_partition_date决策、APDR 持久化,最终由调度器在_resolve_partition_date中与键解码结果合并裁决,落到消费者 DagRun 上,并进入任务模板上下文。对于依赖IdentityMapper或混合 mapper 的多跳分区流水线,partition_date不再因中间一跳无法解码而断链,日期形状的分区从此可以在整条消费链的任务模板中稳定使用;遇到上游时间语义冲突时,系统选择显式置空并记录警告,保证值不"看起来对但实际不稳"。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考