讲师中心 微信公众号
AI工具推荐 视频效率加速

c# Orleans 的 Stream 和 Observer 并发模型

老辰小哥_2025

老辰小哥_2025

发布时间:2026-01-03 12:48:08

|

178人浏览过

|

来源于php中文网

原创

IAsyncObserver.OnNextAsync 不会并发调用,因其由 StreamPullingAgent 单消费者拉取机制与 grain 单线程调度保证严格串行;需异步实现避免阻塞,否则导致流延迟堆积。

c# orleans 的 stream 和 observer 并发模型

Orleans 的 IAsyncStream<T>IAsyncObserver<T> 不是线程安全的并发模型,而是单线程、顺序交付的虚拟流抽象——所有事件按发布顺序、在 grain 或 observer 所属的逻辑上下文中串行处理。

为什么 IAsyncObserver.OnNextAsync 不会并发调用

Orleans 流设计强制“每订阅一个 observer,其 OnNextAsync 调用严格串行化”,即使底层有多个 producer 并发推送,或多个 silo 同时投递消息。这是由 StreamPullingAgent 的单消费者拉取机制 + grain 激活上下文的单线程调度保证的。

  • 你无需加锁或手动同步 OnNextAsync 内部状态(比如累加计数器、更新缓存)
  • 但这也意味着:如果某个 OnNextAsync 执行耗时(如同步 DB 写入、阻塞 IO),整个流订阅将被阻塞,后续消息延迟堆积
  • 错误示例:
    public Task OnNextAsync(MyEvent e, StreamSequenceToken token = null)
    {
        // ❌ 危险:同步写数据库,阻塞流处理
        _dbContext.Events.Add(e);
        _dbContext.SaveChanges(); // 阻塞线程,拖慢整个订阅
        return Task.CompletedTask;
    }
  • 正确做法:始终用异步路径,让控制权及时交还调度器
    public async Task OnNextAsync(MyEvent e, StreamSequenceToken token = null)
    {
        // ✅ 推荐:异步持久化,不阻塞
        await _dbContext.Events.AddAsync(e);
        await _dbContext.SaveChangesAsync();
    }

GetStream<T> 生成句柄是本地的,但流语义跨集群共享

调用 streamProvider.GetStream<T>(streamId) 只是创建一个轻量级、无网络开销的逻辑句柄;真正的流生命周期、消息路由、订阅管理由 Orleans 运行时在集群中协调。这意味着:

在SEO发布前,从路由清单生成XML网站地图和robots.txt
在SEO发布前,从路由清单生成XML网站地图和robots.txt

当代理已经知道网站路由或内容URL,并且在启动前需要有效的sitemap XML、sitemap索引或robots.txt引用时,请使用sitemap。这是一个发布构件技能,而不是爬虫或SEO平台。

