【问题标题】:Observable not reacting to queue changed on different threadObservable 对不同线程上的队列更改没有反应
【发布时间】:2015-01-16 12:16:55
【问题描述】:

我有以下代码:

static void Main()
    {
        var holderQueue = new ConcurrentQueue<int>(GetInitialElements());

        Action<ConcurrentQueue<int>> addToQueueAction = AddToQueue;
        var observableQueue = holderQueue.ToObservable();
        IScheduler newThreadScheduler = new NewThreadScheduler();

        IObservable<Timestamped<int>> myQueueTimestamped = observableQueue.Timestamp();

        var bufferedTimestampedQueue = myQueueTimestamped.Buffer(TimeSpan.FromSeconds(3), TimeSpan.FromSeconds(3), newThreadScheduler);

        var t = new TaskFactory();
        t.StartNew(() => addToQueueAction(holderQueue));

        using(bufferedTimestampedQueue.SubscribeOn(newThreadScheduler).Subscribe(currentQueue =>
        {
            Console.WriteLine("buffer time elapsed, current queue contents is: {0} items.", currentQueue.Count);
            foreach(var item in currentQueue)
                Console.WriteLine("item {0} at {1}", item.Value, item.Timestamp);

            Console.WriteLine("holderqueue has: {0}", currentQueue.Count);
        }))
        {
            Console.WriteLine("started observing queue");

            Console.ReadLine();
        }
    }

    private static void AddToQueue(ConcurrentQueue<int> concurrentQueue)
    {
        while(true)
        {
            var x = new Random().Next(1, 10);
            concurrentQueue.Enqueue(x);
            Console.WriteLine("added {0}", x);
            Console.WriteLine("crtcount is: {0}", concurrentQueue.Count);
            Thread.Sleep(1000);
        }
    }

    private static IEnumerable<int> GetInitialElements()
    {
        var random = new Random();
        var items = new List<int>();
        for (int i = 0; i < 10; i++)
            items.Add(random.Next(1, 10));

        return items;
    }

意图如下:

holderQueue 对象最初填充了一些元素(GetInitialElements),然后在不同的线程上使用更多元素(通过方法AddToQueue)进行更改,并且 observable 应该检测到这种更改并做出反应因此,当它的时间过去时(所以每 3 秒),通过在其订阅中执行该方法。

简而言之,我希望Subscribe 正文中的代码每 3 秒执行一次,并向我显示队列中的更改(在不同的线程上更改)。相反,Subscribe 正文只执行一次。为什么?

谢谢

