ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Orleans 流式编程 API 完全指南:流标识、生产者、消费者与显式/隐式订阅

Orleans 流式编程 API 完全指南:流标识、生产者、消费者与显式/隐式订阅 后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载导读本篇技术指南围绕 Orleans 流式编程的核心 API 展开详细讲解如何通过 Provider name、Stream namespace 与 Stream key 三元组唯一定位一条流如何让 Grain 与 Orleans Client 作为生产者发布消息、作为消费者订阅消息以及如何根据业务订阅语义在ExplicitGrainBasedAndImplicit、ExplicitGrainBasedOnly与ImplicitOnly三种发布/订阅模式之间做出选择。读完本文你将掌握流身份的统一构造方式、显式订阅在 Grain 激活失败后的恢复生命周期、基于[ImplicitStreamSubscription]与IStreamSubscriptionObserver的隐式订阅机制以及 Stateless Worker Grain 参与流消费的约束与行为并理解这些 API 在 Orleans 源码中的实现依据。流标识Stream identity一条 Orleans 流由三个应用层选择共同决定Provider name选择一个已配置的IStreamProvider流提供程序决定事件实际传输与投递的后端Stream namespace将相关流分组并参与隐式订阅的匹配Stream key在命名空间内唯一标识一条流可以是字符串、GUID 或整数。命名空间与 key 共同构成Orleans.Runtime.StreamId。官方推荐把流身份的构建逻辑放在共享代码中避免生产者和消费者各自拼装导致的不一致。以下示例来自 BasicStreaming.cs展示了如何通过共享工厂方法统一获取流public static class TemperatureStreams { public const string ProviderName Telemetry; public const string Namespace device-telemetry; public static IAsyncStreamTemperatureReading Get( Grain grain, string deviceId) { var provider grain.GetStreamProvider(ProviderName); var streamId StreamId.Create(Namespace, deviceId); return provider.GetStreamTemperatureReading(streamId); } }IStreamProvider.GetStream*返回的是一个强类型句柄IAsyncStreamT。获取 Provider 和流句柄是本地操作不会创建任何 broker 实体也不会产生网络往返。生产者与消费者必须就消息类型T达成一致payload 遵循 Orleans 的序列化规则例如示例中的TemperatureReading使用[GenerateSerializer]与[Id(0)]/[Id(1)]标记字段。这也意味着流身份与事件类型解耦同一个StreamId下的消息类型只要双方约定一致即可。生产者Producers任何 Grain 或已配置的 Orleans Client 都可以向流发布消息一条流可以拥有多个生产者。核心发布操作是IAsyncStreamT.OnNextAsync见 BasicStreaming.cs 中的生产者 Grain 示例public sealed class TemperatureProducerGrain : Grain, ITemperatureProducerGrain { private IAsyncStreamTemperatureReading _stream null!; private IGrainTimer? _timer; public override Task OnActivateAsync(CancellationToken cancellationToken) { _stream TemperatureStreams.Get(this, this.GetPrimaryKeyString()); return Task.CompletedTask; } public Task StartAsync() { _timer ?? this.RegisterGrainTimer( PublishAsync, dueTime: TimeSpan.Zero, period: TimeSpan.FromSeconds(5)); return Task.CompletedTask; } private Task PublishAsync(CancellationToken cancellationToken) _stream.OnNextAsync(new TemperatureReading { Celsius Random.Shared.Next(-20, 45), ObservedAt DateTimeOffset.UtcNow, }); }关键语义如下当发布顺序重要时必须await每次OnNextAsync调用。串行等待会按调用顺序提交但并发生产者之间不存在共享顺序返回的 Task 表示Provider 已接受该消息例如持久化 Provider 通常指后端服务已确认并不代表所有消费者都完成了处理Provider 失败以及语义模糊的超时acknowledgement 可能丢失会带来不确定结果需要应用层自行设计重试与去重策略——这一点在 delivery-semantics.md 中有更深入的讨论。消费者Consumers消费者通过实现IAsyncObserverT接口或传入回调委托来订阅流。OnNextAsync回调接收消息项并在 Provider 提供时附带一个StreamSequenceTokenpublic Task OnNextAsync( TemperatureReading item, StreamSequenceToken? token null)消费者回调的完成时机非常关键只有当应用已经承担了对该消息的责任后才应完成OnNextAsync返回的 Task。持久化 Provider 正是依据这个完成信号推进投递位置或在失败后重试避免在回调中阻塞线程。异步的消费者工作会自然地对该订阅施加背压backpressure阻止 Provider 继续投递更多消息。流的投递是**多播multicast**的每个订阅都会收到同一条消息。同一个 Grain 可以对同一条流创建多个显式订阅每个订阅拥有自己独立的StreamSubscriptionHandleT从而拥有独立的投递进度与生命周期。显式订阅与隐式订阅Provider 的StreamPubSubType决定了可用的订阅模型。该枚举定义在 IStreamProviderRuntime.cs三种取值如下Value何时选择该值代价/权衡ExplicitGrainBasedAndImplicit默认应用同时使用两种模型或预期订阅需求会演化灵活性最高。生产者注册时既检查基于 Grain 的 rendezvous也检查隐式 Grain 元数据显式部分需要PubSubStore存储ExplicitGrainBasedOnly所有消费者包括 Client 消费者都使用运行时创建的订阅发现机制聚焦显式订阅。订阅与生产者变更都使用 rendezvous Grain 与PubSubStore应用需要自己管理订阅句柄与恢复ImplicitOnly所有消费者都是声明了[ImplicitStreamSubscription]的 Grain以集群 Grain 元数据作为订阅目录零 rendezvous-Grain 调用、零PubSubStore操作发布/订阅控制平面开销最低凡是需要运行时创建、可单独移除或 Client 参与的订阅都必须切换到支持显式的模式模式决定了订阅发现、协调与存储的工作量而事件传输与投递仍由所选的流 Provider 负责。选择何种模式应基于业务所需的订阅语义并在应用的实际负载下测量效果。修改配置的模式通过Orleans.Hosting.PersistentStreamConfiguratorExtensions.ConfigureStreamPubSub修改 pub/sub 类型后必须重启所有使用该命名 Provider 的 Silo 与 Client并作为协调一致的部署进行确保每个宿主使用相同取值、计算出相同的订阅集合。从源码看ConfigureStreamPubSub的默认参数就是ExplicitGrainBasedAndImplicit见 ClusterClientPersistentStreamConfigurator.cs。切换模式时各订阅模型保持各自的生命周期隐式订阅来自 Grain 元数据。ExplicitGrainBasedOnly只选择显式记录而两种支持隐式的模式都会应用匹配的元数据显式订阅记录遵循所配置PubSubStore的持久性。ImplicitOnly只选择元数据派生的订阅显式记录仍保留在存储中。切换模式前必须核算这些记录若日后在相同 service ID、Provider name 与持久化PubSubStore下重新启用显式支持保留的记录会再次可用每个已激活的消费者会恢复其订阅句柄。当订阅需求预期会演化且持续支持两种模型带来的额外 pub/sub 控制平面开销可以接受时使用默认的组合模式。显式订阅当由应用行为决定是否订阅、何时订阅时使用显式订阅。IAsyncObservableT.SubscribeAsync每次调用都会创建新订阅因此激活代码必须恢复既有订阅句柄而不是再次订阅。示例见 ExplicitSubscriptions.cspublic sealed class ExplicitTelemetryGrain : Grain, IExplicitTelemetryGrain, IAsyncObserverTemperatureReading { private readonly ILoggerExplicitTelemetryGrain _logger; private IAsyncStreamTemperatureReading _stream null!; 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(); } } public Task OnNextAsync( TemperatureReading item, StreamSequenceToken? token null) { /* ... */ } public Task OnErrorAsync(Exception ex) { /* ... */ } }这段代码体现了显式订阅的完整生命周期管理订阅属于 Grain 身份grain identity而非某个激活activation。Grain 激活后通过GetAllSubscriptionHandles()获取全部既有句柄再用ResumeAsync(this)把当前新的 observer 实例挂到每个句柄上首次订阅判断handles.Count 0后才调用SubscribeAsync(this)避免重复创建订阅用StreamSubscriptionHandleT.UnsubscribeAsync()移除订阅。这一生命周期只有在配置了持久化的PubSubStore时才能跨集群重启存活详见 pubsub-storage.md。内存版PubSubStore只在该集群状态仍然可用时保留记录。对于可回放rewindable的持久化 Provider还可以在创建订阅时为 选择订阅起始位置StreamSubscriptionStartPosition.Latest或EarliestAvailable决定从新消息开始还是重放 pulling agent 本地队列缓存中保留的匹配消息。注意起始位置只在订阅创建时生效恢复已有句柄应使用序列令牌推进。结束一个显式订阅结束订阅需要对每一个句柄依次await handle.UnsubscribeAsync()。Orleans 流式运行时会在该操作完成前把每个订阅从 pub/sub 存储中移除并通知活跃的生产者。上面示例中的UnsubscribeAsync方法正是按此顺序执行。隐式订阅当流消息应该按流身份激活某个 Grain时使用隐式订阅——Grain 不需要先被激活或预先订阅消息到来会触发对应 Grain 的激活。示例见 ImplicitSubscriptions.cs[ImplicitStreamSubscription(TemperatureStreams.Namespace)] public sealed class DeviceTelemetryGrain : Grain, IDeviceTelemetryGrain, IAsyncObserverTemperatureReading, IStreamSubscriptionObserver { private readonly ILoggerDeviceTelemetryGrain _logger; private double? _latest; public Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) { var handle handleFactory.CreateTemperatureReading(); return handle.ResumeAsync(this); } public Task OnNextAsync( TemperatureReading item, StreamSequenceToken? token null) { _latest item.Celsius; _logger.LogInformation( Device {DeviceId} reported {Temperature} C, this.GetPrimaryKeyString(), item.Celsius); return Task.CompletedTask; } public Task OnErrorAsync(Exception ex) { _logger.LogError(ex, The telemetry subscription failed); return Task.CompletedTask; } public Taskdouble? GetLatestAsync() Task.FromResult(_latest); }核心机制[ImplicitStreamSubscription]选择流命名空间。对于每个匹配的 Grain 类型Orleans 会把流 key 映射为 Grain key——在上例中流StreamId.Create(device-telemetry, deviceId)会自动路由到IGrainWithStringKey且主键等于deviceId的DeviceTelemetryGrain实现Orleans.Streams.Core.IStreamSubscriptionObserver定义于 IStreamSubscriptionObserver.cs后Orleans 会通过OnSubscribed(IStreamSubscriptionHandleFactory)把隐式句柄交给 GrainGrain 只需调用一次ResumeAsync挂接处理逻辑。隐式订阅的进度语义隐式订阅在 observer 保持挂接期间单调推进一次投递调用成功完成后再次对活跃隐式句柄ResumeAsync传入序列令牌会抛出InvalidOperationException会试图重定位该活跃订阅的位置传null则是在当前位置替换 observer新激活的 Grain 可以用持久化的恢复令牌挂接新提供的句柄如果应用逻辑需要有意重定位投递位置或重放已确认事件应改用显式订阅。从实现看ImplicitOnly模式的ImplicitStreamPubSub通过元数据表解析隐式订阅 ID任何试图在ImplicitOnly下创建显式订阅的操作都会抛出ArgumentOutOfRangeException错误信息明确提示需要改用ExplicitGrainBasedAndImplicit或ExplicitGrainBasedOnly见 ImplicitStreamPubSub.cs。隐式订阅的特性约束隐式订阅声明在 Grain 元数据中不通过SubscribeAsync创建无法在运行时单独移除同一 Grain 绑定不支持多个订阅不存在逐条管理句柄的显式模型。客户端ClientsOrleans Client 在IClientBuilder上配置了 Provider 后可以生产消息并显式消费流。示例配置见 Configuration.cspublic static IHostApplicationBuilder AddStreamingClient( this IHostApplicationBuilder builder) { builder.UseOrleansClient(clientBuilder { clientBuilder.AddMemoryStreams(TemperatureStreams.ProviderName); }); return builder; }Client 订阅的生命周期绑定于其所连接的 Client 进程Client 重连或重启后必须重建订阅。隐式订阅面向 Grain 而非 Client因此 Client 消费必须走显式订阅路径并自行承担句柄管理与恢复责任。Stateless Worker Grain标记了[StatelessWorker]的 Grain 也可以发布和消费流。Stateless Worker 流消费者必须实现IStreamSubscriptionObserver——Orleans 会在每个激活上安装独立的流消费者扩展。投递时优先使用该激活上已有的 observer当没有挂接 observer 时Orleans 调用OnSubscribedGrain 在该回调中调用ResumeAsync挂接一个。从源码文档注释可以确认IStreamSubscriptionObserver.csStateless Worker Grain 会在每次激活时收到该通知以便每个激活都能挂接自己的 observer。关键行为与约束订阅属于 Stateless Worker Grain 身份。对于持久流每个 pulling agent 从其所在 Silo 按该身份投递而正常的 Stateless Worker 放置逻辑会为每次投递尝试选择一个本地激活因此多个并发 pulling agent 可以在不同激活、不同 Silo 上处理消息。每次投递尝试只在一个激活上执行为解码、校验、富化、过滤、转发等无状态转换提供了竞争式消费者competing-consumer执行模型隐式订阅从 Grain 元数据建立 Grain 级订阅显式SubscribeAsync则在运行时建立 Grain 级订阅并挂接调用方激活的 observer。后续激活在首次收到投递时通过OnSubscribed挂接本地 observerUnsubscribeAsync则移除 Grain 级订阅顺序是局部的限于所选激活与 Provider 投递路径。并发投递在不同激活间的完成顺序不定Provider 重试可能选中不同激活因此处理器必须采用无状态或幂等处理并遵循 Provider 的投递保证Stateless Worker 若订阅却未实现IStreamSubscriptionObserver会收到InvalidOperationExceptionStateless Worker observer 以null 序列令牌挂接pulling agent 在投递跨激活移动时为活跃订阅持有进度。向SubscribeAsync或ResumeAsync传入非 null 序列令牌同样产生InvalidOperationException。需要应用管理回退rewind或检查点恢复时应使用常规 Grain 消费者。更进一步关于失败行为与序列令牌请继续阅读 Stream delivery, ordering, replay, and recovery关于显式订阅记录的持久化配置见 Configure PubSub storage关于可回放 Provider 的订阅起始位置选择见 Choose a persistent-stream subscription start position更大的可编译示例可参考 Orleans 测试套件中的SampleStreamingGrain.cs其中包含SampleStreaming_ProducerGrain通过GetStreamProviderGetStreamint获取流并用RegisterGrainTimer周期性发布、SampleStreaming_ConsumerGrainSubscribeAsync显式订阅 UnsubscribeAsync停止消费以及SampleStreaming_InlineConsumerGrain以回调委托OnNextAsync/OnErrorAsync/OnCompletedAsync内联订阅三种典型形态可作为本文所述 API 在真实测试场景中的对照实现。赞分享后端微服务【免费下载链接】orleansCloud Native application framework for .NET项目地址https://gitcode.com/gh_mirrors/or/orleans点击查看免费下载相关推荐cargo-mobile2 安卓开发全攻略从环境配置到 APK 打包的完整流程cargo mobile2 安卓开发全攻略从环境配置到 APK 打包的完整流程 cargo mobile2 是一款让 Rust 移动开发变得简单的工具它为开Julia通道编程掌握Channel与生产者消费者模式的终极指南Julia通道编程掌握Channel与生产者消费者模式的终极指南 Julia作为一种高性能的编程语言提供了强大的并发编程能力其中通道Channel是实操作系统固件驱动开发Apache Pulsar 消息机制完全指南消息、生产者、消费者、订阅与投递语义详解Apache Pulsar 消息机制完全指南消息、生产者、消费者、订阅与投递语义详解 本文以 Apache Pulsar 官方文档 concepts mess消息队列后端流处理上一篇Cubism.js Horizon图表教程如何优雅展示多维度时间序列数据下一篇XXShield单元测试与持续集成确保防Crash代码质量的关键步骤创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表