【问题标题】:Reactive Extensions: SubscribeOn take no effect when compose with Publish反应式扩展:使用 Publish 撰写时,SubscribeOn 无效
【发布时间】:2012-10-12 05:24:19
【问题描述】:

我是响应式扩展的新手。

我有一个基于房间的 MMOFPS 游戏服务器,其中有很多房间,但只有一个监听套接字。我创建了一个cold observable,它表示从网络接收到的消息,并希望将其转换为hot 以便在多个房间之间共享,以便每个房间都可以过滤和处理它们的相关消息。我的方法正确吗?

另一个问题,在将冷的 observable 转换为热的之后,我意识到 SubscribeOn 失去了它的效果。例如:

var observable = Observable.Return(1).Publish(); 
observable 
    .Where( 
        x => 
            { 
                Console.WriteLine("Filter on {0}", Thread.CurrentThread.ManagedThreadId); 
                return true; 
            }) 
    .SubscribeOn(Scheduler.NewThread) 
    .ObserveOn(Scheduler.NewThread) 
    .Subscribe(x => Console.WriteLine("Received on {0}", Thread.CurrentThread.ManagedThreadId)); 
observable.Connect(); 
Console.WriteLine("End {0}", Thread.CurrentThread.ManagedThreadId); 
Console.ReadLine(); 

结果:
12 日收到
结束 10 过滤 10

没有发布:

结果: 结束 10 过滤 11 12日收到

但是当我使用 Publish.RefCount 自动连接时,它按预期工作。

我错过了什么吗? ...

【问题讨论】:

  • 您希望SubscribeOn 会做什么?
  • 很有可能,当您更改消息的套接字源时,就不需要像现在这样热了。
  • @Enigmativity 我认为 SubscribeOn(Scheduler.NewThread) 将在新线程而不是当前线程上执行“Where”。 “结束 10”和“过滤 10”。应该是有区别的。而且我的套接字并不是真正的.Net Socket,我使用Lidgren,并且没有机会访问它使用的.Net 套接字。我的方法是在后台线程中不断轮询 Lidgren 以获取传入的数据包并产生无穷大的 observable。我的方法正确吗?我真的是响应式扩展的新手。
  • 您需要使用ObserveOn 而不是SubscribeOn。后者仅导致观察者的订阅发生在提供的调度程序上。前者是安排Where 的地方。
  • 您需要使用某种形式的FromEventPattern/FromAyncPattern 来获取可观察到的消息来源。如果你这样做,它们会很热。

标签: .net system.reactive


【解决方案1】:

应用运算符的顺序很重要。如果你想共享 SubscribeOn 的订阅副作用,那么你必须在 Publish 之前应用它。例如:

var observable = Observable
    .Return(1)
    .SubscribeOn(Scheduler.NewThread)
    .Publish();

更多信息: http://social.msdn.microsoft.com/Forums/pl-PL/rx/thread/29de2890-8303-404d-a6f8-f0fe0f716a86

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-02-10
    • 2011-08-27
    相关资源
    最近更新 更多