【问题讨论】:

    标签: multithreading c#-4.0 queue system.reactive observer-pattern


    【解决方案1】:

    ToObservable 方法采用 IEnumerable&lt;T&gt; 并将其转换为可观察对象。结果,它将获取您的并发队列并立即枚举它,遍历所有可用项目。您稍后修改队列以添加其他项目这一事实对从并发队列的 GetEnumerator() 实现返回的已枚举 IEnumerable&lt;T&gt; 没有影响。

    【讨论】:

    • 我明白。在这种情况下,队列本身就不能被“观察”吗?
    • 否;但是,听起来您真正想要的是Subject。您可以在主题上调用 OnNext 以“入队”项目,这些项目将由 Rx 管道接收,然后按照您的需要进行缓冲。如果您希望能够在连接管道之前将项目放到主题上,请使用ReplaySubject,它将项目记录到内部缓冲区中,然后在它们连接时将它们重放给观察者。
    • 或者,如果您没有选择只能使用ConcurrentQueue,您可以编写自定义代码来轮询队列尾部的新项目并将它们推送到可观察的管道。这将非常效率低下,但不幸的是队列没有公开任何挂钩来通知代码何时添加新项目。
    • 使用BufferBlock。构造它后,使用AsObservable 获取可用于观察值的可观察对象。生产者可以Post新项目到缓冲区。缓冲区就像一个 FIFO 队列。
    • @DavidPfeffer 同样重要的是要注意,即使 OP 使用 Subject 或 ReplaySubject,或自定义 ObservableQueue&lt;T&gt; 实现,通知也必须按照 §4.2 合同进行序列化,@ 987654322@。这意味着要将ConcurrentQueue&lt;T&gt; 语义转换为可与Rx 运算符一起使用的IObservable&lt;T&gt;,OP 还必须确保防止对OnNext 的重叠调用。在某些情况下,在转换为 observable 之后应用 Rx 的 Synchronized 运算符会更容易。
    【解决方案2】:

    根据 David Pfeffer 的回答,仅使用 .ToObserverable() 无法满足您的需求。

    但是,当我查看您的代码时,我看到了几件事:

    1. 您正在使用NewThreadScheduler
    2. 您正在通过任务添加到队列中
    3. 您正在使用ConcurrentQueue&lt;T&gt;

    我认为只要改变一些事情,你就可以实现你在此设定的目标。首先,我认为您实际上是在寻找BlockingCollection&lt;T&gt;。我知道这似乎不太可能,但你可以让它像线程安全队列一样工作。

    接下来,您已经指定了一个线程来处理NewThreadScheduler,为什么不让它进行轮询/从队列中拉取?

    最后,如果您使用BlockingCollection&lt;T&gt;.GetConsumingEnumerable(CancellationToken) 方法,您实际上可以返回并使用.ToObservable() 方法!

    那么让我们看看重写后的代码:

    static void Main()
    {
        //The processing thread. I try to set the the thread name as these tend to be long lived. This helps logs and debugging.
        IScheduler newThreadScheduler = new NewThreadScheduler(ts=>{
            var t =  new Thread(ts);
            t.Name = "QueueReader";
            t.IsBackground = true;
            return t;
        });
    
        //Provide the ability to cancel our work
        var cts = new CancellationTokenSource();
    
        //Use a BlockingCollection<T> instead of a ConcurrentQueue<T>
        var holderQueue = new BlockingCollection<int>();
        foreach (var element in GetInitialElements())
        {
            holderQueue.Add(element);
        }
    
        //The Action that periodically adds items to the queue. Now has cancellation support
        Action<BlockingCollection<int>,CancellationToken> addToQueueAction = AddToQueue;
        var tf = new TaskFactory();
        tf.StartNew(() => addToQueueAction(holderQueue, cts.Token));
    
        //Get a consuming enumerable. MoveNext on this will remove the item from the BlockingCollection<T> effectively making it a queue. 
        //  Calling MoveNext on an empty queue will block until cancelled or an item is added.
        var consumingEnumerable = holderQueue.GetConsumingEnumerable(cts.Token);
    
        //Now we can make this Observable, as the underlying IEnumerbale<T> is a blocking consumer.
        //  Run on the QueueReader/newThreadScheduler thread.
        //  Use CancelationToken instead of IDisposable for single method of cancellation.
        consumingEnumerable.ToObservable(newThreadScheduler)
            .Timestamp()
            .Buffer(TimeSpan.FromSeconds(3), TimeSpan.FromSeconds(3), newThreadScheduler)
            .Subscribe(buffer =>
                {
                    Console.WriteLine("buffer time elapsed, current queue contents is: {0} items.", buffer.Count);
                    foreach(var item in buffer)
                        Console.WriteLine("item {0} at {1}", item.Value, item.Timestamp);
    
                    Console.WriteLine("holderqueue has: {0}", holderQueue.Count);
                },
                cts.Token);
    
    
        Console.WriteLine("started observing queue");
    
        //Run until [Enter] is pressed by user.
        Console.ReadLine();
    
        //Cancel the production of values, the wait on the consuming enumerable and the subscription.
        cts.Cancel();
        Console.WriteLine("Cancelled");
    }
    
    private static void AddToQueue(BlockingCollection<int> input, CancellationToken cancellationToken)
    {
        while(!cancellationToken.IsCancellationRequested)
        {
            var x = new Random().Next(1, 10);
            input.Add(x);
            Console.WriteLine("added '{0}'. Count={1}", x, input.Count);
            Thread.Sleep(1000);
        }
    }
    
    private static IEnumerable<int> GetInitialElements()
    {
        var random = new Random();
        var items = new List<int>();
        for (int i = 0; i < 10; i++)
            items.Add(random.Next(1, 10));
    
        return items;
    }
    

    现在我想你会得到你期望的结果:

    added '9'. Count=11
    started observing queue
    added '4'. Count=1
    added '8'. Count=1
    added '3'. Count=1
    buffer time elapsed, current queue contents is: 14 items.
    item 9 at 25/01/2015 22:25:35 +00:00
    item 5 at 25/01/2015 22:25:35 +00:00
    item 5 at 25/01/2015 22:25:35 +00:00
    item 9 at 25/01/2015 22:25:35 +00:00
    item 7 at 25/01/2015 22:25:35 +00:00
    item 6 at 25/01/2015 22:25:35 +00:00
    item 2 at 25/01/2015 22:25:35 +00:00
    item 2 at 25/01/2015 22:25:35 +00:00
    item 9 at 25/01/2015 22:25:35 +00:00
    item 3 at 25/01/2015 22:25:35 +00:00
    item 9 at 25/01/2015 22:25:35 +00:00
    item 4 at 25/01/2015 22:25:36 +00:00
    item 8 at 25/01/2015 22:25:37 +00:00
    item 3 at 25/01/2015 22:25:38 +00:00
    holderqueue has: 0
    added '7'. Count=1
    added '2'. Count=1
    added '5'. Count=1
    buffer time elapsed, current queue contents is: 3 items.
    item 7 at 25/01/2015 22:25:39 +00:00
    item 2 at 25/01/2015 22:25:40 +00:00
    item 5 at 25/01/2015 22:25:41 +00:00
    holderqueue has: 0
    Cancelled
    

    【讨论】:

      猜你喜欢
      • 2015-03-15
      • 1970-01-01
      • 1970-01-01
      • 2022-07-06
      • 2020-08-30
      • 1970-01-01
      • 1970-01-01
      • 2015-10-10
      • 1970-01-01
      相关资源
      最近更新 更多