
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_dateAirflow 的分区资产调度中一个分区感知partition-aware的消费者 DAG 由PartitionedAssetTimetable驱动上游资产事件携带partition_keymapper 将其映射为消费者侧的分区键调度器据此为每个分区创建独立的 DagRun。在partition_date透传引入之前DagRun 上的两个分区字段语义并不对称partition_key始终存在是分区身份的字符串标识partition_date日期形状的分区锚点period-start datetime此前只有当消费者的 mapper 能从键本身解码出时间时才能被调度器推导出来。对于IdentityMapper这类键即原样透传的 mapper键本身不含时间语义调度器无法反推出日期——即使 producer 侧明明知道这批数据属于哪个日期。本功能解决的就是这一缺口把 producer DagRun 已经携带的partition_date一路带到消费者 DagRun再暴露给任务模板。变更说明原文为Propagatepartition_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_idasset_model.id, dags_to_queuedags_to_queue, partition_keypartition_key, partition_datepartition_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随 AssetPartitionDagRunAPDR 行一起持久化对应数据库迁移 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, totalFalse): 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_idclean_and_combine_player_stats, schedulePartitionedAssetTimetable( assetsteam_a_player_stats team_b_player_stats team_c_player_stats, default_partition_mapperStartOfHourMapper(), ), catchupFalse, ): task(outlets[combined_player_stats]) def combine_player_stats(dag_runNone): 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 委托到时间型子 mapperbase.py L113-L124、temporal.py L205producer 透传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.pycarry_partition_date的默认行为返回None与 identity 透传行为test_manager.pyregister_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),仅供参考