后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载本篇文章系统讲解 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负责记录每个流上谁订阅了。提供程序组成与生命周期PersistentStreamProviderPersistentStreamProvider.cs是所有持久化流提供程序的公共实现。它通过按提供程序名称注册的IQueueAdapterFactory创建以下组件IQueueAdapter定义入队enqueue与接收receive语义IStreamQueueMapper定义流到队列的映射IStreamQueueBalancer定义队列到 Silo 的所有权分配IQueueAdapterCache为每个代理提供缓存可选的失败处理器failure handler、过滤器filter与退避backoff提供程序。从源码可以看出提供程序实现了IStreamProvider, IInternalStreamProvider, IControllable, IStreamSubscriptionManagerRetriever, ILifecycleParticipantILifecycleObservable等多个接口并通过Participate把自己挂载到 Orleans 生命周期上PersistentStreamProvider.cs初始化阶段Init默认位于ServiceLifecycleStage.ApplicationServices从 DI 中按名称解析IQueueAdapterFactory并调用adapterFactory.CreateAdapter(token)创建适配器如果 pub/sub 类型包含显式订阅还会获取显式订阅管理器PersistentStreamProvider.cs启动阶段Start默认位于ServiceLifecycleStage.Active当适配器方向为ReadOnly或ReadWrite时初始化拉取管理单元InitializePullingAgents并依据StreamLifecycleOptions.StartupState决定是否立即StartAgentsPersistentStreamProvider.cs关闭阶段Close先停止拉取管理单元再提交关闭状态——注意即使生命周期取消令牌超时stopTask.Ignore()也保证管理单元的后台清理继续执行PersistentStreamProvider.cs。默认行为是拉取代理自动启动显式基于 Grain 的订阅与隐式订阅同时启用。对应选项类StreamPubSubOptions.PubSubType的默认值为ExplicitGrainBasedAndImplicitStreamLifecycleOptions.StartupState的默认值为AgentsStarted见 PersistentStreamProviderOptions.cs 与 PersistentStreamProviderOptions.cs。此外PersistentStreamProvider实现了IControllable可通过ExecuteCommand下发StartAgents、StopAgents、GetAgentsState、GetNumberRunningAgents等命令并把AdapterCommandStartRange10000 起、AdapterFactoryCommandStartRange等区段留给自定义适配器/工厂扩展PersistentStreamProvider.cs。拉取代理的关键选项拉取代理的行为由StreamPullingAgentOptions控制PersistentStreamProviderOptions.cs选项默认值说明InitialSubscriptionStartPositionLatest未指定序列令牌/显式起始位置的新订阅从何处开始EarliestAvailable表示从拉取代理本地队列缓存中尚保留的最早消息开始接收器检查点与队列位置不变具体序列令牌或显式起始位置优先级更高BatchContainerBatchSize1每个批容器batch container的批大小GetQueueMsgsTimerPeriod100 ms轮询队列消息的间隔InitQueueTimeout5 s队列初始化超时MaxEventDeliveryTime1 min单次事件投递的最长时间StreamInactivityPeriod30 min流不活动判定周期需要强调的是BatchContainerBatchSize 1与 100 ms 空轮询周期只是运行时默认值并不代表通用吞吐建议在高吞吐场景下通常需要按具体适配器与业务进行调优。队列映射与所有权队列映射器流到队列的确定性映射IStreamQueueMapper把流标识StreamId确定性地映射到一个队列。同一提供程序的所有生产者与消费者必须使用兼容的映射否则事件可能被写入没有代理读取的队列导致写入了却没人消费。队列均衡器与拉取管理单元IStreamQueueBalancer负责把队列分配给各 Silo并发布带序号的队列所有权变更通知。PersistentStreamPullingManagerPersistentStreamPullingManager.cs是一个 Silo 本地的 SystemTarget它在Initialize中调用queueBalancer.Initialize(mapper)并订阅队列分布变化事件SubscribeToQueueDistributionChangeEvents随后获取本 Silo 的队列列表序列化队列变更通知由于通知以 grain 方法调用的形式到达SystemTarget 本身不可重入管理单元再借助nonReentrancyGuarantorAsyncSerialExecutor串行执行并通过latestRingNotificationSequenceNumber忽略过期旧序号的通知PersistentStreamPullingManager.cs为每个归属队列启动/停止一个拉取代理AddNewQueues/RemoveQueues。当集群成员变化Silo 加入或故障时队列会在管理单元之间转移代理本身不是虚拟对象不会迁移——队列交接走的是停止旧代理 → 在新管理单元上初始化新代理的路径。StopAgents会把managerState置为AgentsStopped此时后续到达的分布变化通知会被直接跳过避免在代理未运行时做无意义的重均衡PersistentStreamPullingManager.cs。拉取代理循环每个PersistentStreamPullingAgentPersistentStreamPullingAgent.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 操作完成该订阅身份的清理即使游标对账与其持久化重叠。缓存与游标不变量IQueueCacheOrleans.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外部队列的读/写与可回退性rewindabilityIQueueAdapterReceiver接收、确认ack与关闭IStreamQueueMapper稳定的分区映射IStreamQueueBalancer集群内的队列所有权IQueueCache及其游标缓冲与安全清理safe purgeIStreamFailureHandler投递失败策略具体适配器的编写模式托管与校验可参考 provider-authoring而一个完整的实现实例是 Azure Queue 流 文档对应的适配器仓库中的 Event Hubs、Kinesis、SQS 等提供程序位于 src/Azure/Orleans.Streaming.EventHubs 与 src/AWS也是同一套契约在不同队列技术上的落地可作为对照阅读。小结Orleans 的持久化流拉取架构可以用一句话概括队列是持久边界Silo 上的拉取代理是搬运工本地缓存是解耦层pub-sub 是订阅簿游标与检查点共同决定投递进度。理解PersistentStreamProvider、PersistentStreamPullingManager、PersistentStreamPullingAgent、IQueueCache/SimpleQueueCache之间的协作关系是排查流延迟、背压、重复投递与检查点异常的前提也是编写自定义持久化流提供程序的基础。赞分享后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载相关推荐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 云原生应用框架持久化后端微服务上一篇如何利用Ray Adapter在华为鲲鹏和昇腾硬件上获得3倍性能提升终极迁移指南下一篇StratoVirt基于Rust的下一代轻量级虚拟化平台如何实现安全与性能的完美平衡创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
