【问题标题】:How to start a cold observable using Task.Run() from the task parallel library?如何使用任务并行库中的 Task.Run() 启动冷可观察对象?
【发布时间】:2018-06-28 08:26:00
【问题描述】:

我们有一种情况,我们希望在 C# 应用程序中启动后台“轮询”操作,该操作使用响应式扩展定期返回值。我们要实现的流程如下:

  1. 调用者调用类似Poll() 的方法返回IObservable
  2. 调用者订阅了上述 observable,并启动了一个后台线程/任务,该线程/任务与硬件交互以在某个时间间隔内检索值
  3. 调用者完​​成后,它会处理订阅并自动停止后台线程/任务

尝试 #1

为了证明这一点,我编写了以下控制台应用程序,但这并不是我所期望的:

public class OutputParameters
{
    public Guid Id { get; set; }
    public int Value { get; set; }
}

public class Program
{
    static void Main(string[] args)
    {
        Console.WriteLine("Requesting the polling operation");
        var worker1 = Poll();

        Console.WriteLine("Subscribing to start the polling operation");

        var sub1 = worker1.Subscribe(
            value => { Console.WriteLine($"Thread {value.Id} emitted {value.Value}"); },
            ex => { Console.WriteLine($"Thread threw an exception: {ex.Message}"); },
            () => { Console.WriteLine("Thread has completed"); });


        Thread.Sleep(5000);

        sub1.Dispose();

        Console.ReadLine();
    }


    private static IObservable<OutputParameters> Poll()
    {
        return Observable.DeferAsync(Worker);
    }


    private static Task<IObservable<OutputParameters>> Worker(CancellationToken token)
    {
        var subject = new Subject<OutputParameters>();

        Task.Run(async () =>
        {
            var id = Guid.NewGuid();
            const int steps = 10;

            try
            {
                for (var i = 1; i <= steps || token.IsCancellationRequested; i++)
                {
                    Console.WriteLine($"[IN THREAD] Thread {id}: Step {i} of {steps}");
                    subject.OnNext(new OutputParameters { Id = id, Value = i });

                    // This will actually throw an exception if it's the active call when
                    //  the token is cancelled.
                    //
                    await Task.Delay(1000, token);
                }
            }
            catch (TaskCanceledException ex)
            {
                // Interestingly, if this is triggered because the caller unsibscribed then
                //  this is unneeded...the caller isn't listening for this error anymore
                //
                subject.OnError(ex);
            }

            if (token.IsCancellationRequested)
            {
                Console.WriteLine($"[IN THREAD] Thread {id} was cancelled");
            }
            else
            {
                Console.WriteLine($"[IN THREAD] Thread {id} exiting normally");
                subject.OnCompleted();
            }
        }, token);

        return Task.FromResult(subject.AsObservable());
    }
}

上面的代码实际上似乎几乎立即取消了后台任务,因为这是输出:

Requesting the polling operation
Subscribing to start the polling operation
[IN THREAD] Thread a470e6f4-2e62-4a3c-abe6-670bce39b6de: Step 1 of 10
Thread threw an exception: A task was canceled.
[IN THREAD] Thread a470e6f4-2e62-4a3c-abe6-670bce39b6de was cancelled

尝试 #2

然后我尝试对Worker 方法进行小幅更改,使其异步并等待Task.Run 调用,如下所示:

    private static async Task<IObservable<OutputParameters>> Worker(CancellationToken token)
    {
        var subject = new Subject<OutputParameters>();

        await Task.Run(async () =>
        {
            ...what happens in here is unchanged...
        }, token);

        return subject.AsObservable();
    }

这里的结果使后台任务看起来像是完全控制,因为它在被取消之前确实运行了大约 5 秒,但是订阅回调没有输出。这是完整的输出:

Requesting the polling operation
Subscribing to start the polling operation
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a: Step 1 of 10
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a: Step 2 of 10
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a: Step 3 of 10
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a: Step 4 of 10
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a: Step 5 of 10
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a: Step 6 of 10
[IN THREAD] Thread cf416d81-c3b7-41fe-8f5a-681da368452a was cancelled

我的问题

所以很明显,我不完全理解这里发生了什么,或者在这种情况下使用 DeferAsync 是可观察对象的正确创建方法。

