【问题标题】:Passing in IObservable<T> in method在方法中传入 IObservable<T>
【发布时间】:2015-09-02 18:47:53
【问题描述】:

我正在尝试将 Observable 项目传递到存储库层。

我有界面

public interface Repository{
    IObservable<INotification<bool>> Save(IObservable<T> objects);
}

我在测试这个时遇到了麻烦,并且没有找到任何人做这样的事情的例子。

这种糟糕的设计 Rx 是否明智?我想保存一个结果流,并让存储库根据自己的语义缓冲它们。例如,存储库的实现由 Buffer() 控制。

这样做的部分动机是允许 save 方法在对象流关闭时缓冲/刷新最后一项。它可能会更频繁地这样做,但我不在乎更高级别。

编辑:

我对测试感到非常愚蠢,显然我是在模拟对可观察对象的特定实例的特定调用,如果我使用重放,甚至这样做,实际上是一个对象的新实例,导致我的模拟返回 null .

不过,我仍然对这种模式感到好奇。

【问题讨论】:

  • “我在测试这个时遇到了问题” - 遇到问题是因为/什么问题?
  • 理论上是可以的。在实践中,它可能毫无用处,因为 IObservable 很难用常见的 .NET 应用程序集成工具(WCF、OData 等)进行序列化,除非您打算在一个应用程序中使用此接口
  • 是的,这违背了序列化的目的(我现在想要它,我的应用程序正在关闭!)因为Observerable控制了完成的速度。
  • 问题是如果我事先在 observable 上调用“do”,它会导致 Rx 抛出一个空异常“source is null”。实际调用 save 被嘲笑,并且似乎与可能更早地提取值有关,尽管我似乎无法使用 replay/publish/connect 来解决这个问题。 @paulpdaniels。你有什么建议吗?。我想使用 Observable 因为它提供了一个很好的信号,表明一切都已完成,这样我就不需要给接口一个“保存”方法。所以我可以在 save 的实现中改变缓冲,以及它如何持久化流入的数据流。

标签: c# system.reactive


【解决方案1】:

我认为值得考虑你的界面可能的语义。

public interface Repository
{
    IObservable<INotification<bool>> Save<T>(IObservable<T> objects);
}

可观察的objects 可能是热的或冷的,它可能包含连续值之间的延迟,它可能是无限的。您可能不会设想您的消费者会用“不利”的可观察物来称呼它,但他们可能会。

此外,结果是可观察的。那么,下面的代码应该是什么意思呢?

var results = repository.Save(myObservableObjects);

这可能会触发保存,也可能什么都不做!

您可能必须执行以下操作才能真正触发保存:

 results.Subscribe(...);

如果两个或多个观察者订阅结果 observable 会是什么结果?它会只返回现有结果吗?它会触发源 observable 中对象的新保存吗?

这种存储库的语义实在是太多样化了。

在我看来,这是你需要的界面:

public interface Repository
{
    Task<INotification<bool>> Save<T>(T item);
}

那么你可以这样做:

IObservable<INotification<bool>> results =
    from item in items
    from result in Observable.FromAsync(() => repository.Save(item))
    select result;

这有效地为您提供了您最初想要的完整签名,但您可以完全控制查询的执行,因此您知道保存的内容和时间。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-06-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多