Apache Airflow Shared-Stream 触发器生产者端 Ack 通道:让 Kafka / SQS / Pub/Sub / Service Bus 在事件被消费并持久化后安全提交
本指南基于 Apache Airflow 仓库中的新功能片段(airflow-core/newsfragments/67523.feature.rst)展开:当多个触发器共享同一条消息流(shared-stream)时,Airflow 现在支持在生产者侧增加一条 ack 通道,让 Kafka、SQS、Pub/Sub、Service Bus 等消息中间件在所有订阅者都已接受事件,且派生出的触发器事件已持久化到元数据库之后,才在 broker 句柄上执行 commit、delete 或 ack。读完本文,你将理解 ack 模式的完整工作流程、核心 API 与配置项,以及如何在自己的事件触发器上启用这一能力。
背景:Shared-Stream 触发器为什么需要 ack 通道
在 Airflow 的 triggerer 组件中,事件触发器(BaseEventTrigger)负责监听外部消息源。传统写法下,每个触发器实例各自运行一个独立的 run() 轮询循环,即使多个触发器监听的是同一个 topic / queue,也会产生重复的轮询连接与重复消费。
Shared-stream 机制改变了这一点:当多个触发器实例声明了相同的非 None 的 shared_stream_key 时,triggerer 会通过 SharedStreamManager 把它们路由到同一个共享轮询循环(shared poll loop)上——一条底层轮询产生原始事件,广播给所有参与该流的触发器,每个触发器再用自己的 filter_shared_stream 把原始事件转换成自己的 TriggerEvent(见 shared_stream.py 模块 docstring 与 base.py)。
但共享轮询带来一个新的难题:轮询循环从 broker 拉取事件后,什么时机才应该向 broker 确认(commit/delete/ack)? 如果过早确认,而某个订阅者还没来得及处理、或派生事件还没持久化,一旦 triggerer 崩溃,事件就永久丢失;如果永不确认,broker 会不断重复投递。
新功能片段描述的正是解法:为共享流增加生产者端 ack 通道(producer-side ack channel)。它的语义非常明确——只有当:
- 事件广播时在线的每一个订阅者都已"越过"该事件(拉取下一条原始事件,或取消订阅);
- 订阅者从该事件派生的所有触发器事件都已确认写入元数据库;
事件才算"完全解决(fully resolved)",随后才轮到生产者去调用 broker 句柄上的 commit / delete / ack。
快速路径与 ack 模式:两种运行方式
根据 shared_stream.py 的实现,共享流有两种运行模式:
- 快速路径(fast path):触发器只实现
open_shared_stream(打开共享流并逐条产出原始事件),不覆盖create_shared_stream_producer。此时没有事件 ID、没有解决(resolution)记账,订阅者拿到原始事件后直接跑filter_shared_stream。 - ack 模式(ack mode):触发器覆盖了类方法
create_shared_stream_producer,SharedStreamManager为该流切换到 ack 模式(源码通过遍历 MRO 判断子类是否覆盖了该方法,见_is_ack_required,shared_stream.py)。
两种模式下,订阅者看到的流形状完全一致——同一份 filter_shared_stream 代码在快速路径和 ack 模式下无需任何改动即可运行,因此把一个触发器切换到 ack 模式是向后兼容的(backward-compatible)。区别只在于"解决记账"与"broker 推进"两部分,这两者全部由 manager 从消费进度推导,不增加任何订阅者侧的 API。
ack 通道的核心流程
1. 工厂与流:create_shared_stream_producer + open_stream
ack 模式下,工厂方法 create_shared_stream_producer(kwargs) 在每个组只被调用一次,返回一个 SharedStreamProducer 实例,该实例在一次轮询的生命周期内独占 broker 连接(shared_stream.py)。
生产者实现两个抽象成员:
open_stream():以 async generator 形式打开 broker 连接,逐条产出(raw_event, broker_payload)元组。其中broker_payload是不透明的对象,用于后续在 broker 上推进——例如 Kafka 的 offset、SQS 的 receipt handle、Pub/Sub 的 ack ID。契约是"要么永远 yield,要么抛异常"——提前返回会被视为错误并传播给所有订阅者。advance(batch):接收一个批次(Sequence[AdvanceItem]),在 broker 上执行 commit / delete / ack(shared_stream.py)。
2. 广播与快照:snapshot-at-fan-out
manager 驱动生产者的 open_stream,每拿到一个原始事件,就把这个事件广播进每个订阅者的队列。广播时刻在线的订阅者集合被冻结为快照(snapshot),事件广播之后才加入的订阅者不会进入该事件的待解决集合(pending set)——这就是"fan-out 时刻快照"语义(shared_stream.py)。
每个事件在 _outstanding 中登记为一条待解决记录,记录其待解决的订阅者集合、broker_payload、所属 lane、创建时间以及每个订阅者的绑定状态(_SubscriberBinding,记录窗口是否关闭与未确认的序列号集合)。
3. 订阅者如何"解决"事件:binding window 与持久化门控
一个订阅者要在一个事件上被记为 acked,必须同时满足两个条件(shared_stream.py):
- 绑定窗口关闭(binding window closed):订阅者已越过该事件——具体地,
_ack_drain会在订阅者从队列里拉取下一条原始事件的那一刻关闭上一条事件的窗口(shared_stream.py)。窗口关闭就是"订阅者接受了该事件",无需显式调用。 - 所有派生触发器事件已确认持久化(persist confirmations):订阅者的 filter 从原始事件派生出的每个
TriggerEvent,在离开 runner 时被分配一个序列号(seq);supervisor 在事件写入元数据库后确认该 seq,确认消息在下一次状态同步时回到 runner(通常一个同步周期,约一两秒,远小于 ack 超时)。bind_pending_event把刚产出的触发器事件绑定到当前打开的窗口,confirm_persisted记录持久化确认;只有窗口关闭且无未确认 seq 时,订阅者才真正 resolve 为acked(shared_stream.py)。
这里有一个关键前提:filter 必须在拉取下一条原始事件之前 yield 出所有派生的触发器事件——这也是编写 filter 的自然方式,绑定关系正是依赖这一点建立(shared_stream.py)。如果 filter 直接跳过某事件(不 yield 任何东西),该事件会在 filter 循环回下一条原始事件时立即干净地解决。
4. 推进:advance(batch) 与 AdvanceOutcome
当广播时在线的每一个订阅者都已解决某个事件,该事件"完全解决",进入 broker 推进环节:manager 调用 await producer.advance(batch),批次是该 lane 上连续已解决的前缀(shared_stream.py)。
每个 AdvanceItem 携带事件的 broker_payload 与一个 AdvanceOutcome——包含每个订阅者的解决统计(shared_stream.py):
| 字段 | 含义 |
|---|---|
acked |
订阅者越过了事件且所有派生触发器事件已确认持久化 |
failed |
被 manager 强制失败(ack 超时——包括持久化确认始终未到达——或队列溢出) |
rejected |
订阅者通过 reject_shared_stream_event 主动拒绝,终结性拒绝 |
is_clean 属性判断"是否所有订阅者都接受了事件":failed == 0 and rejected == 0 and acked > 0。值得注意:事件广播时若没有任何订阅者在线,计数全零,此时不算 clean——没有任何东西被接受,生产者不应该提交。
框架只负责报告计数,per-broker 的最终决策完全落在 advance 实现里。例如:
- Service Bus 生产者:
rejected非零则 dead-letter;只有failed非零则 abandon(redeliver);全部接受则 complete; - Pub/Sub 生产者:reject 时
nack,其余情况ack; - Kafka 生产者:同 lane 内按序提交 offset(见下文 lane 讨论)。
advance 一旦抛异常,整个共享流组被终止,每个订阅者收到失败哨兵,broker 从从未提交的 offset 重新投递——"终止"是安全默认,因为细粒度按 lane 重试需要额外追踪哪些 offset 可以安全重提交(shared_stream.py)。
5. Lane:控制推进的并发与顺序
get_advance_lane(broker_payload) 为每个事件指定推进 lane(shared_stream.py):
- 同一 lane 内的事件严格按照 fan-out 顺序推进:前一批
advance返回后才会等待下一批; - 不同 lane 互不等待,但全局任意时刻最多只有一个
advance调用在 await; - 默认实现把所有事件放进同一个 lane,保持原始全局顺序;
- 该方法在 fan-out 前被同步调用一次,必须 O(1) 且不能阻塞;若抛异常,整个 poll 视为失败。
典型的 Kafka 用法是返回 (topic, partition)——累积式 offset 提交只需要分区内有序,慢分区就不会拖住其他分区的提交。
订阅者拒绝:reject_shared_stream_event
除了 yield 触发器事件,订阅者的 filter 还可以调用 reject_shared_stream_event() 来终结性地拒绝当前正在处理的原始事件(shared_stream.py):
- 拒绝会立即解决该事件(没有需要先持久化的派生事件),并计入
AdvanceOutcome.rejected; - 拒绝与被动
failed语义不同:生产者应对 reject 做 dead-letter /nack,而对 failure 做 redeliver; - 实现依赖一个任务本地(task-local)的 context variable 携带"打开的绑定窗口",因此 filter 必须与驱动绑定窗口的 asyncio 任务运行在同一个任务中——若通过
asyncio.to_thread或新建任务驱动 filter 迭代,窗口将不可见,所有 reject 都会变成 no-op(并记录一条 warning); - 在快速路径、独立
run()或两条原始事件之间调用它,仅记录 warning 并忽略。
超时与慢订阅者:ack 超时、队列溢出
ack 模式引入了一个后台任务 _run_ack_timeout_loop,按 max(0.01, ack_timeout / 10) 的节奏扫描未解决事件(shared_stream.py)。任何在 ack_timeout 秒内未完成处理的订阅者——无论是还停在事件上,还是派生事件未获得持久化确认——都会通过既有的 _PollFailure 路径以 AckTimeout 异常被强制失败。其他订阅者不受影响,它们解决后生产者照常推进。
同理,shared_stream_subscriber_queue_size 在 ack 模式下仍然表示"每个订阅者未处理的原始事件数上限":manager 不会等待未解决的记账才拉下一条事件,背压纯粹是队列级的——队列满的订阅者以 _SubscriberOverflow 强制失败,兄弟订阅者不受影响;但队列主要防护的是 filter 还没来得及运行时的突发投递。broker 推进由单一 pump 任务按 lane 顺序派发:当某 lane 头部事件仍有未解决订阅者时,同一 lane 的后续事件都会等待——已解决的事件会累积进下一批(ack 超时正是这个 per-lane 队头阻塞的兜底边界,shared_stream.py)。
triggerer 重启语义与优雅停机
解决状态只在内存中(shared_stream.py):
- triggerer 重启后,broker 会重新投递从未被推进的事件,因此订阅者必须具有幂等性;
- 组停止时若有事件仍在等待持久化确认(例如最后一个订阅者在刚产出事件后立即退订),待推进事件被放弃,broker 重新投递;
- 多个共享同一 key 的触发器一起重启时,第一个重新订阅的触发器会创建全新的组并立即开始轮询,稍后重订阅的触发器作为普通迟到订阅者加入,不计入更早事件的快照,因此可能错过两者之间的窗口期提交的事件。
为此提供了 shared_stream_cohort_grace_period:设为正数秒,在组创建后延迟启动轮询,给并发重订阅留出加入窗口(config.yml)。这是 best-effort 窗口:超过宽限期才重订阅的触发器仍可能错过轮询开始后提交的事件。
配置项一览
三个相关配置项均在 [triggerer] 节下(config.yml),均为 Airflow 3.3.0 新增:
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
shared_stream_subscriber_queue_size |
integer | 1024 |
共享流触发器(通过 BaseEventTrigger.shared_stream_key 加入共享轮询)的每订阅者缓冲大小;慢订阅者填满队列时只有它自己失败,兄弟订阅者不受影响 |
shared_stream_ack_timeout |
float | 300.0 |
ack 模式下每个事件的 ack 超时(秒);订阅者在窗口内未完成处理(越过事件且派生触发器事件确认持久化)即被强制失败 |
shared_stream_cohort_grace_period |
float | 0 |
组创建后延迟启动轮询的秒数;设为 2.0–5.0 可降低 triggerer 重启时同 key 触发器并发重订阅的丢事件概率 |
对应源码中的常量:DEFAULT_SUBSCRIBER_QUEUE_MAX = 1024、DEFAULT_ACK_TIMEOUT = 300.0(shared_stream.py)。其中 shared_stream_ack_timeout 可在创建 manager 时用 SharedStreamManager(ack_timeout=...) 覆盖,triggerer 中则由配置项驱动。
如何启用 ack 通道:实现一个生产者
在 base.py 中,BaseEventTrigger 的四个关键成员构成 ack 模式的完整拼图:
class MyBrokerEventTrigger(BaseEventTrigger):
# 1) 共享 key:相同 key 的触发器实例共用一条轮询
def shared_stream_key(self) -> Hashable | None:
return ("my_broker", self.queue_name)
# 2) 打开共享流并产出 (raw_event, broker_payload) —— 仅 ack 模式使用
@classmethod
async def open_shared_stream(cls, kwargs):
raise NotImplementedError # ack 模式下不再需要
@classmethod
def create_shared_stream_producer(cls, kwargs) -> SharedStreamProducer:
return MyProducer(kwargs)
# 3) 把原始事件转换为自己的 TriggerEvent —— 两种模式共用,无需改动
async def filter_shared_stream(self, shared_stream):
async for raw_event in shared_stream:
yield TriggerEvent(payload=transform(raw_event))
生产者类则继承 SharedStreamProducer(shared_stream.py):
class MyProducer(SharedStreamProducer):
async def open_stream(self):
# 在此打开 broker 连接(不要在工厂或 __init__ 中打开)
async for raw_event, handle in self._broker.pull():
yield raw_event, handle # handle: Kafka offset / SQS receipt handle / Pub/Sub ack ID
async def advance(self, batch: Sequence[AdvanceItem]) -> None:
for item in batch:
if item.outcome.is_clean:
self._broker.ack(item.broker_payload) # 全部接受且已持久化
elif item.outcome.rejected:
self._broker.dead_letter(item.broker_payload) # 终结性拒绝
else:
self._broker.redeliver(item.broker_payload) # 被动失败,交给 broker 重投
def get_advance_lane(self, broker_payload):
return ("default",) # 默认单 lane 保序;Kafka 可返回 (topic, partition)
对照 新功能片段 的表述:Kafka / SQS / Pub/Sub / Service Bus 正是这种"持有 broker 句柄,在全部订阅者接受且派生触发器事件持久化到元数据库之后才 commit / delete / ack"的上游消息中间件场景。上述代码里的 advance 即"在 broker 句柄上提交"的落点。
实战要点小结
- 幂等性是硬要求:triggerer 重启或组终止后,broker 会重投从未推进的事件,订阅者必须能够安全处理重复;
- filter 保持兼容:把触发器从快速路径切到 ack 模式不需要改
filter_shared_stream,流的形状不变; - 持久化门控是核心卖点:ack 不会早于"派生触发器事件写入元数据库",从根源上避免"broker 已确认但事件尚未持久化"的丢失窗口;代价是推进等待一次状态同步(约一两秒),远小于默认 300 秒的 ack 超时;
- 区分 reject 与 fail:reject 是终结性拒绝(dead-letter / nack),fail 是临时失败(redeliver),
AdvanceOutcome让生产者可以分别处理; - 调整三个配置项应对不同负载:突发流量调大
shared_stream_subscriber_queue_size,慢消费调大shared_stream_ack_timeout,triggerer 滚动重启场景给shared_stream_cohort_grace_period设 2–5 秒。
延伸阅读
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0632
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00