有没有合适的方法来实现这种方法?

【问题讨论】:

  • Subject&lt;OutputParameters&gt; 必须在有人订阅后立即启动轮询的异步工作。研究如何在 c# 中跟踪订阅者。您可以使用事件或使用观察者模式,我只会使用事件。然后,一旦所有订阅者都取消订阅,发送一个取消令牌以停止轮询。大部分代码应该在Subject&lt;OutputParameters&gt; 中,因此它被封装而不是在Program 类中。
  • 另外,首先使用 Windows 窗体应用程序尝试它,一旦它工作,在控制台应用程序中搜索 Task&lt;T&gt; 并研究它,因为使用控制台有点棘手。
  • 出于好奇,投反对票有什么原因吗?
  • 如果你问我,我自己也在想这个。我认为你的问题措辞得当而且很清楚。
  • @SamStorie - 不要像您的问题那样使用Subject。要么将其包装在 Observable.Defer 中,要么尝试完全避免它。将它作为变量放在方法中会导致它在您的 observable 中被捕获,从而创建一个 run-once observable - 并且发送的任何错误或完成信号都将永远杀死您的 observable。

标签: c# task-parallel-library system.reactive


【解决方案1】:

如果只有 RX 的解决方案就足够了,这就可以了。如果你问我,那就更清楚了......

static IObservable<OutputParameters> Poll()
{
    const int steps = 10;
    return Observable.Defer<Guid>(() => Observable.Return(Guid.NewGuid()))
        .SelectMany(id => 
            Observable.Generate(1, i => i <= steps, i => i + 1, i => i, _ => TimeSpan.FromMilliseconds(1000))
                .ObserveOn(new EventLoopScheduler())
                .Do(i => Console.WriteLine($"[IN THREAD] Thread {id}: Step {i} of {steps}"))
                .Select(i => new OutputParameters { Id = id, Value = i })
        );
}

解释:

  • Generate 就像 Rx 的 for 循环。最后一个参数控制何时发射项目。这相当于你的 for 循环 + Task.Delay
  • ObserveOn 控制观察到 observable 的位置/时间。在这种情况下,EventLoopScheduler 将为每个订阅者启动一个新线程,并且该 observable 中的所有项目都将在新线程上观察到。

来自谜团:

static IObservable<OutputParameters> Poll()
{
    const int steps = 10;
    return Observable.Defer<OutputParameters>(() =>
    {
        var id = Guid.NewGuid();
        return Observable.Generate(1, i => i <= steps, i => i + 1, i => i,
                _ => TimeSpan.FromMilliseconds(1000), new EventLoopScheduler())
            .Do(i => Console.WriteLine($"[IN THREAD] Thread {id}: Step {i} of {steps}"))
            .Select(i => new OutputParameters { Id = id, Value = i });
    });
}

【讨论】:

  • 您应该将 observable 包裹在 Observable.Defer 中以启用对 id 的捕获,否则多个订阅者将获得相同的 id
  • @SamStorie - EventLoopScheduler 确保每个 observable 和安排在该 observable 上的每个操作都使用一个且只有一个线程,然后确保没有执行并发访问 - 只有一个线程所以两件事不能同时运行。这是允许多线程代码针对非线程安全代码运行的一种非常巧妙的方法。
  • @Enigmativity,谢谢。更新了答案以包括 Defer 中的 Guid 初始化,以及 EventLoopScheduler 优先于 NewThreadScheduler
  • @Shlomo - 很抱歉编辑你的答案。我以为我会将您的答案变体作为对您答案的编辑,而不是发布我自己的答案并使一切变得混乱。请随意使用或删除它。我发布它的原因是Return/SelectMany/ObserveOn 组合将在编组到EventLoopScheduler 之前首先使用线程池中的线程。我写它的方式将避免所有这些并坚持单线程。
  • 再次感谢你们的帮助。我们已将其放入我们的实际应用程序中,并且该方法效果很好。我做的一个小改变是让我通过调度程序,这样我就可以使用 TestScheduler 更轻松地进行测试......这些答案也帮助我探索。赞一个!
猜你喜欢
  • 1970-01-01
  • 2017-09-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多