首页/新闻资讯/正文详情

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

发布时间:2026/9/24 21:41:35 来源:云帆数科 栏目:资讯中心
配置 Orleans PubSub 存储:让流订阅元数据在集群重启后依然存活
配置 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),仅供参考

相关推荐

Akka Classic Cluster Metrics 扩展实战指南:集群指标采集、自适应负载均衡与 Sigar 配置
Akka Classic Cluster Metrics 扩展实战指南:集群指标采集、自适应负载均衡与 Sigar 配置

后端并发编程异步编程 【免费下载链接】akka-core A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments. 项目地址: https://gitcode.com/gh_mirrors/ak/akka-core 点击查看 免费下载 导读 本文基… · 2026/9/24 21:41:35

@formily/reactive 的 markRaw:彻底掌控响应式劫持边界的权威指南
@formily/reactive 的 markRaw:彻底掌控响应式劫持边界的权威指南

formily/reactive 的 markRaw:彻底掌控响应式劫持边界的权威指南 【免费下载链接】formily 📱🚀 🧩 Cross Device & High Performance Normal Form/Dynamic(JSON Schema) Form/Form Builder -- Support React/React Native/Vu… · 2026/9/24 21:41:28

JVM中的klass与Class对象:类加载后内存里到底放了什么?
JVM中的klass与Class对象:类加载后内存里到底放了什么?

上个月帮朋友做模拟面试,我问了个自认为很基础的问题:“你天天用的HashMap.class和它背后方法区里的类元数据,到底是同一个东西吗?”对方想了一会儿,答:“Class 对象不是存在方法区吗?”这个答案… · 2026/9/24 21:41:21

AI日报系统设计原理与工程实践指南
AI日报系统设计原理与工程实践指南

我无法生成关于“AI 日报(2026年9月17日)”的博文内容。原因如下:该标题缺乏实质性项目信息——无具体技术动作、无明确功能目标、无实际场景描述、无代码/工具/流程线索,仅是一个时间戳加泛称“AI 日报”,且配套的【项… · 2026/9/24 22:11:39

云端大模型+本地小模型:用GPT-6编排MiniCPM构建研究智能体实战
云端大模型+本地小模型:用GPT-6编排MiniCPM构建研究智能体实战

面壁智能公开点赞GPT-6,还把自家的MiniCPM5-2B交给GPT-6来编排,组合成一个本地研究智能体,这消息在Agent开发者圈子里传开后,我看到最多的评论是:GPT-6都强成那样了,为什么还要在本地跑一个2B的小模型&… · 2026/9/24 22:11:39

NeoHorse-1黑马解析:Harness与RSI如何让Agent稳定自愈
NeoHorse-1黑马解析:Harness与RSI如何让Agent稳定自愈

1. 这匹“黑马”到底踩中了什么痛点AI 圈每隔几个月就会冒出一个新名字,但大多数热闹三天就散了。NeoHorse-1 这次能被讨论,我认为核心不在于它跑分多高,而在于它把RSI、Harness、Agent这三个原本各说各话的概念拧成了一股绳。先说结论&#… · 2026/9/24 22:11:39

314张绝缘子缺陷图训练YOLO:小样本目标检测与边缘部署避坑指南
314张绝缘子缺陷图训练YOLO:小样本目标检测与边缘部署避坑指南

简介:本资源面向电力巡检与计算机视觉方向的开发者、研究生及算法工程师,提供一套可直接用于YOLO目标检测训练的绝缘子缺陷数据集,帮助解决输电线路绝缘子缺陷样本稀缺、标注成本高的问题。压缩包共1257个文件,约121.86MB&#xf… · 2026/9/24 22:11:39

专精特新企业系统化跃迁:从单点冠军到系统冠军的实战框架
专精特新企业系统化跃迁:从单点冠军到系统冠军的实战框架

这些年我走访过不少制造型企业,从几十人的隐形冠军“苗子”到营收过亿的省级专精特新企业都有接触。大家普遍有一个共同的困惑:企业做到一定规模后,单点突破的红利吃完了,再往上走感觉处处是瓶颈——产品有竞争力但客户结构单一&a… · 2026/9/24 22:11:33

机房可视化运维实战:从监控架构到告警降噪的主动管理指南
机房可视化运维实战:从监控架构到告警降噪的主动管理指南

机房管理这事儿,干得越久越觉得它像个“盲盒”。平时看着一切正常,LED灯全绿,空调嗡嗡响,你也说不清到底哪个环节会在哪一刻出幺蛾子。直到业务群开始有人喊“系统好慢”“连不上了”,你才带着笔记本一路小跑冲进机房&… · 2026/9/24 22:11:33

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程
基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为… · 2026/9/24 0:00:13

1D-CNN时间序列建模实战:从Conv1d原理到工业落地
1D-CNN时间序列建模实战:从Conv1d原理到工业落地

简介:面向时间序列数据建模的一维卷积神经网络完整实现,适合深度学习入门者及需要快速验证时序模型的研究者,能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小,只有3KB,内含3个Python脚… · 2026/9/24 0:00:26

柔软的L:汉语语流中被忽视的舌肌张力控制
柔软的L:汉语语流中被忽视的舌肌张力控制

1. 这个“L”不是字母表里的L,而是舌尖上的L最近在几个方言群和语音教学社群里,反复看到有人发一句:“也说字母L:柔软的长舌”。初看以为是英语发音课笔记,点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码