配置 Orleans PubSub 存储让流订阅元数据在集群重启后依然存活【免费下载链接】orleansCloud Native application framework for .NET项目地址: https://gitcode.com/gh_mirrors/or/orleansOrleans 流Stream通过 pub/sub 汇合点rendezvous连接生产者和消费者而名为PubSubStore的 grain 存储提供程序负责持久化显式订阅元数据。本文基于官方文档 pubsub-storage.md 展开讲解PubSubStore的持久化取舍、Azure Table Storage 生产级配置以及订阅生命周期管理并结合仓库源码说明其内部实现原理帮助你在开发与生产环境中做出正确的配置决策。理解 PubSubStore 在 Orleans 流架构中的角色Orleans 的流系统是一个虚拟流virtual stream实现生产者向逻辑流 ID 发送事件消费者订阅该流 ID两者互不知道对方的存在。连接它们的正是发布/订阅汇合点pub/sub rendezvous——它维护着哪个流被谁订阅了的映射关系。在默认配置下这个汇合点由 grain 承载因此需要把订阅元数据持久化到某个 grain 存储提供程序中这个提供程序被命名为PubSubStore。该名称在源码中被定义为一个常量任何需要默认 pub/sub 存储的组件都会引用它// src/Orleans.Core.Abstractions/Providers/ProviderConstants.cs public const string DEFAULT_PUBSUB_PROVIDER_NAME PubSubStore;从源码结构看PubSubStore是 Orleans 流系统中的约定俗成的存储槽位流订阅管理器StreamSubscriptionManagerAdmin.cs在构造时请求ExplicitGrainBasedAndImplicit类型的 pub/sub 运行时该运行时内部的订阅记录 grain 使用PubSubStore作为其持久化状态流检查点 grainStreamCheckpointGrain.cs通过[PersistentState(StateName, ProviderConstants.DEFAULT_PUBSUB_PROVIDER_NAME)]直接绑定到PubSubStore用于持久化持久流persistent stream队列的消费位置grain 检查点器GrainStreamQueueCheckpointer.cs默认使用PubSubStore作为检查点存储。也就是说只要你使用持久流提供程序如 Azure Event Hubs、Azure Queue、Kinesis、SQS并启用 grain 检查点UseGrainCheckpointerPubSubStore就承担着订阅元数据 队列消费位置双重持久化职责。即使你只使用AddMemoryStreams流提供程序在默认情况下也期望存在一个名为PubSubStore的存储提供程序。三种持久化形态的取舍形态配置方式持久性适用场景内存存储AddMemoryGrainStorage(PubSubStore)集群内存状态丢失即订阅记录丢失开发、单元测试、演示持久存储AddAzureTableGrainStorage(PubSubStore, ...)等跨 silo 重启、集群重启存活生产环境隐式订阅[ImplicitStreamSubscription]特性由 grain 元数据派生不产生订阅记录订阅与流 ID 存在确定映射关系时其中隐式订阅值得特别说明它不经过显式订阅记录而是从 grain 的元数据[ImplicitStreamSubscription(namespace)]中推导出该 grain 订阅了哪个流因此不写入PubSubStore也不受存储持久性影响。仓库中的示例 ImplicitSubscriptions.cs 展示了典型写法[ImplicitStreamSubscription(TemperatureStreams.Namespace)] public sealed class DeviceTelemetryGrain : Grain, IDeviceTelemetryGrain, IAsyncObserverTemperatureReading, IStreamSubscriptionObserver { public Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) { var handle handleFactory.CreateTemperatureReading(); return handle.ResumeAsync(this); } // ... }开发环境用内存存储快速起步在开发与测试阶段使用内存存储是官方推荐的方式。仓库中的流配置示例 Configuration.cs 给出了完整的 silo 配置// memory_silo 片段 builder.UseOrleans(siloBuilder { siloBuilder .AddMemoryStreams(TemperatureStreams.ProviderName) .AddMemoryGrainStorage(PubSubStore); });对应的客户端侧同样只需要添加流提供程序客户端本身不直接使用PubSubStore// memory_client 片段 builder.UseOrleansClient(clientBuilder { clientBuilder.AddMemoryStreams(TemperatureStreams.ProviderName); });需要注意内存存储的订阅记录在集群状态丢失时会一并消失。如果 silo 全部重启且没有其他持久化副本之前创建的显式订阅会丢失消费者需要重新执行订阅逻辑。生产环境以 Azure Table Storage 持久化 PubSubStore对于生产环境官方文档推荐使用 Azure Table Storage 作为PubSubStore的持久化后端并优先使用托管标识managed identity而非连接字符串。方式一托管标识推荐// pubsub_managed_identity 片段 var endpoint new Uri(configuration[AZURE_TABLE_STORAGE_ENDPOINT]!); var credential new DefaultAzureCredential(); hostBuilder.UseOrleans(siloBuilder { siloBuilder.AddAzureTableGrainStorage( PubSubStore, options options.TableServiceClient new TableServiceClient(endpoint, credential)); });托管标识方式通过DefaultAzureCredential依次尝试环境凭据、托管标识、Azure CLI 等多种认证链避免了在配置文件中硬编码密钥适合部署在 Azure 容器应用、AKS、VM 等支持托管标识的环境中。AZURE_TABLE_STORAGE_ENDPOINT指向 Table 服务的终结点形如https://account.table.core.windows.net/。方式二连接字符串// pubsub_connection_string 片段 hostBuilder.UseOrleans(siloBuilder { siloBuilder.AddAzureTableGrainStorage( PubSubStore, options options.TableServiceClient new TableServiceClient(connectionString)); });连接字符串方式适合本地开发、测试以及无法使用托管标识的受限环境。同样的配置模式也适用于其他持久化后端例如AddDynamoDBGrainStorageOrleans.Clustering.DynamoDB、AddAdoNetGrainStorageOrleans.Persistence.AdoNet或AddCosmosGrainStorageOrleans.Persistence.Cosmos只需将存储提供程序名称指定为PubSubStore即可。集群身份与存储的对应关系官方文档强调了一条关键原则使用稳定的 Orleans service ID并在集群重启之间保持相同的持久化配置。服务 IDservice ID是 Orleans 集群的逻辑标识。pub/sub 订阅记录的存储键派生自服务 ID因此修改 service ID → 从 pub/sub 系统的角度看订阅注册表变成逻辑上全新的旧订阅记录不再被新集群识别删除或重建底层表 → 订阅记录全部丢失相当于新建注册表存储配置不一致 → 不同 silo 可能读写不同的表导致订阅状态不一致。生产环境中应把 service ID 视为需要刻意维护、保持不变的部署标识。订阅生命周期激活、恢复与移除持久化PubSubStore只保证订阅记录谁订阅了哪个流得以保存但它不保存消费者的 observer 实例。文档明确指出即使使用持久化的PubSubStore显式消费者在激活后也必须调用StreamSubscriptionHandleT.ResumeAsync()将当前 observer 实例挂接到订阅句柄上。同样持久化的事件存储durable event storage也不会让订阅记录自动变得持久——事件存储与订阅元数据是两个独立层次需要根据恢复需求分别配置。仓库中的显式订阅示例 ExplicitSubscriptions.cs 完整展示了这一生命周期管理是可直接套用的实战模板public override async Task OnActivateAsync(CancellationToken cancellationToken) { _stream TemperatureStreams.Get(this, this.GetPrimaryKeyString()); var handles await _stream.GetAllSubscriptionHandles(); foreach (var handle in handles) { await handle.ResumeAsync(this); } } public async Task SubscribeAsync() { var handles await _stream.GetAllSubscriptionHandles(); if (handles.Count 0) { await _stream.SubscribeAsync(this); } } public async Task UnsubscribeAsync() { var handles await _stream.GetAllSubscriptionHandles(); foreach (var handle in handles) { await handle.UnsubscribeAsync(); } }这段代码体现了显式订阅的三个核心操作SubscribeAsync首次订阅时创建订阅记录并持久化到PubSubStore。代码先检查GetAllSubscriptionHandles()是否已有句柄避免重复创建订阅ResumeAsyncgrain 激活OnActivateAsync时从存储中取回所有订阅句柄并重新挂接当前实例实现grain 重启后恢复订阅UnsubscribeAsync不再需要时移除订阅。文档建议在订阅不再被需要时主动调用它防止PubSubStore中堆积无用的订阅记录。运维指南备份、命名与替换官方文档给出四条直接可执行的运维建议像备份其他应用元数据一样备份并监控PubSubStore。订阅记录属于业务元数据丢失后显式订阅需要逐个重建保持提供程序名称稳定。从 pub/sub 系统的视角看提供程序名称是流身份stream identity的一部分重命名提供程序等于改变了流身份会导致既有订阅失效及时清理订阅。通过StreamSubscriptionHandleT.UnsubscribeAsync()移除不再需要的订阅控制存储增长与系统开销替换PubSubStore前先规划好显式订阅的重建方案。无论是更换存储后端还是迁移到新集群都要预先设计如何重新创建既有显式订阅例如在 grain 激活逻辑中通过SubscribeAsync幂等重建。底层原理从源码看 PubSub 汇合点的运作为了更稳妥地配置PubSubStore有必要理解它在 Orleans 流实现中的位置。流提供程序的 pub/sub 类型由StreamPubSubOptions控制其默认值是ExplicitGrainBasedAndImplicit显式基于 grain 隐式// src/Orleans.Streaming/PersistentStreams/Options/PersistentStreamProviderOptions.cs public class StreamPubSubOptions { public StreamPubSubType PubSubType { get; set; } DEFAULT_STREAM_PUBSUB_TYPE; public const StreamPubSubType DEFAULT_STREAM_PUBSUB_TYPE StreamPubSubType.ExplicitGrainBasedAndImplicit; }该配置通过ConfigureStreamPubSub扩展方法应用到持久流提供程序上ClusterClientPersistentStreamConfigurator.cs。基于 grain 的显式订阅意味着每个流 ID 对应一个订阅管理器 grain其状态持久化在PubSubStore中——这正是本文配置项存在的原因。此外使用持久流提供程序如 Event Hubs、Kinesis、Azure Queue时队列消费位置checkpoint也是关键状态。若启用UseGrainCheckpointer检查点会默认存入PubSubStore见 PersistentStreamConfiguratorExtension.cs 的UseGrainCheckpointer及其对GrainStreamQueueCheckpointerOptions.StorageProviderName默认值为PubSubStore的说明见 GrainStreamQueueCheckpointerOptions.cs。因此在规划持久化时要同时覆盖订阅元数据与检查点数据两个层面。总结一套配置决策清单决策点建议开发/测试环境AddMemoryGrainStorage(PubSubStore)接受重启即丢失生产环境持久化后端 稳定 service ID 保持配置一致优先托管标识订阅恢复grain 激活时调用ResumeAsync重新挂接 observer订阅清理不再需要时调用UnsubscribeAsync存储替换提前规划显式订阅的重建视同元数据迁移PubSubStore虽小却是 Orleans 流系统可靠性的基石之一。理解它的持久化边界、正确配置存储后端并配合完整的订阅生命周期管理才能在集群重启、滚动升级等场景下保持流的投递连续性。若要深入理解流系统内部的汇合点与 pulling agent 设计可继续阅读 Orleans streams implementation。【免费下载链接】orleansCloud Native application framework for .NET项目地址: https://gitcode.com/gh_mirrors/or/orleans创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考 SEO 优化官网定制响应式建站教育培训建站