【问题标题】:Going from Observable to Enumerable/Value without Blocking从 Observable 到 Enumerable/Value 没有阻塞
【发布时间】:2012-12-02 06:31:16
【问题描述】:

我遇到了一些阻塞问题。我正在尝试从推到拉。即我想在这里访问我的数据,在这种情况下,它是通过我的 Observable 处理后的数组。

type HistoryBar = 
    {Open: decimal; High: decimal; Low: decimal; Close: decimal; Time: DateTime; Volume: int; RequestId: int; Index: int; Total: int}

let transformBar =
    client.HistoricalData
    |> Observable.map(fun args ->
        {
            Open = args.Open
            High = args.High
            Low =  args.Low
            Close = args.Close
            Time = args.Date
            Volume = args.Volume
            RequestId = args.RequestId
            Index = args.RecordNumber
            Total = args.RecordTotal
       }
    )

let groupByRequest (obs:IObservable<HistoryBar>) = 
    let bars = obs.GroupByUntil((fun x -> x.RequestId), (fun x -> x.Where(fun y -> y.Index = y.Total - 1)))
    bars.SelectMany(fun (x:IGroupedObservable<int, HistoryBar>) -> x.ToArray())

let obs = transformBar |> groupByRequest

client.RequestHistoricalData(1, sym, DateTime.Now, TimeSpan.FromDays(10.0), BarSize.OneDay, HistoricalDataType.Midpoint, 0)

如果我订阅 obs,那么只要我打电话给client.RequestHistoricalData 一切都会正常。我想做的是将 obs 转换为基础类型,在本例中为HistoryBar []。我试过使用wait、ToEnumberable,但没有成功。提取我最后创建的数据的正确方法是什么?

编辑,添加人为的 C# 示例代码以显示库的正常工作方式。我在这里真正想了解的是如何从可观察到标准列表或数组。我不确定是否需要一个可变结构才能做到这一点。如果我不得不猜测,我会说不。

static void Main(string[] args)
{
...
client.HistoricalData += client_HistoricalData;
client.RequestHistoricalData(1, sym, DateTime.Today, TimeSpan.FromDays(10), BarSize.OneDay, HistoricalDataType.Midpoint, 0);
....
}

static void client_HistoricalData(object sender, HistoricalDataEventArgs e)
{

Console.WriteLine("Open: {0}, High: {1}, Low: {2}, Close: {3}, Date: {4}, RecordId: {5}, RecordIndex: {6}", e.Open, e.High, e.Low, e.Close, e.Date, e.RequestId, e.RecordNumber);
}

【问题讨论】:

  • client.RequestHistoricalData 是做什么的?请说明。在SelectMany 内部调用ToArray 会导致阻塞对整个序列的急切评估。
  • 添加了更多细节。 client.RequestHistoricalData 将进行调用以执行数据获取,这将调用与 client.HistoricalData 关联的事件处理程序
  • C# 中的人为代码仍然是人为的。
  • 不知道为什么ToEnumerable 在这里不起作用?您将必须阻止,因为IEnumerator&lt;T&gt; 无法产生任何东西,直到观察到某些东西向您屈服。

标签: f# system.reactive


【解决方案1】:

这个问题并没有很清楚地首先如何加载数据(无论是惰性/时变等),所以我只是假设它是一个有时间限制的值流。

从您的代码看来,您似乎想在流完成时找到流中的最后一个值。 Last 方法为您提供了流中完成时推送的最后一个值 - 然而,这是同步的并且阻塞直到流完成。非阻塞版本,LastAsync 返回一个 Observable,它会在源代码完成时产生一个值。

let from0To4 = 
    Observable.Interval(TimeSpan.FromSeconds(0.1)).Take(5)

let lastValue =
    from0To4.LastAsync()

let disposable =
    lastValue |> Observable.subscribe(log)

要将 Observable 转换为列表而不部分阻塞,您可以使用 Buffer 方法。要在 Observable 完成之前缓冲所有值,请使用 ToList。

let fullBuffer =
    from0To4.ToList()

let disposable =
    fullBuffer |> Observable.subscribe(fun ls -> printfn "Buffer(%d): %A" ls.Count ls)

输出:

缓冲区(5):seq [0L; 1升; 2升; 3升; ...]

【讨论】:

  • Observable.ToList() 在这里比调用Buffer(Observable.Never()) 更合适,注意它们的工作方式本质上完全相同。 Observable.ToList() 是非阻塞的,并在底层完成时返回单个通知(列表)。
  • 不确定这有什么关系。我要从调用client.RequestHistoricalData 后触发的.NET 事件处理程序client.HistoricalData 开始
  • @davewolfs 您的目标是将事件的所有实例转换为列表吗?
猜你喜欢
  • 2019-07-17
  • 2013-05-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-07-06
  • 2013-08-19
  • 1970-01-01
  • 2015-03-15
相关资源
最近更新 更多