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