【问题标题】:Reactive Extensions - raising async events and subscribing on specific threads反应式扩展 - 引发异步事件并订阅特定线程
【发布时间】:2011-11-12 19:14:19
【问题描述】:

我有一系列使用 RX 发布/订阅模型的模块。

这里是事件注册码(每个订阅模块重复):

_publisher.GetEvent<DataEvent>()                 
    .Where(sde => sde.SourceName == source.Name) 
    .ObserveOn(Scheduler.TaskPool)
    .Subscribe(module.OnDataEvent);  

发布者简单,感谢José Romaniello's code

public class EventPublisher : IEventPublisher
{
    private readonly ConcurrentDictionary<Type, object> _subjects = 
        new ConcurrentDictionary<Type, object>(); public IObservable<TEvent> GetEvent<TEvent>() 
    { 
        var subject = (ISubject<TEvent>)_subjects.GetOrAdd(typeof(TEvent), t => new Subject<TEvent>()); 
        return subject.AsObservable(); 
    }
    public void Publish<TEvent>(TEvent sampleEvent)
    {
        object subject; 
        if (_subjects.TryGetValue(typeof(TEvent), out subject)) 
        { 
            ((ISubject<TEvent>)subject).OnNext(sampleEvent); 
        }
    }
}

现在我的问题:正如您在上面看到的,我使用 .ObserveOn(Scheduler.TaskPool) 方法为每个模块、每个事件从池中分离出一个新线程。这是因为我有很多事件和模块。当然,问题在于事件在时间顺序上混淆了,因为一些事件彼此靠近触发,然后最终以错误的顺序调用 OnDataEvent 回调(每个 OnDataEvent 都带有一个时间戳)。

有没有一种简单的方法来使用 RX 来确保事件的正确顺序?或者我可以编写自己的调度程序来确保每个模块按顺序获取事件?

当然,事件是按正确的顺序发布的。

提前致谢。

【问题讨论】:

    标签: c# multithreading system.reactive publish-subscribe observer-pattern


    【解决方案1】:

    尝试使用EventPublisher的这个实现:

    public class EventPublisher : IEventPublisher
    {
        private readonly EventLoopScheduler _scheduler = new EventLoopScheduler(); 
        private readonly Subject<object> _subject = new Subject<object>();
    
        public IObservable<TEvent> GetEvent<TEvent>() 
        { 
            return _subject
                .Where(o => o is TEvent)
                .Select(o => (TEvent)o)
                .ObserveOn(_scheduler);
        }
    
        public void Publish<TEvent>(TEvent sampleEvent)
        {
            _subject.OnNext(sampleEvent);
        }
    }
    

    它使用EventLoopScheduler 来确保所有事件按顺序在同一个后台线程上发生。

    从您的订阅中删除ObserveOn,因为如果您在另一个线程上观察,您可能会再次以错误的顺序发生事件。

    这能解决您的问题吗?

    【讨论】:

    • 为了与原始示例保持一致,您可以通过对位进行观察并在最初的订阅模块中声明调度程序来为每个观察者设置一个线程,只需为每个模块使用一个 eventloopscheduler 而不是任务池.
    • 谢谢大家,我会试试这个并报告!
    【解决方案2】:

    尝试同步方法,如:

    _publisher.GetEvent<DataEvent>()                 
        .Where(sde => sde.SourceName == source.Name) 
        .ObserveOn(Scheduler.TaskPool).Synchronize()
        .Subscribe(module.OnDataEvent);
    

    虽然我使用相同的代码尝试了您的方案,但发现数据按顺序到达并且不重叠。可能这是您的应用程序特有的东西。

    【讨论】:

    • Synchronize 方法除了确保底层的 observable 遵守 observable 合约之外什么都不做——即OnNext*(OnError|OnCompleted)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-05-15
    • 1970-01-01
    • 2023-03-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多