DYNAMIC EXPLAINER / ARTICLE FLOW

把章节变成可单步观察的过程

01 / 14
STEP 01 / CODE / LOGIC

先说结论

事件队列把“事件发生”和“监听者何时处理”分开,能稳定顺序并避免深层同步重入;代价是结果不再立刻可见,还必须定义延迟、容量和事件风暴策略。

队列初始:[Damage(30)]
处理 Damage:HP 100 -> 70,并入队 HealthChanged(70)
当前队列:[HealthChanged(70)]
处理 HealthChanged:UI 更新为 70
队列结束:[]

分类:工程设计

先说结论

事件队列把“事件发生”和“监听者何时处理”分开,能稳定顺序并避免深层同步重入;代价是结果不再立刻可见,还必须定义延迟、容量和事件风暴策略。

A 事件产生 B 事件

队列初始:[Damage(30)]
处理 Damage:HP 100 -> 70,并入队 HealthChanged(70)
当前队列:[HealthChanged(70)]
处理 HealthChanged:UI 更新为 70
队列结束:[]

若采用同步递归发布,UI 监听会在 Damage 处理尚未完全结束时进入。队列模式则让当前事件先完成,再处理新事件。若每个 HealthChanged 又发布另一个 HealthChanged,队列会持续增长,这就是需要监控的事件风暴。

先用生活模型理解事件

想象一块公告板:

“三号房间已开放。”

这句话描述的是一个已经发生的事实。

看到公告的人可以有很多:

  • 引导员更新路线;
  • 记录员写入日志;
  • 显示屏刷新状态;
  • 也可能暂时没有人关注。

公告不会要求某个监听者替它完成“开放房间”这件事,因为开放已经发生了。

事件系统解决的是:

事实已经发生,怎样让感兴趣的部分知道?

事件与命令的职责差异

事件和命令不能只靠“都是消息”混为一谈。

对比项 事件 命令
语义 已发生的事实 希望执行的请求
命名示例 DoorOpened OpenDoor
责任方 可以有多个监听者,也可以无人监听 通常有一个明确执行者
是否可拒绝 不能把事实拒绝成未发生 可以因校验失败而拒绝执行
结果依赖 发布者通常不靠监听返回值决定事实 调用方常需要明确执行结果

推荐边界:

提交 OpenDoorCommand
    ↓
执行成功,门的状态已改变
    ↓
发布 DoorOpenedEvent

如果“发布事件后等某个监听者修改核心状态”,事件就承担了命令职责,处理结果和责任方都会变得不明确。

为什么需要事件队列

最直接的事件分发是同步调用:

发布事件
    ↓
立刻逐个调用监听者
    ↓
全部返回后,发布函数才返回

同步分发简单,而且调用栈能直接展示因果关系。

事件队列则把发布与处理分开:

当前阶段发布事件 -> 入队
稍后的明确处理阶段 -> 出队并分发

它适合:

  • 避免在状态修改过程中重入监听逻辑;
  • 统一在某个逻辑阶段处理通知;
  • 控制每轮处理量;
  • 记录事件顺序并辅助调试;
  • 为明确的延迟语义提供容器。

代价是事实发生时间与监听时间分离,调用栈不再直接连接两者。

事件队列的逐步状态

一个教学模型:

Created:事件对象已创建
    ↓
Queued:加入待处理队列
    ↓
Dispatching:从队列取出,建立监听者快照
    ↓
Handled:所有监听者完成

失败路径:

Dispatching
    ↓
某监听者抛出异常
    ↓
本次分发失败,异常保持可见

“异常后是否继续调用剩余监听者”会改变行为。本文的最小实现选择让异常直接向上传播,不吞错、不返回默认成功,也不擅自继续。

最小 C# 示例

下面示例使用类型作为事件分类键,并在分发前复制监听列表,避免监听者在回调中订阅或退订时破坏当前枚举。

public interface IEvent
{
}

public sealed record TextAdded(string Text) : IEvent;

public sealed class EventQueue
{
    private readonly Queue<IEvent> _pending = new();
    private readonly Dictionary<Type, List<Delegate>> _listeners = new();

    public void Subscribe<TEvent>(Action<TEvent> listener)
        where TEvent : IEvent
    {
        ArgumentNullException.ThrowIfNull(listener);

        Type eventType = typeof(TEvent);
        if (!_listeners.TryGetValue(eventType, out List<Delegate>? list))
        {
            list = new List<Delegate>();
            _listeners.Add(eventType, list);
        }

        list.Add(listener);
    }

    public bool Unsubscribe<TEvent>(Action<TEvent> listener)
        where TEvent : IEvent
    {
        ArgumentNullException.ThrowIfNull(listener);

        return _listeners.TryGetValue(typeof(TEvent), out List<Delegate>? list)
            && list.Remove(listener);
    }

    public void Publish(IEvent eventData)
    {
        ArgumentNullException.ThrowIfNull(eventData);
        _pending.Enqueue(eventData);
    }

    public void DispatchNext()
    {
        if (_pending.Count == 0)
        {
            throw new InvalidOperationException("没有待分发事件。");
        }

        IEvent eventData = _pending.Dequeue();
        Type eventType = eventData.GetType();

        if (!_listeners.TryGetValue(eventType, out List<Delegate>? listeners))
        {
            return; // “允许无人监听”是此教学事件模型的明确契约。
        }

        Delegate[] snapshot = listeners.ToArray();
        foreach (Delegate listener in snapshot)
        {
            // 不捕获异常:监听失败会沿调用栈向上传播。
            listener.DynamicInvoke(eventData);
        }
    }
}

var events = new EventQueue();

Action<TextAdded> printLength = message =>
    Console.WriteLine(message.Text.Length);

events.Subscribe(printLength);
events.Publish(new TextAdded("示例"));
events.DispatchNext();

教学示例中的 "示例" 只是演示输入。示例代码明确表达:

  • 发布只负责入队;
  • DispatchNext 才执行监听者;
  • 无人监听是合法状态;
  • 监听异常不会被替换为成功结果。

DynamicInvoke 便于展示统一容器,但存在反射式调用开销,并会把目标异常包装为调用异常。实际设计也可以为每种事件保存类型安全的分发表;选择取决于规模与调试需求。

监听者生命周期

事件发布者通常持有监听委托。若监听委托引用了对象,发布者可能间接延长该对象的生命周期。

正常路径

对象创建
    ↓
订阅事件
    ↓
接收事件
    ↓
对象离开有效期前退订

失败路径

对象逻辑上已结束
    ↓
忘记退订
    ↓
发布者仍持有委托
    ↓
对象继续收到事件,或无法按预期释放

订阅管理常见方式:

  • 明确的 Subscribe / Unsubscribe 成对调用;
  • 返回一个可释放的订阅令牌;
  • 由具有相同生命周期的容器统一管理。

是否使用弱引用属于架构决策。弱引用会改变监听者的存活语义,不能作为忘记管理生命周期的通用补丁。

分发期间修改监听列表

监听者可能在处理事件时退订自己或订阅新监听者。直接枚举原列表会导致集合修改异常或不明确的本轮可见性。

本文示例建立快照,因此契约是:

本次分发使用开始分发时的监听者集合。
分发中发生的订阅变化,从后续事件开始生效。

复制快照需要 O(L) 时间和空间,L 为该类型的监听者数量。

延迟事件的时间语义

“延迟”至少可能表示三种不同含义:

延迟到下一逻辑轮;
延迟指定逻辑时长;
延迟到某个阶段结束。

它们不能只用一个模糊的 Delay 字段混在一起。

基于逻辑轮的教学模型

教学示例:

当前逻辑轮 = 10
事件计划在逻辑轮 12 分发

数值 1012 仅用于解释。

逐步状态:

轮 10:创建事件,目标轮 12,进入延迟容器
轮 11:目标未到,不分发
轮 12:移动到就绪队列
轮 12 的事件阶段:按规则分发

需要明确:

  • 目标轮是“开始时”还是“结束时”触发;
  • 同一目标轮的事件按什么顺序;
  • 调度一个已经过期的目标轮是报错还是立即处理;
  • 暂停、变速或回放是否影响逻辑轮;
  • 延迟事件是否允许取消。

这些都属于行为契约,不能靠容器默认值猜测。

延迟容器

若需要不断取出最早到期事件,可以使用按目标时间排序的最小堆:

加入延迟事件:O(log n)
取出最早事件:O(log n)
查看最早到期时间:O(1)

如果延迟范围是有限且离散的逻辑轮,也可以研究时间轮等结构。数据结构应由时间语义和规模决定。

事件顺序与重入

假设处理事件 A 时又发布事件 B

可能有两种语义:

立即重入:当前监听栈内立刻处理 B。
排队:先完成 A 的所有监听者,再在队列顺序中处理 B。

事件队列通常用于选择第二种,但仍要明确:

B 是在本轮继续处理,还是下一轮处理?
一次排空队列是否允许处理新加入的事件?

如果采用“直到队列为空”,监听者不断发布新事件时,当前轮可能永远无法结束。

一种可观察的处理边界是:在阶段开始时记录本批数量,只处理这批事件,新事件留到下一阶段。另一种是设置明确预算并在超限时报告错误。选择哪一种会影响响应时序,需要由系统契约决定。

怎样调试事件系统

队列切断了直接调用栈,因此需要补充可观察信息。

建议记录的事实包括:

事件类型;
事件实例或追踪标识;
发布时间;
计划分发时间;
实际分发时间;
发布来源;
监听者数量;
处理顺序;
失败监听者与原始异常。

