【问题标题】:Observable LINQ inconsistent exceptions thrown抛出可观察到的 LINQ 不一致异常
【发布时间】:2016-02-01 22:11:02
【问题描述】:

在编写股票市场交易员IObserver 的过程中,我遇到了三个错误,主要来自Reactive Extensions 库。

我有以下CompanyInfo 类:

public class CompanyInfo
{
    public string Name { get; set; }

    public double Value { get; set; }
}

还有一个叫做StockMarketIObservable<CompanyInfo>

public class StockMarket : IObservable<CompanyInfo>

我的Observer 如下所示:

public class StockTrader : IObserver<CompanyInfo>
{
    public void OnCompleted()
    {
        Console.WriteLine("Market Closed");
    }

    public void OnError(Exception error)
    {
        Console.WriteLine(error);
    }

    public void OnNext(CompanyInfo value)
    {
        WriteStock(value);
    }

    private void WriteStock(CompanyInfo value) { ... }
}

我运行以下代码:

StockMarket market = GetStockMarket();
StockTrader trader = new StockTrader();

IObservable<CompanyInfo> differential = market  //[F, 1], [S, 5], [S, 4], [F, 2]
    .GroupBy(x => x.Name)                       //[F, 1], [F, 2]; [S, 5], [S, 4]
    .SelectMany(x => x                  //4, 8, 2, 3
        .Buffer(2, 1)                   //(4, 8), (8, 2), (2, 3), (3)
        .SkipLast(1)                    //(4, 8), (8, 2), (2, 3)
        .Select(y => new CompanyInfo    //(+100%), (-75%), (+50%)
        {
            Name = x.Key,
            Value = (y[1].Value - y[0].Value) / y[0].Value
        })                                      //[F, +100%]; [S, -20%]
    );

using (IDisposable subscription = differential.Subscribe(trader))
{
    Observable.Wait(market);
}

