【问题标题】:Cancel RX.Net Observer's ongoing OnNext methods取消 RX.Net Observer 正在进行的 OnNext 方法
【发布时间】:2014-11-03 14:43:55
【问题描述】:

正如我最初的问题中所述(请参阅Correlate interdependent Event Streams with RX.Net),我有一个 RX.net 事件流,只要不触发某个其他事件(基本上是“将更改-* 事件处理为只要系统已连接,则在断开连接时暂停,并在系统重新连接后重新开始处理 Change-* 事件。

但是,虽然这适用于新事件,但我将如何取消/向正在进行的 .OnNext() 调用发出取消信号?

【问题讨论】:

  • 为了澄清,处理方法是 async/await-able 并且已经将 CancellationToken 作为参数。也许可以将 RX 和这个 async 方法结合起来。
  • 你能发布一些代码吗?你现在为CancellationToken 传递什么?
  • @Brandon 将在今天晚些时候进行,目前正在运行中,该项目在我家的机器上。现在我只是在玩时提交一个 default(CancellationToken) / CancellationToken.None,但不相信使用和外部 CancellationTokenSource 字段(我将在原始 Stream 中触发并随后重新创建)。

标签: c# events system.reactive reactive-programming cancellation


【解决方案1】:

由于您的观察者已经被写入接受CancellationToken,我们可以修改您的 Rx 流以提供一个连同事件数据。我们将使用 Rx CancellationDisposable,每当取消订阅流时,我们都会处理它。

// Converts into a stream that supplies a `CancellationToken` that will be cancelled when the stream is unsubscribed
public static IObservable<Tuple<CancellationToken, T>> CancelOnUnsubscribe<T>(this IObservable<T> source)
{
    return Observable.Using(
        () => new CancellationDisposable(),
        cts => source.Select(item => Tuple.Create(cts.Token, item)));
}

把这个和另一个问题的解决方案放在一起:

DataSourceLoaded
    .SelectMany(_ => DataSourceFieldChanged
        .Throttle(x)
        .CancelOnUnsubscribe()
        .TakeUntil(DataSourceLoaded))
    .Subscribe(c => handler(c.Item1, c.Item2));

TakeUntil 子句被触发时,它将取消订阅CancelOnUnsubscribe observable,这反过来又会处理CancellationDisposable 并导致令牌被取消。发生这种情况时,您的观察者可以观察此令牌并停止其工作。

【讨论】:

  • 感谢您的解决方案!我发现您对类似问题 (stackoverflow.com/questions/18477018/…) 的回答读起来更优雅一些,并将其组合如下:gist.github.com/jbattermann/914d48f4011f0842a03b
  • 更新了要点,它现在可以完全工作了。感谢 Brandon,也感谢有关 CancellationDisposable 的提示 - 也不知道那个。
  • 啊,是的,这与您在此处发布的问题略有不同。您的要点(和该解决方案)取消当前的工作任务并在新的DataSourceFieldChanged 事件到达时启动一个新任务。听起来这实际上就是您想要的,尽管如此优惠:)
【解决方案2】:

有一个 SelectMany 的异步重载,但诚然,如果存在类似的 Do 重载,它在语义上会更合适。

var subscription = 
  (from _ in DataSourceLoaded
   from __ in DataSourceFieldChanged
     .Throttle(x)
     .SelectMany(DataSourceFieldChangedAsync)
     .TakeUntil(DataSourceUnloaded)
   select Unit.Default);
  .Subscribe();  // Subscribing for side effects only.

...

async Task<Unit> DataSourceFieldChangedAsync(Field value, CancellationToken cancel);

这很好,因为它也将取消与订阅联系在一起。

调用任一

subscription.Dispose()

DataSourceUnloaded.OnNext(x);

将导致CancellationToken 被取消。

【讨论】:

  • 我创建了一个工作项。目前正在讨论替代方法:github.com/Reactive-Extensions/Rx.NET/issues/…
  • 谢谢戴夫。我采用了您和 Brandon 的解决方案(更准确地说,Brandon 对基本相同但早先提出的 SO 问题的回答)并将它们组合成:gist.github.com/jbattermann/914d48f4011f0842a03b
  • 这只是“几乎”正确工作:如果 DataSourceFieldChanged 仍在异步处理并且同时发生 DataSourceUnloaded ,则后者不会触发 WorkerTask(..) 方法的 cancelToken - 对于明显的原因,但还不是完全“在那里”。
  • 好吧,仔细看看 Select/Switch 和 TakeUntil 语句的逻辑安排,当前的 Gist 版本完全按预期工作 \o/
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-11-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-02-13
相关资源
最近更新 更多