题干与适用场景
上游任务每天产生多个数据资产,下游报表和质量检查希望在依赖资产更新后自动运行。请说明如何用 Airflow 的 asset-aware scheduling 建模生产者、消费者、分区和失败恢复,并解释它与时间调度、外部传感器的边界。
面试官考察点
- 是否把 asset 当成可被任务更新的逻辑数据依赖,而不是简单的文件名或 cron 标签。
- 是否能区分 asset 事件、DAG 的时间表、表达式组合和事件触发器。
- 是否考虑重复事件、迟到数据、分区粒度、数据质量门禁和回填影响。
- 是否能说明监控、权限、幂等、失败重试与暂停恢复策略。
回答前需要澄清的问题
- 资产代表整张表、某个分区还是对象存储路径?更新事件是否带分区和批次信息?
- 下游需要所有上游资产都更新,还是任一资产更新即可?是否存在跨 DAG 的资产表达式?
- 事件到达但数据质量检查失败时,应该阻止消费 DAG、重试生产者,还是允许人工放行?
- 需要补历史数据、重放事件或兼容现有 cron 调度吗?重复触发的代价是什么?
30 秒回答框架
我会先把业务数据产品建模成资产及其更新契约,再让生产 DAG 在成功写入并完成质量检查后发出资产更新事件。消费 DAG 使用资产依赖或表达式决定触发条件,而不是用更短的 cron 猜测数据是否就绪。随后定义分区、幂等、重复事件、迟到更新和回填规则,并监控事件延迟、待触发 DAG、失败重试和数据新鲜度。若事件来自外部系统,则评估 Airflow event-driven trigger 的可恢复性与安全边界。
分步骤深入解答
1. 建立资产契约
在 Airflow Assets 官方文档中,资产是 DAG 之间共享的数据依赖。生产任务应在数据成功提交且契约满足后更新资产;仅创建临时文件、任务开始或部分写入不应被当成可消费事件。资产 URI、所有权和分区粒度需要稳定,否则下游无法判断新旧数据。
2. 选择触发逻辑
消费 DAG 可以依赖一个或多个资产。根据 Asset-Aware Scheduling 官方文档,多个资产的逻辑组合要表达“全部更新”或“任一更新”等业务条件;时间表仍可作为独立约束。设计时应明确事件和时间条件同时满足时的行为,避免把表达式误当成数据过滤器。
3. 处理分区与重复
资产事件不自动替代数据分区水位。事件应携带可追踪的批次或分区信息,消费任务用水位表、唯一键或事务写入保证幂等。重复事件、任务重试和调度器恢复都可能再次评估依赖;下游应能安全重跑,而不是依赖“每个事件只来一次”的假设。
4. 外部事件、质量与恢复
如果更新来自消息队列或外部系统,可参考 Event-driven scheduling 官方文档选择事件触发器,但要核对触发器是否支持可恢复检查、连接凭据和资源释放。数据质量失败时保留资产未就绪状态,告警并重试或人工放行;回填时明确是否发出新事件、是否隔离历史分区,以及如何防止消费 DAG 与实时更新互相覆盖。
高质量示范回答
我会先定义资产契约:每个资产的 URI、分区键、生产者、质量门槛和更新批次。生产 DAG 只有在原子写入和质量检查通过后才更新资产,消费 DAG 用资产依赖或表达式表达“全部上游完成”与“任一上游完成”,必要时再叠加时间表。事件中记录分区和批次,消费端用水位和幂等键处理重复、重试与调度器恢复。外部更新采用受控的事件触发器并限制凭据和资源生命周期。监控事件延迟、待触发任务、数据新鲜度和失败重放;回填与实时路径使用不同批次边界,防止历史重放覆盖最新结果。
常见错误
- 把 asset 当作任意文件路径,未定义谁在什么时刻拥有更新权限。
- 只缩短 cron 间隔,却没有数据就绪、分区和质量门槛。
- 认为资产事件天然只触发一次,忽略重试、重复发布和调度器恢复。
- 把资产依赖表达式当成数据内容过滤器,遗漏实际的分区水位判断。
- 外部事件触发器长期持有连接或凭据,却没有超时、释放和告警策略。
- 回填时直接复用实时 DAG,导致历史事件重复覆盖在线结果。
追问及应对
如果两个上游资产到达时间不同怎么办?
明确下游是等待全部资产,还是允许先处理部分结果。等待全部资产时记录每个分区的到达水位和超时告警;允许部分结果时输出版本号,让下游知道结果仍可被后续资产更新。
如何测试重复事件?
在测试环境重复发布同一资产更新、重试生产任务并重启调度器,验证消费任务的唯一键、水位和事务边界。检查结果行数、版本号和外部副作用,确保重跑不会重复发送通知或扣费。
什么时候保留外部传感器?
当依赖系统无法发布 Airflow 资产事件、需要轮询一个受控接口,或迁移期间必须兼容旧契约时可以暂时使用传感器。应记录迁移期限与轮询成本,最终让上游提供可验证的资产更新信号。