出现三个错误之一:

  • Reactive Extensions 中抛出以下ArgumentException

    System.ArgumentException:已添加具有相同键的项目。 在 System.ThrowHelper.ThrowArgumentException(ExceptionResource 资源) 在 System.Collections.Generic.Dictionary`2.Insert(TKey key, TValue value, Boolean add) 在 System.Reactive.Linq.Observable.GroupBy'3._.OnNext(TSource 值)

  • 以下IndexOutOfRangeException

    参数名称:索引 在 System.ThrowHelper.ThrowArgumentOutOfRangeException(ExceptionArgument 参数,ExceptionResource 资源) 在 System.Collections.Generic.List'1.get_Item(Int32 索引) 在 StockMarketTests.c__DisplayClass0_0.b__2(IList'1 y) 在 System.Reactive.Linq.Observable.Select'2._.OnNext(TSource 值)

  • Console 的文字偶尔会发生调整(颜色应该一致):

什么会导致这些奇怪的症状?

【问题讨论】:

    标签: c# .net linq system.reactive observable


    【解决方案1】:

    Reactive Extensions 概念的最大优点之一是能够订阅“某处”发生的“发生” (IObservable) 和在这个'发生'上应用面向对象的概念——这不需要知道那个'somewhere'在哪里。

    这种方式Reactive Extensions 大大简化了面向eventproducer-consumer problems 的编程

    订阅IObservable 而不知道观察到的数据来源的能力迫使订阅者假设通知是不可预测的。换句话说,当观察IObservable 你应该假设通知可以同时发送

    由于Reactive Externsions 的行为契约,IObservables 应该一次生产一个项目。通常情况下会发生这种情况,但有时外部实现不遵循该约定。

    让我们来看看这三个问题中的每一个:

    GroupBy 不是线程安全的


    GroupBy 通过返回一个IObservable&lt;IGroupedObservable&lt;T&gt;&gt; 工作,它的OnNext 方法调用外部IObservableOnNext 与匹配当前通知的IGroupedObservable&lt;T&gt;。它通过为Dictionary 中的每个键保留一个IGroupedObservable&lt;T&gt; (更准确地说是一个Subject&lt;T&gt; 来实现这一点 - 这并不奇怪 - 不是ConcurrentDictionary em>。这意味着两个邻近的通知可能会导致重复插入

    Select 并不孤单


    Select 的线程安全由其提供的委托决定。在上面的例子中,提供给Select 方法的委托依赖于Buffer(2, 1) 将提供一个大小为2 的列表。Buffer 包含一个Queue,它不是并发,因此当从多个线程迭代时 - BufferQueue 可以为我们提供一些意想不到的结果

    另一个Exception 可能出于同样的原因抛出,如果y 将被提供null,则NullReferenceException,或者QueueInvalidOperationException 可以在迭代时被修改。

    即使是基本的观察也不安全


    最后但同样重要的是,即使您只进行基本观察,StockTraderOnNext 方法也会在非atomic operation 中修改控制台,这会导致奇怪的文本布局。

    那你能做什么?


    Synchronize 方法的存在使您能够验证您正在订阅 linear IObservable&lt;T&gt;,这意味着 OnNext 方法的调用不超过一次可以同时发生

    由于即使是GroupBy 扩展方法也不是线程安全的,所以Synchronize 方法需要在链的开头调用:

    IObservable<CompanyInfo> differential = market  //[F, 1], [S, 5], [S, 4], [F, 2]
        .Synchronize()
        .GroupBy(x => x.Name)                       //[F, 1], [F, 2]; [S, 5], [S, 4]
        .SelectMany(x => x                  //4, 8, 2, 3
            .Buffer(2, 1)                   //(4, 8), (8, 2), (2, 3), (3)
            .SkipLast(1)                    //(4, 8), (8, 2), (2, 3)
            .Select(y => new CompanyInfo    //(+100%), (-75%), (+50%)
            {
                Name = x.Key,
                Value = (y[1].Value - y[0].Value) / y[0].Value
            })                                      
        );                                          //[F, +100%]; [S, -20%]
    

    请注意,Synchronize 为您的查询添加了另一个代理 Observable,因此它会使查询慢一点,因此 您应该避免在不需要时使用它

    【讨论】:

    • 这就像你把 SO 变成了你自己的个人博客文章。
    • 哈哈,帮了遇到这个问题的人,想到这里分享一下解决方法……
    • -1 我不同意这个建议。 @Enigmativity 很赚钱。这篇长文掩盖了问题的核心。 IObservable&lt;T&gt; 合同不能完全由 .NET 强制执行。它包括对OnXXX 调用进行序列化的要求。结束。如果您收到未序列化的调用,则这是 observable 中的错误。现在要么你拥有或信任可观察的,要么你不拥有,然后你就可以适当地编码。假设所有的 observable 都充满了 bug 并不是一个好的建议——最终的游戏是在每个 observable 之后放置 Synchronize。代码膨胀和性能下降。
    • ...您需要务实并意识到这种可能性,并做出适当的决定。这一切对我来说都不像新闻快讯。由于性能影响,Rx 的实现者做出了故意决定不强制执行或检查序列化。真正的危险是,这个建议可能会导致某人走上错误的道路,并导致他们编写一堆不必要的过于谨慎的代码。希望这些 cmets 能提供一些平衡。
    • 这个答案对于 Rx 是否适合完成其工作非常具有误导性。 Rx 的操作符是我所见过的一样可靠的代码。它们非常健壮,并且根据 Rx 合约牢固地构建。只有当源 observable 不遵循合同时才会出现问题。代码的问题在于StockMarket的实现。如果它是按照合同编写的,那一切都会好起来的。这应该是重点。
    【解决方案2】:

    您的代码的问题不在于查询,也不在于 Rx 本身。问题可能来自您实际的StockMarketStockTrader 实现。

    现在,问题很可能是因为您正在为您的 market observable 创建两个订阅。

    当你写这个时:

    using (IDisposable subscription = differential.Subscribe(trader))
    {
        Observable.Wait(market);
    }
    

    ...您将获得两个market 的订阅。一个在differential.Subscribe(trader),另一个因为Observable.Wait(market);

    我怀疑两个并发订阅导致了您的问题,但没有看到 StockMarket 的实现,我们无法判断它为什么会抛出。

    这是实现您自己的 observable 和observer 实现的危险。你应该避免这样做。最好有一个属性 IObservable&lt;CompanyInfo&gt; CompanyValues { get; } 悬挂在使用标准 Rx 运算符构建的 CompanyInfo 之外。

    而且您应该始终避免像.Wait(...) 这样的阻塞操作。

    作为一项快速测试,我会将您当前的Observable.Wait(market); 替换为Thread.Sleep(?),并具有足够长的睡眠时间,以查看您的代码是否正常运行。当然,您需要确保在后台调度程序上生成值(例如 Scheduler.Default)。

    我运行了这段代码来测试你的查询:

    public class CompanyInfo
    {
        public string Name { get; set; }
    
        public double Value { get; set; }
    }
    
    public class StockTrader : IObserver<CompanyInfo>
    {
        public void OnCompleted()
        {
            Console.WriteLine("Market Closed");
        }
    
        public void OnError(Exception error)
        {
            Console.WriteLine(error);
        }
    
        public void OnNext(CompanyInfo value)
        {
            WriteStock(value);
        }
    
        private void WriteStock(CompanyInfo value) { Console.WriteLine($"{value.Name} = {value.Value}"); }
    }
    
    public class StockMarket : IObservable<CompanyInfo>
    {
        private CompanyInfo[] _values = new CompanyInfo[]
        {
            new CompanyInfo() { Name = "F", Value = 1 },
            new CompanyInfo() { Name = "S", Value = 5 },
            new CompanyInfo() { Name = "S", Value = 4 },
            new CompanyInfo() { Name = "F", Value = 2 },
        };
    
        public IDisposable Subscribe(IObserver<CompanyInfo> observable)
        {
            return _values.ToObservable().ObserveOn(Scheduler.Default).Subscribe(observable);
        }
    }
    

    ...用这个:

    StockMarket market = new StockMarket();
    StockTrader trader = new StockTrader();
    
    IObservable<CompanyInfo> differential = market  //[F, 1], [S, 5], [S, 4], [F, 2]
        .GroupBy(x => x.Name)                       //[F, 1], [F, 2]; [S, 5], [S, 4]
        .SelectMany(x => x                  //4, 8, 2, 3
            .Buffer(2, 1)                   //(4, 8), (8, 2), (2, 3), (3)
            .SkipLast(1)                    //(4, 8), (8, 2), (2, 3)
            .Select(y => new CompanyInfo    //(+100%), (-75%), (+50%)
            {
                Name = x.Key,
                Value = (y[1].Value - y[0].Value) / y[0].Value
            })                                      //[F, +100%]; [S, -20%]
        );
    
    IDisposable subscription = differential.Subscribe(trader);
    
    Thread.Sleep(10000);
    

    我没有一次让它崩溃或导致任何异常。

    【讨论】:

    • 我明天会给你两个实现,解释为什么这是错误的,问题是并发和并发本身。您的实现不是从多个线程触发的,因此您没有同样的问题。另外,这是时间问题,如果时间不对,可能会得到不同的结果。
    • 另外,一次订阅也会发生这种情况。
    • @TamirVered - Rx 有一个行为契约,一次产生一个值。如果源数据不遵守合同,则可能需要.Synchronize(),但我已经证明问题不在于查询,而在于源。
    • 你没有表现出来。你的实现是一个无状态的cold Observable,当然它不会有线程问题,你也不会从多个线程推送通知。如果我只从一个线程向我的StockMarket 推送通知,我没有问题。稍后我将分享我的GetStockMarket 方法,以便您看到我所说的内容。
    • @TamirVered - 如果源 observable 没有按照合同行事,您只会遇到线程问题。冷不冷都无所谓。我不知道你说的无国籍是什么意思。如果您的源行为本身,您应该能够毫无问题地从多个线程推送多个值。如果没有,那么根据您的回答,输入.Synchronize() 就可以了。
    猜你喜欢
    • 2012-01-15
    • 2014-01-11
    • 1970-01-01
    • 2014-10-10
    • 2013-10-10
    • 2016-01-19
    • 1970-01-01
    • 2017-03-23
    相关资源
    最近更新 更多