配置 Orleans PubSub 存储:让流订阅元数据在集群重启后依然存活
2026/9/24 21:41:31 网站建设 项目流程

配置 Orleans PubSub 存储:让流订阅元数据在集群重启后依然存活

【免费下载链接】orleansCloud Native application framework for .NET项目地址: https://gitcode.com/gh_mirrors/or/orleans

Orleans 流(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作为其持久化状态;
  • 流检查点 grain(StreamCheckpointGrain.cs)通过[PersistentState(StateName, ProviderConstants.DEFAULT_PUBSUB_PROVIDER_NAME)]直接绑定到PubSubStore,用于持久化持久流(persistent stream)队列的消费位置;
  • grain 检查点器(GrainStreamQueueCheckpointer.cs)默认使用PubSubStore作为检查点存储。

也就是说,只要你使用持久流提供程序(如 Azure Event Hubs、Azure Queue、Kinesis、SQS)并启用 grain 检查点(UseGrainCheckpointer),PubSubStore就承担着订阅元数据 + 队列消费位置双重持久化职责。即使你只使用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, IAsyncObserver<TemperatureReading>, IStreamSubscriptionObserver { public Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) { var handle = handleFactory.Create<TemperatureReading>(); 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)); });

连接字符串方式适合本地开发、测试以及无法使用托管标识的受限环境。同样的配置模式也适用于其他持久化后端,例如AddDynamoDBGrainStorage(Orleans.Clustering.DynamoDB)、AddAdoNetGrainStorage(Orleans.Persistence.AdoNet)或AddCosmosGrainStorage(Orleans.Persistence.Cosmos),只需将存储提供程序名称指定为PubSubStore即可。

集群身份与存储的对应关系

官方文档强调了一条关键原则:使用稳定的 Orleans service ID,并在集群重启之间保持相同的持久化配置

服务 ID(service ID)是 Orleans 集群的逻辑标识。pub/sub 订阅记录的存储键派生自服务 ID,因此:

  • 修改 service ID → 从 pub/sub 系统的角度看,订阅注册表变成"逻辑上全新"的,旧订阅记录不再被新集群识别;
  • 删除或重建底层表 → 订阅记录全部丢失,相当于新建注册表;
  • 存储配置不一致 → 不同 silo 可能读写不同的表,导致订阅状态不一致。

生产环境中应把 service ID 视为需要刻意维护、保持不变的部署标识。

订阅生命周期:激活、恢复与移除

持久化PubSubStore只保证订阅记录("谁订阅了哪个流")得以保存,但它不保存消费者的 observer 实例。文档明确指出:

即使使用持久化的PubSubStore,显式消费者在激活后也必须调用StreamSubscriptionHandle<T>.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(); } }

这段代码体现了显式订阅的三个核心操作:

  1. SubscribeAsync:首次订阅时创建订阅记录并持久化到PubSubStore。代码先检查GetAllSubscriptionHandles()是否已有句柄,避免重复创建订阅;
  2. ResumeAsync:grain 激活(OnActivateAsync)时,从存储中取回所有订阅句柄并重新挂接当前实例,实现"grain 重启后恢复订阅";
  3. UnsubscribeAsync:不再需要时移除订阅。文档建议在订阅不再被需要时主动调用它,防止PubSubStore中堆积无用的订阅记录。

运维指南:备份、命名与替换

官方文档给出四条直接可执行的运维建议:

  1. 像备份其他应用元数据一样备份并监控PubSubStore。订阅记录属于业务元数据,丢失后显式订阅需要逐个重建;
  2. 保持提供程序名称稳定。从 pub/sub 系统的视角看,提供程序名称是流身份(stream identity)的一部分,重命名提供程序等于改变了流身份,会导致既有订阅失效;
  3. 及时清理订阅。通过StreamSubscriptionHandle<T>.UnsubscribeAsync()移除不再需要的订阅,控制存储增长与系统开销;
  4. 替换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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询