- 后端
- 微服务
【免费下载链接】orleans
Cloud Native application framework for .NET
本篇文章系统讲解 Orleans 中"持久化流提供程序(Persistent Stream Provider)"的运行时实现机制:生产者如何通过适配器将事件写入持久化队列,各 Silo 上的拉取代理(Pulling Agent)如何认领队列分区、批量读取、缓存、发现订阅并将事件以普通 Orleans 调用投递给消费者。读完本文,你将掌握持久化流的组件构成与生命周期、队列映射与所有权分配、拉取代理的循环与关闭交接、缓存/游标的不变量、pub-sub 握手协议,以及"至多一次/至少一次"投递语义背后的决定性因素,并能基于源码定位各环节的精确实现位置。
本文属于 Orleans 运行时实现层面(implementation)的剖析;如需流的 API 用法与提供程序选型,请参考 streaming 文档。
总体架构:从生产者到消费者的一条拉取链路
持久化流提供程序(persistent stream provider)把 Orleans 流与一个持久化队列技术(例如 Azure Queue、Event Hubs、Kinesis、SQS 等)连接起来。其核心思路是"由 Silo 主动从队列拉取,而不是由队列推送到 Silo":
整条链路的分工如下:
- 生产者(Producer)通过提供程序获取流句柄后,经
IQueueAdapter入队; - 持久化队列(Queue)是事件的真正存储边界;
- 队列均衡器(IStreamQueueBalancer)决定每个队列由哪个 Silo 负责;
- 拉取管理单元(PersistentStreamPullingManager)是每个 Silo 上的本地 SystemTarget,负责启动/停止归属于本 Silo 的拉取代理;
- 拉取代理(Pulling Agent)本身也是 SystemTarget,从队列批量读数据、写入本地缓存、向 pub-sub 查询订阅者,再以普通 Orleans 调用把事件投递给消费者;
- 队列缓存(IQueueCache)解耦"队列读取"与"消费者投递";
- pub-sub负责记录每个流上"谁订阅了"。
提供程序组成与生命周期
PersistentStreamProvider(PersistentStreamProvider.cs)是所有持久化流提供程序的公共实现。它通过按提供程序名称注册的IQueueAdapterFactory创建以下组件:
IQueueAdapter:定义入队(enqueue)与接收(receive)语义;IStreamQueueMapper:定义流到队列的映射;IStreamQueueBalancer:定义队列到 Silo 的所有权分配;IQueueAdapterCache:为每个代理提供缓存;- 可选的失败处理器(failure handler)、过滤器(filter)与退避(backoff)提供程序。
从源码可以看出,提供程序实现了IStreamProvider, IInternalStreamProvider, IControllable, IStreamSubscriptionManagerRetriever, ILifecycleParticipant<ILifecycleObservable>等多个接口,并通过Participate把自己挂载到 Orleans 生命周期上(PersistentStreamProvider.cs):
- 初始化阶段(
Init,默认位于ServiceLifecycleStage.ApplicationServices):从 DI 中按名称解析IQueueAdapterFactory并调用adapterFactory.CreateAdapter(token)创建适配器;如果 pub/sub 类型包含显式订阅,还会获取显式订阅管理器(PersistentStreamProvider.cs); - 启动阶段(
Start,默认位于ServiceLifecycleStage.Active):当适配器方向为ReadOnly或ReadWrite时,初始化拉取管理单元(InitializePullingAgents),并依据StreamLifecycleOptions.StartupState决定是否立即StartAgents(PersistentStreamProvider.cs); - 关闭阶段(
Close):先停止拉取管理单元,再提交关闭状态——注意即使生命周期取消令牌超时,stopTask.Ignore()也保证管理单元的后台清理继续执行(PersistentStreamProvider.cs)。
默认行为是:拉取代理自动启动;显式基于 Grain 的订阅与隐式订阅同时启用。对应选项类StreamPubSubOptions.PubSubType的默认值为ExplicitGrainBasedAndImplicit,StreamLifecycleOptions.StartupState的默认值为AgentsStarted(见 PersistentStreamProviderOptions.cs 与 PersistentStreamProviderOptions.cs)。
此外,PersistentStreamProvider实现了IControllable,可通过ExecuteCommand下发StartAgents、StopAgents、GetAgentsState、GetNumberRunningAgents等命令,并把AdapterCommandStartRange(10000 起)、AdapterFactoryCommandStartRange等区段留给自定义适配器/工厂扩展(PersistentStreamProvider.cs)。
拉取代理的关键选项
拉取代理的行为由StreamPullingAgentOptions控制(PersistentStreamProviderOptions.cs):
| 选项 | 默认值 | 说明 |
|---|---|---|
InitialSubscriptionStartPosition | Latest | 未指定序列令牌/显式起始位置的新订阅从何处开始(EarliestAvailable表示从拉取代理本地队列缓存中尚保留的最早消息开始,接收器检查点与队列位置不变;具体序列令牌或显式起始位置优先级更高) |
BatchContainerBatchSize | 1 | 每个批容器(batch container)的批大小 |
GetQueueMsgsTimerPeriod | 100 ms | 轮询队列消息的间隔 |
InitQueueTimeout | 5 s | 队列初始化超时 |
MaxEventDeliveryTime | 1 min | 单次事件投递的最长时间 |
StreamInactivityPeriod | 30 min | 流不活动判定周期 |
需要强调的是,BatchContainerBatchSize = 1与 100 ms 空轮询周期只是运行时默认值,并不代表通用吞吐建议;在高吞吐场景下通常需要按具体适配器与业务进行调优。
队列映射与所有权
队列映射器:流到队列的确定性映射
IStreamQueueMapper把流标识(StreamId)确定性地映射到一个队列。同一提供程序的所有生产者与消费者必须使用兼容的映射,否则事件可能被写入没有代理读取的队列,导致"写入了却没人消费"。
队列均衡器与拉取管理单元
IStreamQueueBalancer负责把队列分配给各 Silo,并发布带序号的队列所有权变更通知。PersistentStreamPullingManager(PersistentStreamPullingManager.cs)是一个 Silo 本地的 SystemTarget,它:
- 在
Initialize中调用queueBalancer.Initialize(mapper)并订阅队列分布变化事件(SubscribeToQueueDistributionChangeEvents),随后获取本 Silo 的队列列表; - 序列化队列变更通知:由于通知以 grain 方法调用的形式到达,SystemTarget 本身不可重入,管理单元再借助
nonReentrancyGuarantor(AsyncSerialExecutor)串行执行,并通过latestRingNotificationSequenceNumber忽略过期(旧序号)的通知(PersistentStreamPullingManager.cs); - 为每个归属队列启动/停止一个拉取代理(
AddNewQueues/RemoveQueues)。
当集群成员变化(Silo 加入或故障)时,队列会在管理单元之间转移;代理本身不是虚拟对象,不会迁移——队列交接走的是"停止旧代理 → 在新管理单元上初始化新代理"的路径。StopAgents会把managerState置为AgentsStopped,此时后续到达的分布变化通知会被直接跳过,避免在代理未运行时做无意义的重均衡(PersistentStreamPullingManager.cs)。
拉取代理循环
每个PersistentStreamPullingAgent(PersistentStreamPullingAgent.cs)都是 SystemTarget,因此它的所有代码都在 Orleans 调度器的单线程调度下执行(RunOrQueueTask保证串行)。其主循环为:
- 向适配器接收器(
IQueueAdapterReceiver)请求一个批次; - 把批次容器加入队列缓存;
- 按流对缓存条目分组;
- 解析并缓存 pub-sub 注册信息;
- 独立推进每个订阅的游标;
- 通过 Orleans 消息投递事件;
- 记录投递进度与失败;
- 只清除缓存判定为"可以安全移除"的数据。
代理在Initialize时创建队列缓存(queueAdapterCache.CreateQueueCache(QueueId))与接收器(queueAdapter.CreateReceiver(QueueId)),并注册一个周期定时器驱动RunQueuePump;定时器首轮带有随机偏移(RandomTimeSpan.Next(GetQueueMsgsTimerPeriod)),避免所有代理在同一时刻同时轮询(PersistentStreamPullingAgent.cs)。接收器初始化失败不会阻止代理开始泵取——IQueueAdapterReceiver被要求自行负责初始化重试。
关闭与队列交接
当代理停止(Shutdown,见 PersistentStreamPullingAgent.cs)时,执行以下有序清理:
- 关闭准入(AdmissionGate):
_workAdmission.CloseAsync()拒绝新的后台工作,同时停止轮询定时器; - 等待在途工作完成:等待接收器初始化任务、活跃的队列泵取任务,以及已接受的 producer 注册、订阅握手与投递完成;
- 已接受的工作完成令牌记账,释放注册"pin"与批次保护,期间缓存与接收器仍然可用;未完成的调用保留其原有消息超时与重试限制;
- 上报最终投递进度给缓存(
NotifyDeliveryProgress),释放订阅游标,关闭接收器——这样提供程序特定的检查点刷新能观察到完整的进度; - 关闭时仍在进行的注册(
hasPendingRegistrations)会保留现有检查点,因为其订阅位置尚不确定;producer 注销在接收器清理之后进行; - 当管理单元把同一代理复用于一条新分配的队列时,
Initialize会先等待上述完整清理(await shutdownTask),再重新开放准入(_workAdmission = new())。
关于订阅握手的时序细节(源码注释与实现明确体现):
- 显式订阅通知会立即返回确认(
AddSubscriber不阻塞调用方),代理通过完成事件跟踪其异步握手;这允许订阅方消费者先结束当前调用、再响应握手(PersistentStreamPullingAgent.cs); - 握手完成后建立订阅的当前游标与回放位置;最新请求的握手拥有对账权,被取代的旧请求的响应不会覆盖该所有权;旧握手代(handshake generation)的投递完成与错误处理释放各自的工作,同时保留替代位置,使最终检查点进度反映被接受的 rewind;
- 失败的重新握手(re-handshake)会使订阅位置不确定——即使该订阅此前已注册;代理在空闲清理期间仍保留该流条目并保留现有检查点,直到一次成功握手完成位置对账(参见
HasUnresolvedHandshake、HandshakeRequestId的判定逻辑); - 订阅移除会撤销在途握手与投递的所有权;在有效所有权下发出的终态 pub-sub 操作完成该订阅身份的清理,即使游标对账与其持久化重叠。
缓存与游标不变量
IQueueCache(Orleans.Streaming 公共缓存接口)把"从队列读取"与"向消费者投递"解耦。每个订阅拥有一个独立的IQueueCacheCursor,因此慢消费者不会直接阻塞位于更后游标处的快消费者。
关键不变量:
- 缓存跟踪所有活跃订阅中最早的投递进度;清理(purge)绝不能移除仍被任一游标需要的条目;
- 默认实现
SimpleQueueCache使用压力桶(pressure buckets):当滞后增长时停止或放慢读取,而不是丢弃未投递事件;其默认容量为4,096 个批容器(SimpleQueueCacheOptions.CacheSize,默认4096,见 SimpleQueueCacheOptions.cs,且校验器要求CacheSize > 0); - 缓存容量不是持久性:队列始终是持久边界,具体取决于适配器的确认(acknowledgement)契约;接收器关闭后,提供程序特定的检查点(checkpoint)才把最终进度落盘。
Pub-sub 握手
代理会为每个流注册为 producer,并从流 pub-sub 获取订阅记录。关键点:
- 代理持有pin 游标,在订阅握手完成期间防止缓存清理越过请求的起始令牌;
- 新的订阅通知会更新代理的本地 pub-sub 缓存(
pubSubCache); - 序列令牌(sequence token)允许可回退(rewindable)的适配器从受支持的历史位置开始;而
IQueueAdapter.IsRewindable为false的适配器必须拒绝不支持的令牌,而不是假装支持(源码中GetCacheCursorAtStartPosition对不支持的位置会抛NotSupportedException,参见 PersistentStreamPullingAgent.cs)。
投递与失败语义
代理通常先等待投递完成再推进订阅游标,从而形成"按订阅"的背压。当投递失败时,代理调用配置的IStreamFailureHandler;根据提供程序策略,显式订阅可能被 fault 并移除。IStreamFailureHandler接口(IStreamFailureHandler.cs)暴露:
ShouldFaultSubsriptionOnError:出错时是否应 fault 订阅;OnDeliveryFailure(...):事件投递给消费者的一切手段用尽后被调用;OnSubscriptionFailure(...):建立订阅失败时被调用。
持久化流并不是普遍"恰好一次"(exactly-once)。最终语义取决于:
- 外部队列在何时认为消息已被确认(acknowledged);
- 适配器在接收器或 Silo 故障后能否重投(redeliver);
- 缓存检查点行为;
- 消费者的幂等性;
- 提供程序特定的序列令牌。
因此,一条队列消息在所有权变更或故障之后可能被再次投递;执行持久性副作用(如写数据库)的消费者应当具备幂等性。相关测试(如 PersistentStreamPullingAgentTests.cs 与 PullingAgentManagementTests.cs)覆盖了代理生命周期、队列管理与恢复等场景,可作为行为参照。
扩展契约:各接口的职责边界
提供程序作者应保持以下职责分离(全部定义在 Orleans.Streaming 命名空间):
| 接口 | 职责 |
|---|---|
IQueueAdapter | 外部队列的读/写与可回退性(rewindability) |
IQueueAdapterReceiver | 接收、确认(ack)与关闭 |
IStreamQueueMapper | 稳定的分区映射 |
IStreamQueueBalancer | 集群内的队列所有权 |
IQueueCache及其游标 | 缓冲与安全清理(safe purge) |
IStreamFailureHandler | 投递失败策略 |
具体适配器的编写模式(托管与校验)可参考 provider-authoring,而一个完整的实现实例是 Azure Queue 流 文档对应的适配器;仓库中的 Event Hubs、Kinesis、SQS 等提供程序(位于 src/Azure/Orleans.Streaming.EventHubs 与 src/AWS)也是同一套契约在不同队列技术上的落地,可作为对照阅读。
小结
Orleans 的持久化流拉取架构可以用一句话概括:队列是持久边界,Silo 上的拉取代理是搬运工,本地缓存是解耦层,pub-sub 是订阅簿,游标与检查点共同决定投递进度。理解PersistentStreamProvider、PersistentStreamPullingManager、PersistentStreamPullingAgent、IQueueCache/SimpleQueueCache之间的协作关系,是排查流延迟、背压、重复投递与检查点异常的前提,也是编写自定义持久化流提供程序的基础。
- 后端
- 微服务
【免费下载链接】orleans
Cloud Native application framework for .NET
相关推荐
Orleans 流投递语义详解:投递保证、顺序、重放与恢复
Orleans 流投递语义详解:投递保证、顺序、重放与恢复 本文是 Orleans 流(Streams)编程模型的实战指南,聚焦"投递语义"这一核心主题:生产者
后端微服务SeaTunnel Edge Agent 架构深度解析:WAL 出站队列、EdgeSocket 协议与投递语义边界
SeaTunnel Edge Agent 架构深度解析:WAL 出站队列、EdgeSocket 协议与投递语义边界 SeaTunnel Edge Agent 是
数据集成ETL大数据批处理流处理变更数据捕获Orleans 持久化流自定义队列适配器(Custom Queue Adapter)完整开发指南
Orleans 持久化流自定义队列适配器(Custom Queue Adapter)完整开发指南 导读 本文是 Orleans(.NET 云原生应用框架)持久化
后端微服务
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考