【问题标题】:In Rx.NET, how to enrich values from a (async/observable) lookup在 Rx.NET 中,如何丰富(异步/可观察)查找中的值
【发布时间】:2014-02-06 13:58:49
【问题描述】:

假设我有一个投标流 - 我想用投标人的名字来丰富它:

[
  { bidder: 'user/7', bet: 20  },
  { bidder: 'user/8', bet: 21 }, 
  { bidder: 'user/7', bet: 25  }, 
  /*...., 2 seconds later */
  { bidder: 'user/8', bet: 25  },
  { bidder: 'user/9', bet: 30  },
  ...

投标人名称来自网络服务:

GET '/users?id=7&id=8' =>
[{ user: 'user/7', name: 'Hugo Boss'}, { user: 'user/8', name: "Karl Lagerfeld"}

我将其包装到反应式读取缓存中:

IObservable<User> Users(IObservable<string> userIds);

撰写

现在我想将其组合成以下输出:

[
  { bidder: 'user/7', bet: 20, name: 'Hugo Boss'  },
  { bidder: 'user/8', bet: 21, name: 'Karl Lagerfeld' }, 
  { bidder: 'user/7', bet: 25, name: 'Hugo Boss'  }, 
  /*...., 2 seconds later */
  { bidder: 'user/8', bet: 25, name: 'Karl Lagerfeld'  },
  { bidder: 'user/9', bet: 30, name: 'Somebody else'  },
  ...

步骤 1

我想我需要将出价流投影到用户 ID 流。简单Select。然后在 read-through-cache 中,我用Buffer(TimeSpan, int) 将它分成块。

现在我有一个出价流和一个用户。

第二步

但是现在,如何将两者结合起来呢?

提示:对于 BID,保持订单是有意义的 - 但在我的真实代码中,我并不关心订单。所以我想要一个不依赖于用户缓存以正确顺序返回用户的解决方案。然后Zip 就可以完成这项工作。

我宁愿在获得用户信息后立即将所有出价发布到我的结果流中。

解决方案?

我很确定我需要以某种方式保持一些临时状态(窗口/缓冲区/...)。但我不知道在哪里以及如何。 也许这应该作为自定义运算符来实现;或者可能已经有一个了?

有什么想法吗?

编辑: 似乎不可能在流之上实际编写它。相反,我需要为 userid->user 函数获得一个承诺(Task 或 IObservable),并将其留给承诺批量加载和/或缓存用户,如果合适的话。

【问题讨论】:

  • 为什么Users 方法采用IObservable&lt;string&gt;?为什么不是可枚举的?
  • 因为可枚举不是响应式的 - 它会阻塞等待下一个用户 ID 的线程(可能明天、一秒钟甚至永远不会出现。
  • 对我来说,IObservable in IObservable out 感觉很奇怪。 IMO,拥有 IObserver/IObservable 对感觉更自然。在这种情况下,我可以想象一个方法可能采用单个值,并且缓存被抽象掉。在您的 Users 方法中,OnNexts 请求。缓存实际上是一种优化,它们应该是可选的。
  • 我试图描述的与@bradgonesurfing 给出的答案相似。
  • 您必须简化您的要求。简单的事情是,您向缓存请求single 用户,并且缓存为您提供了一个承诺,Task&lt;User&gt;,它将在未来的某个时间提供该用户。在内部,缓存可以缓冲缓存未命中和对服务器的批量请求。然后,用户将Task&lt;User&gt; 与 Bid 对象配对,并等待任务完成或出错。任务完成后,您现在拥有一对用户和投标。此模式可以使用 RX 或 TPL 建模,但签名 IObservable&lt;User&gt; Users(IObservable&lt;string&gt; userIds); 不起作用

标签: system.reactive


【解决方案1】:

类似下面的东西应该可以工作。将缓存表示为 ISubject 并且您有一个异步缓存。在缓存内部,您可以在一个时间窗口内缓冲查询并将它们批量请求到服务器。

一些虚拟类

public struct User
{
    public string id;
}

public struct Bid
{
    public string userId;
    public int bid;
}

缓存对象本身

/// <summary>
/// Represent the cache by a subject that get's notified of
/// user requests and produces users. Internally this can
/// be buffered and chunked to the server. 
/// </summary>
public static ISubject<string, User> UsersCacheSubject;

分派和异步查询到数据库的方法

public static Task<User> UsersCache(string id)
{
    var r = UsersCacheSubject
        .Where(user=>user.id==id)
        .Replay(1)
        .Take(1);


    UsersCacheSubject.OnNext(id);

    return r.FirstAsync().ToTask();
}

试试看

 public void TryItOut()
{
    IObservable<Bid> bidderObservable = Observable.Repeat(new Bid());
     var foo = from bidder in bidderObservable
              from user in UsersCache(bidder.userId)
              where bidder.userId == user.id
              select new {bidder, user};
 }

【讨论】:

  • IMO UsersCache(string id) 揭示了不应公开的实现细节(主题名称应保持不变)。客户端不需要关心它的获取、计算或缓存。您可能会在 UserRepository 上看到 Get。虽然,这是一个小问题。
  • 没有公开实现细节。调用 UsersCache(string id) 你会得到一个 Task 可以在完成后继续。该任务的完成方式是一个实现细节,但不会向调用者公开。
  • 不,我的意思是,它被缓存的事实。名字本身。缓存实际上是一种优化,它们应该是可选的。我希望以这种方式命名一个具体的类。
  • 好的。只是名字 :) 我可以很容易地改变它。
【解决方案2】:

这不就是像下面这样简单吗?

var query =
    from b in bids
    from u in Users(Observable.Return(b.Bidder))
    select new
    {
        b.Bidder,
        b.Bet,
        u.Name
    };

【讨论】:

  • 用户应该得到一个流 - 然后能够进行优化(块加载,缓存,++)......使用你的代码,每次投注都会调用用户......或者我错过了什么?
  • 实际上这可行,但比它需要的要复杂。用户应该只使用一个 id 并返回一个 Task。如果项目在缓存中,则该任务可以立即完成,如果需要服务器调用,则可以稍后完成。好处是用户的实现可以缓冲对服务器的请求,并且对任务完成的顺序没有要求。我试图在上面的回答中详细说明。
  • @LarsCorneliussen - 我正在遵循 OP 给出的签名。然而,长时间保持可观察的开放可能不是一个好主意。它可能导致内存问题。像任何 I/O 资源一样,它们应该打开,然后在不需要时关闭。
  • @bradgonesurfing - 我正在遵循 OP 给出的签名。我更倾向于拥有一个简单的User GetUser(string userId) 签名,然后将其包装在Observable.StartTask 中。我更喜欢 Rx,因为一旦你进入 Rx,它只会在 99% 的情况下将 TPL 吹走。 :-)
  • 好吧,用户提供的 sig 对于他们想要做的事情是错误的,因为无法将对 id 的请求与稍后提供它的事件相关联。有时它有助于指出这一点。在许多情况下,TPL 和 RX 也是可以互换的。如果您肯定只返回一个值,则有时使用 Task 而不是 IOBservable 会更清楚,但这没什么大不了的。 LINQ 很好地吃掉了两者。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-04-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多