配置 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(); } }这段代码体现了显式订阅的三个核心操作:
SubscribeAsync:首次订阅时创建订阅记录并持久化到PubSubStore。代码先检查GetAllSubscriptionHandles()是否已有句柄,避免重复创建订阅;ResumeAsync:grain 激活(OnActivateAsync)时,从存储中取回所有订阅句柄并重新挂接当前实例,实现"grain 重启后恢复订阅";UnsubscribeAsync:不再需要时移除订阅。文档建议在订阅不再被需要时主动调用它,防止PubSubStore中堆积无用的订阅记录。
运维指南:备份、命名与替换
官方文档给出四条直接可执行的运维建议:
- 像备份其他应用元数据一样备份并监控
PubSubStore。订阅记录属于业务元数据,丢失后显式订阅需要逐个重建; - 保持提供程序名称稳定。从 pub/sub 系统的视角看,提供程序名称是流身份(stream identity)的一部分,重命名提供程序等于改变了流身份,会导致既有订阅失效;
- 及时清理订阅。通过
StreamSubscriptionHandle<T>.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),仅供参考