下载
  • streamId 必须全局唯一:靠 StreamId.Create("namespace", guid) 组合保证——guid 建议来自 grain ID(如 this.GetPrimaryKey()),namespace 区分业务域(如 "order-events""chat-messages"
  • 同一 streamId 在不同 silo 上调用 GetStream,拿到的是指向同一个逻辑流的句柄,不是各自独立的副本
  • Producer 和 Consumer 可以部署在完全不同的 silo 上,Orleans 自动完成消息投递和订阅发现(依赖配置的 pub/sub 存储,如 PubSubStore
  • 常见坑:忘记在 silo 配置中启用流提供程序和 pub/sub 存储
    // SiloBuilder 必须包含这两行
    silo.AddMemoryStreams("SimpleStreamProvider")
        .AddMemoryGrainStorage("PubSubStore"); // 否则 SubscribeAsync 会静默失败或超时

背压不是自动的,得靠配置+代码配合

Orleans 默认不拒绝或缓冲过载消息;它依赖两种机制协同实现背压:LoadShedQueueFlowController(CPU 触发限流)和 BatchContainerBatchSize(控制每次拉取条数)。但它们都需显式开启。

  • 默认情况下,OnNextAsync 调用会排队等待执行,没有上限——内存可能被撑爆
  • 启用 CPU 背压:
    options.LoadSheddingEnabled = true;
    options.LoadSheddingLimit = 0.8; // CPU > 80% 时暂停接收新消息
  • 减小批处理大小降低内存压力(尤其对大消息流):
    services.Configure<StreamPullingAgentOptions>(opts =>
    {
        opts.BatchContainerBatchSize = 1; // 默认是 1,设为 5 可提升吞吐但增加延迟
    });
  • 真正关键的背压点在 consumer 端代码里:别让 OnNextAsync 成为瓶颈,该异步就异步,该降频就降频(比如用 RegisterTimer 批量消费)

最易被忽略的一点:流的“顺序性”只保证 per-subscription,不保证 per-producer 或全局全序。如果你从 3 个不同 grain 并发调用 stream.OnNextAsync,consumer 收到的顺序取决于网络延迟、silos 负载、pulling agent 调度时机——想强序,必须用同一个 grain 发送,或引入外部序列号+客户端排序逻辑。

热门AI工具

更多
WorkBuddy

一款AI办公效率工具,主要用于腾讯云推出的AI原生桌面智能体工作台,适合需要提升相关任务效率的用户。

蛙蛙写作

一款AI论文写作工具,主要用于超级AI智能写作助手,适合需要提升相关任务效率的用户。

Atoms
Atoms Hot

Atoms是一款AI智能体工具,第一支自动构建真实业务的 AI 团队。

DeepSeek

DeepSeek是一款面向对话、写作、编程和推理场景的AI大模型工具。

豆包大模型

豆包大模型是一款由字节跳动推出的企业级大语言模型服务平台。

咔片AIPPT

一款在线AI演示文稿制作工具,可根据主题和内容需求辅助生成PPT结构与页面,提高演示材料制作效率。

切问学术

切问学术是一款AI论文写作工具,复旦大学NLP团队推出的AI学术智能体。

讯飞智作

讯飞智作是一款AI视频创作工具,AI文本配音工具,数字人课程、营销视频制作。

火山引擎

火山引擎是一款面向企业的云计算与AI服务平台。

相关专题

更多
堆和栈的区别
堆和栈的区别

堆和栈的区别:1、内存分配方式不同;2、大小不同;3、数据访问方式不同;4、数据的生命周期。本专题为大家提供堆和栈的区别的相关的文章、下载、课程内容,供大家免费下载体验。

4567

2023.07.18

堆和栈区别
堆和栈区别

堆(Heap)和栈(Stack)是计算机中两种常见的内存分配机制。它们在内存管理的方式、分配方式以及使用场景上有很大的区别。本文将详细介绍堆和栈的特点、区别以及各自的使用场景。php中文网给大家带来了相关的教程以及文章欢迎大家前来学习阅读。

2088

2023.08.10

线程和进程的区别
线程和进程的区别

线程和进程的区别:线程是进程的一部分,用于实现并发和并行操作,而线程共享进程的资源,通信更方便快捷,切换开销较小。本专题为大家提供线程和进程区别相关的各种文章、以及下载和课程。

3478

2023.08.10

Buffalo框架数据库开发全教程
Buffalo框架数据库开发全教程

本专题围绕Buffalo框架数据库开发,讲解database.yml多环境配置、soda与fizz迁移生成回滚、模型结构体标签、增删改查与条件查询、一对多与多对多关联、数据校验、回调钩子、事务处理及原生SQL执行能力。

0

2026.09.23

Buffalo框架路由与请求处理实操指南
Buffalo框架路由与请求处理实操指南

本专题讲解Buffalo框架路由与请求处理机制,涵盖路由注册与分组、资源路由、Handler编写规范、Context上下文方法、参数绑定、中间件编写挂载、Session与Cookie读写、Flash消息及错误页面定制方法。

0

2026.09.23

Buffalo框架零基础入门教程
Buffalo框架零基础入门教程

本专题整理Buffalo框架入门内容,涵盖Go环境准备、buffalo CLI安装、新项目生成、目录结构说明、dev热加载启动、数据库连接配置与常见报错排查,帮助新手按约定优于配置的思路跑通第一个Buffalo框架应用。

0

2026.09.23

Conan创建软件包配方指南
Conan创建软件包配方指南

本专题介绍通过conanfile.py创建软件包的方法,讲解包名、版本、依赖和构建设置等基础信息,以及source、build、package、package_info等常用方法的作用及编写思路。

0

2026.09.22

Conan二进制包配置指南
Conan二进制包配置指南

本专题介绍Conan根据操作系统、编译器、架构和构建类型生成二进制包的方法,讲解Profile、Settings、Options及Package ID的作用,帮助管理不同平台和编译环境下的包版本。

0

2026.09.22

Conan私有仓库搭建教程
Conan私有仓库搭建教程

本专题系统的讲解Conan私有仓库的搭建流程,涵盖仓库服务部署、存储目录配置、用户认证、权限划分和远程地址添加,并介绍内部C++依赖包的上传、下载及版本维护方法。

0

2026.09.22

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
相关推荐
/
热门推荐
/
最新课程
关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn