Apache Airflow:为 RollupMapper 配置 wait_policy,让分区汇总在部分窗口满足时提前触发
Apache Airflow为 RollupMapper 配置 wait_policy让分区汇总在部分窗口满足时提前触发【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow在 Airflow 的资产Asset分区调度中RollupMapper负责把多个上游分区键汇总rollup到一个下游分区上。此前的默认行为是死等窗口内全部上游键到齐才触发下游本次特性见 66848.feature.rst为RollupMapper新增了wait_policy参数WaitForAll()默认在所有键到齐时触发MinimumCount(n)则在部分窗口partial window满足时提前触发。读完全文你可以掌握两种等待策略的语义细节、调度器的判定链路、以及如何在 DAG 中编写可运行的提前触发汇总逻辑。背景RollupMapper 的死等问题Airflow 支持基于分区的资产依赖上游 DAG 每产出一个分区键如某小时、某区域下游的PartitionedAssetTimetable就需要判断这一批键是否足以触发我的一次运行。RollupMapper由三个组件组合而成base.pyupstream_mapper把每个上游键归一化到下游粒度例如StartOfDayMapper把小时键归到当天window声明一个下游键需要哪些上游键例如DayWindow声明一天由 24 个小时组成实现见 window.pywait_policy决定在期望窗口与实际已到达键之间何时触发下游运行。在加入wait_policy之前汇总语义只有一种等窗口内所有期望的上游键都到齐才触发。这在生产上有一个明显痛点——只要有一个慢分区或丢失分区下游运行就会被无限期挂起。wait_policy参数正是为此引入的。RollupMapper.__init__的签名base.pydef __init__( self, *, upstream_mapper: PartitionMapper, window: Window, wait_policy: WaitPolicy | None None, max_downstream_keys: int | None None, ) - None: ... if wait_policy is None: wait_policy WaitForAll()不传wait_policy时默认构造WaitForAll()完全保持旧行为存量 DAG 不受影响。两种内置策略的语义策略类定义在 wait_policy.pySDK 侧对应实现在 task-sdk/.../partition_mappers/wait_policy.py用户通过airflow.sdk导入同名类。WaitForAll等到全齐才触发attrs.define(frozenTrue) class WaitForAll(WaitPolicy): def is_satisfied(self, matched: int, expected: int) - bool: return matched expected仅当matched expected窗口内每个期望键都已到达才触发空窗口expected 0按空真vacuously satisfied处理直接视为满足它实现上对键集合做子集短路判断expected matched避免物化完整交集wait_policy.py。MinimumCount绝对下界与相对缺口的两种写法MinimumCount(n)的参数n支持正负两种写法wait_policy.py取值触发条件说明n 0matched n绝对下界到齐 n 个期望键即触发n 0matched max(0, expected n)相对写法至多允许-n个键缺失n 0构造即抛ValueError退化配置永远触发哪怕空窗口被显式拒绝例如对 3 个区域的窗口MinimumCount(2)与MinimumCount(-1)至多缺 1 个等价。负值写法的好处是阈值随窗口大小自适应——窗口从 24 小时扩成别的规模时至多缺 k 个的语义依然成立。n 0的max(0, ...)钳位保证阈值永不为负空窗口下matched 0 0仍满足空真语义。PartitionSatisfaction同时报告满足与永久不可达调度器关心的不只是现在够不够还有这个配置是否永远不可能满足。策略的统一入口是is_satisfied_by_keys()返回结构化结果wait_policy.pyattrs.define(frozenTrue) class PartitionSatisfaction: satisfied: bool # 是否达到触发阈值 unreachable: bool # 阈值是否因窗口基数而永不可达 unreachable_reason: str | None # 不可达时的人读原因MinimumCount.is_unreachable(expected)的逻辑是仅当n 0且n expected时不可达——例如对 60 键的HourWindow配MinimumCount(61)无论上游来多少事件都不可能触发负n经钳位后有上界expected因此永远不会不可达wait_policy.py。调度器如何消费 wait_policy调度链路在 scheduler_job_runner.py 的_resolve_asset_partition_status中对 rollup 类资产调度器取该分区 DAG 运行APDR的partition_key调用mapper.to_upstream(...)得到期望键集合再与实际到达的键比对expected mapper.to_upstream(apdr.partition_key) actual actual_by_asset.get(asset_id, set()) result mapper.wait_policy.is_satisfied_by_keys(matchedactual, expectedexpected) if result.unreachable: self._warn_unreachable_asset_partition(apdrapdr, namename, uriuri, reasonresult.unreachable_reason) return False return result.satisfied这里有三个值得注意的工程细节不可达只告警、不触发策略把永不可达的原因串构造好unreachable_reason调度器负责去重后经_warn_unreachable_asset_partition转发到日志运行保持挂起状态返回False——策略拥有消息内容调度器拥有转发策略见 scheduler_job_runner.py。配置错误不拖垮调度循环mapper 抛异常时按尚未满足处理异常以ERROR级别记入调度日志便于定位误配置。非 rollup 资产不受影响非 rollup 资产只要已有事件日志即视为满足wait_policy只作用于RollupMapper。实战示例多区域分类汇总提前触发官方文档 assets.rst 的 Wait policies 一节给出了完整用法与示例 DAG example_asset_partition.py 中的segment_region_stats_early_rollup对应。场景是上游资产按us/eu/apac三个区域分区产出下游汇总 DAG 不必等三个区域全到到齐任意两个即触发from airflow.sdk import ( DAG, Asset, FixedKeyMapper, MinimumCount, PartitionedAtRuntime, PartitionedAssetTimetable, RollupMapper, SegmentWindow, asset, ) asset( urifile://incoming/player-stats/multi-region.csv, schedulePartitionedAtRuntime(), ) def multi_region_player_stats(self, outlet_events): outlet_events[self].add_partitions([us, eu, apac]) # Consumer: fires once at least two of the three declared region partitions arrive. with DAG( dag_idsegment_region_stats_early_rollup, schedulePartitionedAssetTimetable( assetsAsset.ref(namemulti_region_player_stats), default_partition_mapperRollupMapper( upstream_mapperFixedKeyMapper(all_regions), windowSegmentWindow([us, eu, apac]), wait_policyMinimumCount(2), ), ), catchupFalse, ): ...组合中各部件的职责SegmentWindow([us, eu, apac])声明一个下游周期包含的固定字符串键集合window.pyFixedKeyMapper(all_regions)把所有上游键折叠到单一下游分区键wait_policyMinimumCount(2)把默认的全等语义改为至少到 2 个。文档同时提醒MinimumCount(-1)是同一阈值的相对写法至多缺 1 个对 3 成员窗口与MinimumCount(2)等价如果希望显式表达意图而非依赖默认可以显式传wait_policyWaitForAll()——示例 DAG 中时间汇总部分正是这样做的example_asset_partition.pyStartOfDayMapperDayWindowWaitForAll等待 24 个小时分区全部到齐。另外需要注意适用边界文档指出分类汇总Segment rollup的完整语义在 WAIT_FOR_ALL 语义即默认值下才有意义改用MinimumCount后应确保下游逻辑能容忍缺失分区。序列化与自定义策略的限制wait_policy与upstream_mapper、window一样参与 DAG 序列化RollupMapper.serialize()中通过encode_wait_policy(self.wait_policy)编码base.py编码实现见 encoders.py。两个约束只接受内置策略encode_wait_policy走BUILTIN_WAIT_POLICIES快路径仅识别内置WaitForAll/MinimumCount自定义WaitPolicy子类会抛出WaitPolicyNotSupported。测试 test_rollup_wait_policy.py 中的test_non_builtin_wait_policy_rejected验证了这一点往返一致性测试覆盖了MinimumCount(3)、MinimumCount(-2)的序列化/反序列化往返确认反序列化后实例与原对象相等且MinimumCount.n这类载荷字段不会被静默丢弃。从测试用例看边界行为tests/unit/partition_mappers/test_rollup_wait_policy.py 系统性地验证了本特性的边界MinimumCount(0)构造即抛ValueError(MinimumCount(0) is degenerate)提示改用WaitForAll()或n ! 0不传wait_policy时mapper.wait_policy为WaitForAll实例触发阈值矩阵MinimumCount(5)在期望 60 的窗口中5 个匹配即触发、4 个不触发MinimumCount(-3)在 57/60 匹配时触发、56/60 不触发空窗口下MinimumCount(5)不触发而MinimumCount(-3)因钳位触发不可达判定MinimumCount(5)对期望 4 不可达、对期望 5 可达MinimumCount(61)对 60 键的HourWindow不可达且原因串为wait policy MinimumCount(61) can never be satisfied given the windows cardinality 60——与调度器告警内容一致MinimumCount继承基类默认路径集合 → 计数 →is_satisfiedWaitForAll则覆写为集合级短路判断。小结wait_policy是RollupMapper上的一处小而关键的扩展点默认WaitForAll()保证与既有行为完全兼容MinimumCount(n)提供绝对下界n 0与相对缺口n 0两种写法n 0被显式拒绝策略层额外返回永久不可达信号调度器将其转成一次性告警避免误配置如 61 60导致的静默挂起仅内置策略可参与 DAG 序列化自定义策略在编码阶段即被拒绝。适用前提以上行为以当前仓库代码为准涉及资产分区汇总Asset partitioning RollupMapper的调度场景非 rollup 资产及未使用RollupMapper的 DAG 不受影响。可进一步延伸阅读 assets.rst 中分区资产的完整章节与 example_asset_partition.py 中的时间汇总与分类汇总两类完整示例。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考