因果链

若事件 A 的监听者发布了事件 B,可以让 B 记录 A 的追踪标识作为父标识:

A
├─ B
│  └─ D
└─ C

这样能区分“谁先发生”和“谁导致谁发生”。

环形历史缓冲

调试记录可能不断增长。环形缓冲可以只保留最近 N 条记录,写入通常为 O(1)

容量 N 是行为相关的诊断配置,必须依据可接受内存和需要观察的时间窗口确定。本文不提供生产默认值。

异常可见性

不应这样处理:

try
{
    listener();
}
catch
{
    // 什么也不做
}

这会同时丢失:

  • 哪个监听者失败;
  • 原始异常和调用栈;
  • 后续监听者是否仍被执行;
  • 事件是否被系统视为处理完成。

若确实需要隔离监听者异常,应先明确:

是否继续后续监听者;
如何聚合并上报异常;
本事件最终状态是什么;
失败是否触发重试,以及重试是否会重复副作用。

在规则未确定前,保留原始失败比返回默认成功更可靠。

进阶:事件风暴怎样形成

事件风暴指短时间内事件数量或派生链急剧增长,使队列、CPU 或日志压力失控。

常见形成方式

循环:
A 的监听者发布 B,B 的监听者又发布 A。

扇出:
一个事件触发很多监听者,每个监听者又发布多个事件。

重复:
同一状态在一轮内被多次写入,每次都发布等价事件。

追赶:
处理速度长期低于发布速度,队列持续增长。

用数量理解扇出

教学示例:每个事件派生 3 个新事件,持续 4 层。

第 0 层:1
第 1 层:3
第 2 层:9
第 3 层:27
第 4 层:81

数值仅用于展示指数增长,不代表任何推荐阈值。

发现事件风暴

可以观察:

  • 每轮发布数和处理数;
  • 队列长度变化;
  • 单一事件类型的频率;
  • 单个根事件形成的因果链深度和宽度;
  • 同一发布源的重复事件;
  • 最老事件的等待时间。

控制手段的边界

可能的控制手段包括合并重复事件、按键去重、限制处理预算、背压或拒绝新事件。

但每一种都会改变语义:

合并会丢失中间事实;
去重要求先定义“什么算同一个事件”;
限流会增加延迟;
丢弃会造成事实不可见;
重试可能重复副作用。

因此不能把它们作为默认容错。应先找出风暴来源,再由明确规则决定是否允许合并、延迟或拒绝。

边界与失败路径

无监听者

对于事实广播,无监听者通常可以是合法状态。本文示例明确采用这个契约。若某消息必须有人执行,它更可能是命令,不应依赖事件碰巧被监听。

监听者重复订阅

列表允许重复时,同一委托会被调用多次。是否禁止重复订阅必须明确;不能既允许添加又在分发时悄悄去重。

事件数据可变

事件入队后若仍可被外部修改,监听者看到的可能不是发布时事实。不可变记录可以降低这种风险。

多线程访问

普通 Queue<T>List<T> 不提供任意多线程并发安全。跨线程发布、订阅与分发需要明确所有权、同步和顺序语义。

队列增长

生产速度超过消费速度时,内存会持续增长。最大容量、背压、拒绝或持久化都属于需要明确批准的策略;不能默认丢弃事件并声称成功。

监听者异常

本文示例直接抛出异常。异常发生前已执行的监听者副作用不会自动回滚,后续监听者也不会自动执行。

复杂度与代价

设队列长度为 Q,某事件类型监听者数量为 L

操作 常见时间复杂度 主要代价
事件入队 O(1) 保存事件对象
普通出队 O(1) 不含监听执行成本
建立监听快照 O(L) 复制委托引用
分发事件 O(L) 加监听逻辑成本 间接调用与副作用
订阅 列表尾部通常 O(1) 持有委托引用
退订 O(L) 查找目标委托
最小堆延迟调度 入队/出队 O(log Q) 维护时间顺序

额外工程成本包括:

  • 发布与处理不在同一调用栈,调试更难;
  • 时序、重入和监听快照规则需要文档化;
  • 事件对象与追踪记录增加内存;
  • 异常隔离、重试和背压都需要明确语义。

小结

事件传播已经发生的事实,命令请求执行动作。
队列把发布时间与处理时间分离。
延迟必须使用明确的时间基准和到期边界。
监听生命周期、分发快照和重入规则必须可解释。
异常要保持可见,不能用默认成功掩盖。
事件风暴应先追踪因果,再决定是否允许限流或合并。

一个可调试的事件系统,不只是“能把消息发出去”,还要让人回答:

谁在什么时间发布?
为什么发布?
排队多久?
按什么顺序交给了谁?
在哪个监听者失败?
这个事件又产生了哪些后续事件?

这些答案比事件总线的接口形式更重要。