【问题标题】:Creating a REST client API using Reactive Extensions (Rx)使用响应式扩展 (Rx) 创建 REST 客户端 API
【发布时间】:2010-05-16 23:00:40
【问题描述】:

我正在努力了解反应式扩展 (Rx) 的正确用例。不断出现的示例是 UI 事件(拖放、绘图),以及 Rx 适合异步应用程序/操作(如 Web 服务调用)的建议。

我正在开发一个需要为 REST 服务编写小型客户端 API 的应用程序。我需要调用四个 REST 端点,其中三个是为了获取一些参考数据(机场、航空公司和状态),第四个是主要服务,可为您提供给定机场的航班时间。

我创建了公开三个参考数据服务的类,方法如下所示:

public Observable<IEnumerable<Airport>> GetAirports()
public Observable<IEnumerable<Airline>> GetAirlines()
public Observable<IEnumerable<Status>> GetStatuses()
public Observable<IEnumerable<Flights>> GetFlights(string airport)

在我的 GetFlights 方法中,我希望每个航班都保存一个参考,它是从机场出发的,以及运营该航班的航空公司。为此,我需要 GetAirports 和 GetAirlines 提供的数据。每个机场、航空公司和状态都将被添加到一个字典(即字典)中,这样我就可以在解析每个航班时轻松设置参考。

flight.Airport = _airports[flightNode.Attribute("airport").Value]
flight.Airline = _airlines[flightNode.Attribute("airline").Value]
flight.Status = _statuses[flightNode.Attribute("status").Value]

我当前的实现现在看起来像这样:

public IObservable<IEnumerable<Flight>> GetFlightsFrom(Airport fromAirport)
{
    var airports = new AirportNamesService().GetAirports();
    var airlines = new AirlineNamesService().GetAirlines();
    var statuses = new StatusService().GetStautses();


    var referenceData = airports
        .ForkJoin(airlines, (allAirports, allAirlines) =>
                            {
                                Airports.AddRange(allAirports);
                                Airlines.AddRange(allAirlines);
                                return new Unit();
                            })
        .ForkJoin(statuses, (nothing, allStatuses) =>
                            {
                                Statuses.AddRange(allStatuses);
                                return new Unit();
                            });

    string url = string.Format(_serviceUrl, 1, 7, fromAirport.Code);

    var flights = from data in referenceData
                    from flight in GetFlightsFrom(url)
                    select flight;

    return flights;
}

private IObservable<IEnumerable<Flight>> GetFlightsFrom(string url)
{
    return WebRequestFactory.GetData(new Uri(url), ParseFlightsXml);
}

当前的实现基于 Sergey 的回答,并使用 ForkJoin 来确保顺序执行,并且我引用的数据在航班之前加载。这个实现比我之前的实现必须触发“ReferenceDataLoaded”事件更优雅。

【问题讨论】:

标签: .net system.reactive


【解决方案1】:

我认为,如果您从每个 REST 调用中接收到实体列表,那么您的调用应该有一点不同的签名 - 您不是在观察返回集合中的每个值,而是在观察调用完成的事件。所以对于机场来说,它应该有签名:

public IObservable<Aiports> GetAirports()

下一步是并行运行前三个并等待所有这些:

var ports_lines_statuses = 
    Observable.ForkJoin(GetAirports(), GetAirlines(), GetStatuses());

第三步是用 GetFlights() 组合上面的 abservable:

var decoratedFlights = 
  from pls in ports_lines_statuses
  let airport = MyAirportFunc(pls)
  from flight in GetFlights(airport)
  select flight;

编辑:我仍然不明白为什么您的服务返回

IObservable<Airport> 

而不是

IObservable<IEnumerable<Airport>>

AFAIK,通过 REST 调用,您可以一次获得所有实体 - 但也许您会进行分页? 无论如何,如果你想让 RX 做缓冲,你可以使用 .BufferWithCount() :

    var allAirports = new AirportNamesService()
        .GetAirports().BufferWithCount(int.MaxValue); 
...

然后就可以申请ForkJoin了:

var ports_lines_statuses =  
    allAirports
        .ForkJoin(allAirlines, PortsLinesSelector)
        .ForkJoin(statuses, ...

ports_lines_statuses 将包含时间线上的单个事件,该事件将包含所有参考数据。

编辑:这是另一个,使用新创建的 ListObservable(仅限最新版本):

allAiports = airports.Start(); 
allAirlines = airlines.Start();
allStatuses = statuses.Start();

...
whenReferenceDataLoaded =
  Observable.Join(airports.WhenCompleted()
                 .And(airlines.WhenCompleted())
                 .And(statuses.WhenCompleted())
                 Then((p, l, s) => new Unit())); 



    public static IObservable<Unit> WhenCompleted<T>(this IObservable<T> source)
    {
        return source
            .Materialize()
            .Where(n => n.Kind == NotificationKind.OnCompleted)
            .Select(_ => new Unit());
    }

【讨论】:

  • 我实际上确实希望首先获得“一批”中的所有航空公司、机场和状态,因为当我获得航班时,我需要存在这三个参考数据集,以便我可以将它们链接到航班。所以我需要把 Airports 变成一个像这样 Dictionart 的字典,这样我就可以做到:flight.Airport = airports[flightXml.AirportCode].
  • 我用新的方法签名更新了这个问题。你在哪里对,我真的想一次得到所有的机场、航空公司和状态。我在 PortLinesSelector 中是否正确是一种结合机场和航空公司的方法,然后我需要第二种方法来将先前的结果与 ew 结果结合起来?我尝试为 Silverlight 3/4 下载最新版本的 RX,但在 Observable 上找不到 Start() 方法(仅以开头)。
  • @jonas-folleso 是的,PortsLinesSelector 类似于 (ports, lines) => new { ports, lines },第二个选择器会将第三个结果附加到这个。这里的想法是尽量保持函数式风格,只通过管道传递数据,而不是使用局部变量。可观察的 Start() 仅在 .Net4 版本上,因此您必须等到它被移植到其他人...
  • 好吧,酷。你肯定给我指出了正确的方向,当我得到一些效果很好的东西时,我会将此标记为已回答。谢谢分配!
【解决方案2】:

这里的用例是基于拉的 - IEnumerable 很好。如果您想说,通知新航班的到达位置,然后在 Observable.Generate 中包装基于拉取的 REST 调用可能会有一些价值。

【讨论】:

  • 所以 Rx 不是在我的场景中构建 REST 客户端的好方法吗?由于这是 WP7,我无法使其同步,所以替代方案是:GetAirlinesAsync,并有一个 GetAirlinesCompleted 事件。然后我必须调用 GetAirlinesAsync、GetAirportsAsync 和 GetStatusesAsync,并等待所有三个回调事件触发,然后再调用 GetFlights..?我还计划扩展我的方法,使其每 3 分钟重新调用一次 GetFlights 服务以刷新。因此,在到达时观察新的飞行物体听起来是个好主意..?
  • 如果底层 API 只是基于异步的,那么 RX 就更有意义了。 Observable.Generate...
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-03-25
  • 2020-12-20
  • 1970-01-01
  • 1970-01-01
  • 2012-10-02
  • 1970-01-01
相关资源
最近更新 更多