【发布时间】:2018-06-28 08:26:00
【问题描述】:
我们有一种情况,我们希望在 C# 应用程序中启动后台“轮询”操作,该操作使用响应式扩展定期返回值。我们要实现的流程如下:
- 调用者调用类似
Poll()的方法返回IObservable - 调用者订阅了上述 observable,并启动了一个后台线程/任务,该线程/任务与硬件交互以在某个时间间隔内检索值
- 调用者完成后,它会处理订阅并自动停止后台线程/任务
尝试 #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<OutputParameters>必须在有人订阅后立即启动轮询的异步工作。研究如何在 c# 中跟踪订阅者。您可以使用事件或使用观察者模式,我只会使用事件。然后,一旦所有订阅者都取消订阅,发送一个取消令牌以停止轮询。大部分代码应该在Subject<OutputParameters>中,因此它被封装而不是在Program类中。 -
另外,首先使用 Windows 窗体应用程序尝试它,一旦它工作,在控制台应用程序中搜索
Task<T>并研究它,因为使用控制台有点棘手。 -
出于好奇,投反对票有什么原因吗?
-
如果你问我,我自己也在想这个。我认为你的问题措辞得当而且很清楚。
-
@SamStorie - 不要像您的问题那样使用
Subject。要么将其包装在Observable.Defer中,要么尝试完全避免它。将它作为变量放在方法中会导致它在您的 observable 中被捕获,从而创建一个 run-once observable - 并且发送的任何错误或完成信号都将永远杀死您的 observable。
标签: c# task-parallel-library system.reactive