【问题标题】:Using Observable.Publish with reactive extensions使用带有响应式扩展的 Observable.Publish
【发布时间】:2012-10-22 06:25:20
【问题描述】:

我对使用 Observable.Publish 进行多播处理的生命周期有点困惑。应该如何正确使用连接?与直觉相反,我发现我不需要为多播观察者调用 connect 来开始他们的订阅。

var multicast = source.Publish();
var field0 = multicast.Select(record => record.field0);
var field1 = multicast.Select(record => record.field1);

// Do I need t*emphasized text*o call here?
var disposable = multicast.connect()

// Does calling 
disposable.Dispose();
// unsubscribe field0 and field1?

编辑

我的困惑是为什么我在不打电话的时候订阅成功 显式连接 IConnectableObservable。但是我在打电话 在隐式调用 Connect 的 IConnectableObservable 上等待

Public Async Function MonitorMeasurements() As Task


    Dim cts = New CancellationTokenSource

    Try
        Using dialog = New TaskDialog(Of Unit)(cts)

            Dim measurementPoints = 
                MeasurementPointObserver(timeout:=TimeSpan.FromSeconds(2)).
                TakeUntil(dialog.CancelObserved).Publish()

            Dim viewModel = New MeasurementViewModel(measurementPoints)
            dialog.Content = New MeasurementControl(viewModel)
            dialog.Show()

            Await measurementPoints
        End Using
    Catch ex As TimeoutException
        MessageBox.Show(ex.Message)
    Catch ex As Exception
        MessageBox.Show(ex.Message)
    End Try

End Function

请注意,我的 TaskDialog 公开了一个名为 CancelObserved 的可观察对象 当按下取消按钮时。

解决方案

@asti 在链接中发布了解决方案。这是该链接中 RX 团队的引述

注意,使用 await 会导致订阅发生,从而使可观察序列变热。此版本中包含对 IConnectableObservable 的 await 支持,这会导致将序列连接到其源并订阅它。如果没有 Connect 调用,等待操作将永远无法完成

【问题讨论】:

  • 我想我们都编辑了链接。哦,好吧!

标签: .net system.reactive


【解决方案1】:

Publish 在源上返回一个 IConnectableObservable<T>,它本质上是带有 Connectmethod 的 IObservable<T>。您可以使用Connect 和它返回的IDisposable 来控制对源的订阅。

Rx 被设计成一个“一劳永逸”的系统。在您明确处置它们或它们完成/错误之前,订阅不会终止。

disp0 = field0.Subscribe(...); disp1 = field1.Subscribe(...) - 订阅不会终止,直到disp0, disp1 被显式处理 - 这与多播源的连接无关。

您可以在不干扰下方管道的情况下连接和断开连接。不用担心手动管理连接的一种更简单的方法是使用.Publish().RefCount(),只要至少有一个观察者仍然订阅它,它就会保持连接。这称为预热 observable。


已更新问题中的编辑

OP 在IConnectableObservable<T> 上呼叫await

来自Release notes for Rx:

..await 的使用通过导致 订阅发生。此版本中包含等待支持 对于 IConnectableObservable,这会导致将序列连接到 它的来源以及订阅它。如果没有 Connect 调用, await 操作永远不会完成。

示例(取自同一页面)

static async  void Foo()
{
    var xs = Observable.Defer(() =>
    {
        Console.WriteLine("Operation started!");
        return Observable.Interval(TimeSpan.FromSeconds(1)).Take(10);
    });

    var ys = xs.Publish();

    // This doesn't trigger a connection with the source yet.
    ys.Subscribe(x => Console.WriteLine("Value = " + x));

    // During the asynchronous sleep, nothing will be printed.
    await Task.Delay(5000);

    // Awaiting causes the connection to be made. Values will be printed now,
    // and the code below will return 9 after 10 seconds.
    var y =  await ys;
    Console.WriteLine("Await result = " + y);
}

【讨论】:

  • 我的困惑是我只调用 .Publish() 而不调用结果上的 connect 并且它仍在工作。
  • 啊哈!!我在 IConnectableObservable 上调用 await 以等待序列完成。我猜那是隐式调用 Connect :)
  • @bradgonesurfing 哇。我没有看到这一点。公告如下:social.msdn.microsoft.com/Forums/en-US/rx/thread/… 见第二个帖子。
  • :) """注意,使用 await 会导致订阅发生,从而使 observable 序列变热。此版本中包含 await 对 IConnectableObservable 的支持,这会导致将序列连接到它的源以及订阅它。没有 Connect 调用,等待操作永远不会完成。"""
  • @bradgonesurfing 回答您的问题 - F# 的计算表达式是对一元结构的表达式重写,而异步是将用户代码重写为状态机 - 它没有那么强大。尝试在 C# 中执行此操作将使您与该语言抗争。
【解决方案2】:

发布允许您共享订阅。这显然对于使 Cold 可观察序列 Hot 最有用。即,采用会导致某些订阅副作用(可能是与网络的连接)发生的序列,并确保副作用执行一次,并在消费者之间共享序列的结果。

实际上,您在冷序列上调用发布,订阅您的消费者,然后在订阅之后连接已发布的序列以缓解任何竞争条件。

所以基本上,你在上面做了什么。

对于已经很热的序列来说,这在很大程度上是没有意义的,例如主题、FromEventPattern 或已发布和连接的内容。

从 Connect() 方法中处理值将“断开”序列,从而阻止消费者获得更多值。如果消费者订阅中的任何一个想要提前分离,您也可以处理这些订阅。

说了这么多,你似乎在做正确的事。你看到的问题是什么?我假设您正在连接到一个已经很热的序列。

【讨论】:

  • 连接不是“热的” 订阅后将建立 UDP 连接,取消订阅后应关闭 UDP 连接。然而,解析 UDP 帧会导致每帧有几个字段。每个字段都需要设置为可观察的,因为它们有不同的客户端观察者。如果我不调用 Publish,那么会创建到同一个套接字的多个连接,这当然是一个错误。调用 Publish 给了我正确的结果,但我似乎不必调用 Connect,这是一个难题。
  • 啊...可能你有一个顽皮的方法?可能是您有一个方法返回一个急切(不是懒惰)的 IObservable。当连接到消息传递层并且您没有将代码包装在 Observable.Create(...) 中时,最常见的情况是
  • 看不到我通过@asti 发布解决方案的问题的更新。在 ICOnnectableObservable 上调用 await 会隐式调用 connect
  • 忽略.. 我知道你现在有了答案。是的,您希望 Await 连接 observable,否则它将如何连接?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-05-17
  • 2012-09-21
  • 1970-01-01
  • 1970-01-01
  • 2014-03-25
相关资源
最近更新 更多