后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载Azure Queue Storage 是 Microsoft Orleans 持久化流persistent streams的官方队列后端之一。Microsoft.Orleans.Streaming.AzureStorage包提供了一个基于 Azure 队列存储的流适配器它遵循 Orleans 通用的持久化流拉取架构并实现了 Azure 特有的队列映射、消息编码、接收确认与删除语义。本文以该适配器为主线结合仓库源码src/Azure/Orleans.Streaming.AzureStorage与测试用例test/Extensions/Orleans.Azure.Tests/Streaming讲解注册 API、适配器行为、接收与确认流程、队列映射默认值、可见性调优及扩展点读完你可以独立完成该提供者的配置、调优与自定义数据适配器开发。关于持久化流拉取架构的整体机制队列均衡、拉取代理、缓存、游标、pub-sub、恢复请先阅读 Persistent stream pulling architecture本文聚焦 Azure Queue 这一具体实现。注册 APISilo 与客户端构建器Microsoft.Orleans.Streaming.AzureStorage包通过两个扩展方法注册流提供者Silo 端SiloBuilderExtensions.AddAzureQueueStreams客户端端ClientBuilderExtensions.AddAzureQueueStreams两者的简洁重载都接受一个提供者名称和ActionOptionsBuilderAzureQueueOptions配置委托。完整实现见 silo 注册 与 client 注册。Silo 与客户端配置器都允许替换队列数据适配器通过ConfigureQueueDataAdapter但只有 silo 配置器暴露缓存与拉取代理组件如ConfigureCacheSize因为拉取代理pulling agent和队列缓存只运行在 silo 上——客户端只承担生产者的角色。一个完整的 silo 注册示例参考 README.mdusing Microsoft.Extensions.Hosting; using Orleans.Hosting; using Orleans.Streams; var builder Host.CreateApplicationBuilder(args) .UseOrleans(siloBuilder { siloBuilder .UseLocalhostClustering() .AddAzureQueueStreams( name: AzureQueueStreamProvider, b b.ConfigureAzureQueue(ob ob.Configure((options, dep) { options.QueueServiceClient new QueueServiceClient(UseDevelopmentStoragetrue); options.QueueNames Enumerable.Range(0, 8) .Select(num ${dep.Value.ClusterId}-{num}).ToList(); }))); }); await builder.RunAsync();安装包dotnet add package Microsoft.Orleans.Streaming.AzureStorage通过配置驱动注册ProviderType: AzureQueueStorage除了代码注册该提供者还实现了基于配置的注册AzureQueueStreamProviderBuilder通过程序集特性RegisterProvider(AzureQueueStorage, Streaming, ...)对外暴露silo 与客户端均可使用见 AzureQueueStreamProviderBuilder.cs。配置文件写法如下{ Orleans: { Streaming: { AzureQueueProvider: { ProviderType: AzureQueueStorage, ServiceKey: myQueueClient, // 可选DI 中已注册的 keyed QueueServiceClient ConnectionName: AzureQueue, // 可选指向连接字符串配置 ConnectionString: UseDevelopmentStoragetrue, // 可选直接连接字符串或队列服务 URI MessageVisibilityTimeout: 00:00:37, // 可选TimeSpan 字符串 QueueNames: [ q1, q2 ] // 可选显式队列名列表 } } } }配置解析逻辑见 AzureQueueStreamProviderBuilder.cs按以下优先级构造QueueServiceClientServiceKey从服务容器中取出已注册的 keyedQueueServiceClientConnectionName从根IConfiguration的GetConnectionString读取连接字符串ConnectionString直接作为连接字符串使用。若其能被解析为绝对 URI带查询字符串SAS URI→new QueueServiceClient(uri)不带查询字符串 → 使用DefaultAzureCredential认证否则 → 按普通连接字符串构造。该行为由测试 AzureQueueStreamProviderBuilderTests.cs 覆盖Minimal_Configuration验证UseDevelopmentStoragetrue会生成 AccountName 为devstoreaccount1的客户端Full_Configuration验证MessageVisibilityTimeout: 00:00:37被解析为TimeSpan.FromSeconds(37)Missing_ConnectionString验证缺少凭据时QueueServiceClient为null。AzureQueueOptions 关键属性AzureQueueOptions定义于 AzureQueueStreamOptions.cs属性说明QueueServiceClient访问 Azure Queue Service 的客户端直接赋值即可推荐的现代方式QueueNames用于划分提供者的队列名列表未配置时由 silo/client 配置器根据服务 ID 与提供者名自动生成MessageVisibilityTimeout接收后消息对其他消费者不可见的时长null表示使用 Azure Queue Service 默认值代码中的ConfigureQueueServiceClient(...)各重载与ClientOptions属性均已标记[Obsolete]新代码应直接设置QueueServiceClient属性。AzureQueueOptionsValidator会在启动时校验未配置客户端或QueueNames为空都会抛出OrleansConfigurationException。适配器行为读-写、不可重绕、默认 V2 编码AzureQueueAdapterAzureQueueAdapter.cs是IQueueAdapter的实现Direction为StreamProviderDirection.ReadWrite既生产又消费IsRewindable为false——Azure Queue Storage 不暴露持久的任意流偏移因此传入非 null 的StreamSequenceToken会直接抛出ArgumentException见QueueMessageBatchAsync的第 44 行断言。生产路径生产者调用QueueMessageBatchAsync适配器先用IQueueDataAdapterstring, IBatchContainer将流标识StreamId、事件负载与请求上下文编码为一条 Azure 队列消息再按流到队列的映射写入对应队列。队列客户端按需创建并通过ConcurrentDictionaryQueueId, AzureQueueDataManager缓存。编码层使用IQueueDataAdapterstring, IBatchContainerAzureQueueDataAdapterV2默认使用AzureQueueBatchContainerV2与EventSequenceTokenV2支持 Orleans 自定义序列化器如 JSON见 IAzureQueueDataAdapter.csAzureQueueDataAdapterV1保留用于兼容历史消息不是新提供者的推荐格式。接收端并不携带序列号——AzureQueueAdapterReceiver在解码时用本地单调递增的lastReadMessage为每个批次分配EventSequenceTokenV2见 AzureQueueAdapterReceiver.cs。因此序列令牌是接收时分配的本地序列用于传递与确认而非持久化的外部偏移。接收与确认可见性、pop receipt 与删除AzureQueueAdapterReceiverAzureQueueAdapterReceiver.cs实现了拉取代理与 Azure Queue Storage 之间的交互每次调用GetQueueMessagesAsync最多向 Azure 请求32 条可见消息常量MaxNumberOfMessagesToPeekAzure 的 Get Messages 操作 会让被接收的消息暂时不可见但不删除它同时返回每条消息的pop receipt拉取代理解码并交付这批消息经由队列缓存queue cache与各订阅游标cursor分发给消费者只有被报告为已交付的消息才会通过 Delete Message 操作携带 pop receipt真正删除MessagesDeliveredAsync→DeleteQueueMessage。完整确认时序如下引自原文档优雅的队列所有权交接在队列所有权转移如 silo 退出、队列重新均衡时旧接收者的Shutdown会调用ReleaseMessagesAsync对每一条已预取但尚未确认的消息执行 Update Message 操作将其可见性超时设为TimeSpan.Zero使消息立即可见。这样接管该队列的新接收者无需等待可见性超时即可继续交付见 AzureQueueAdapterReceiver.cs 与 AzureQueueDataManager.ReleaseQueueMessage。该操作需要凭据授权——若使用队列服务 SAS必须包含u更新权限。如果更新失败Orleans 会记录释放失败警告日志AzureQueue_18消息会等现有可见性租约到期后由 Azure 自动恢复可见。失败与重复交付若接收者、拉取代理或 silo 在交接完成前崩溃可见性超时最终到期Azure 会再次返回该消息——消费者必须容忍重复交付。此外若可见性在消息处理期间过期pop receipt 会失效此时删除可能失败接收者在MessagesDeliveredAsync中会对删除异常记录警告并忽略AzureQueue_15消息随后会被重新接收。这与 delivery semantics 中持久化流并非普遍恰好一次的结论一致执行持久化副作用写数据库、发外部通知等的消费者应具备幂等性。该行为由 AzureQueueAdapterTests.cs 中的ShutdownReleasesPendingMessages测试直接验证先由初始接收者收到消息但不确认随后Shutdown再由替换接收者成功重新收到同一条消息并确认删除。队列映射与默认值当AzureQueueOptions.QueueNames未显式提供时配置器会在PostConfigure阶段通过AzureQueueStreamProviderUtils.GenerateDefaultAzureQueueNames(serviceId, providerName)自动生成队列名见 AzureQueueStreamBuilder.cs 与 AzureQueueStreamProviderUtils.cs底层使用HashRingBasedStreamQueueMapper默认哈希环包含8 个队列生成的队列名形如{serviceId}-{queueName}保证不同服务/提供者的队列互不干扰每个被拥有的队列在其当前 silo 上各有一个拉取代理与一个队列缓存。该提供者沿用的持久化流公共默认值设置默认值HashRingStreamQueueMapperOptions.TotalQueueCount8StreamPullingAgentOptions.GetQueueMsgsTimerPeriod100 msSimpleQueueCacheOptions.CacheSize4,096 个批次容器StreamPullingAgentOptions.MaxEventDeliveryTime1 分钟StreamPullingAgentOptions.StreamInactivityPeriod30 分钟队列数量决定了拉取并行度与所有权转移的粒度队列越多均衡越细但每队列的拉取开销也越多。更改队列名会改变物理分区集合应视为数据迁移而非例行调优——新旧队列名之间的事件流无法衔接未消费消息会留在旧队列中。默认队列命名规则在 AzureQueueDataManager.cs 中实现队列名会先被小写化并替换/ \ # ? : . %等非法字符再校验长度3–63、首尾字符、可用字符集字母、数字、连字符与禁止连续--等 Azure 命名约束校验失败会抛出ArgumentException。可见性与缓存进度两类时钟的调优Azure 的可见性超时与 Orleans 的队列缓存保留是两个不同的时钟可见性控制一条未删除的 Azure 消息何时可以再次被接收队列缓存控制拉取代理为活跃订阅游标保留批次的时间长度。调优原则对应原文档MessageVisibilityTimeout若短于最坏情况交付耗时会增加重复接收与失效 pop receipt 的风险消息还在处理中就已重新可见或删除时 receipt 过期若过长则代理故障后恢复会变慢消息长期不可见取值应来自实测交付延迟与故障恢复目标RTO而非拍脑袋。整体运行调优请参考 stream provider guidance。底层还有一个值得注意的默认重试策略AzureQueueDefaultPolicies见 AzureQueueDataManager.cs定义队列操作最多重试 5 次、重试间隔 100 ms、单次操作超时约 3 秒单次慢操作超过阈值会记录AzureQueue_13警告日志便于发现存储端异常。扩展与兼容点自定义数据适配器若要在保留 Azure 传输层的同时改变负载编码可实现自定义IQueueDataAdapterstring, IBatchContainer并通过ConfigureQueueDataAdapter注册siloBuilder.AddAzureQueueStreams(AzureQueueStreamProvider, b { b.ConfigureAzureQueue(ob ob.Configure(options { options.QueueServiceClient new QueueServiceClient(...); })); b.ConfigureQueueDataAdapterMyCustomDataAdapter(); });自定义适配器必须保留流标识StreamId、线协议wire contract定义的请求上下文值以及消费者使用的序列信息。可运行的完整实现与版本化滚动发布指导见 Customize persistent-stream data formats。实验性 JSON 流仓库还提供了实验性的 JSON 序列化入口AddAzureQueueJsonStreamssilo 与客户端均支持见 SiloBuilderExtensions.cs通过AzureQueueJsonDataAdapter实现紧凑 JSON 信封编码并支持在 JSON 与二进制之间按PreferJson/EnableFallback选项切换与回退。该特性标记为[Experimental(StreamingJsonSerializationExperimental)]API 可能在未来版本变化生产环境采用前需评估。测试覆盖配置绑定行为AzureQueueStreamProviderBuilderTests测试文件覆盖缺省连接串、最小配置、完整配置三种 JSON 场景适配器确认与游标行为AzureQueueAdapterTests测试文件验证发送/接收往返、批次中事件类型的拆分int/string 各半、缓存游标从任意令牌开始遍历的顺序性以及关闭时释放待确认消息供新接收者重投。小结Orleans 的 Azure Queue 流适配器把 Azure 队列存储的可见性 pop receipt模型嫁接到持久化流拉取架构上默认 8 个哈希环队列负责分区接收者每次最多拉取 32 条可见消息只有确认交付后才删除所有权交接时通过 Update Message 立即释放未确认消息。配置上推荐直接设置QueueServiceClient并显式或自动生成QueueNames按实测延迟调优MessageVisibilityTimeout且把更改队列名当作数据迁移对待。由于 Azure 队列不提供持久偏移且可能重复投递生产消费者必须做到幂等。赞分享后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载相关推荐EMQX MQTT Message Queue$q/$queue 消息队列深入解析架构、数据流、配置与 HTTP API 实战EMQX MQTT Message Queue$q/$queue 消息队列深入解析架构、数据流、配置与 HTTP API 实战 EMQX 的 MQTT M后端物联网消息队列通信Azure Queue Storage 的 Rust 客户端实战基于 azure_storage_queue 实现队列消息的发送、接收与 RBAC 鉴权Azure Queue Storage 的 Rust 客户端实战基于 azure_storage_queue 实现队列消息的发送、接收与 RBAC 鉴权 导读AI 技能AI 插件EMQX MQTT Ingress 桥接队列订阅$queue与 MQTT 5 Subscription Identifiers 深度解析EMQX MQTT Ingress 桥接队列订阅 $queue 与 MQTT 5 Subscription Identifiers 深度解析 MQTT in后端物联网消息队列通信上一篇实测对比Ling-2.6-flash在BFCL-V4/TAU2-bench等5大基准的SOTA表现下一篇终结iOS假联网Connectivity 6.x全解析与实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考