【问题标题】:Create Rx.NET Observable in SignalR hub在 SignalR hub 中创建 Rx.NET Observable
【发布时间】:2014-04-18 16:59:50
【问题描述】:

我有一个 SignalR 集线器,它侦听客户端请求并使用 Rx.NET 观察数据库表,以便在更新可用时立即将更新发送回请求它们的客户端。但是看起来,根据客户端请求,在集线器中创建的 Observer 实例在方法调用完成后立即被销毁(由 GC?);因此,我没有收到任何更新。

这是我当前的Hub 实现:

public class BookHub : Hub {
    private readonly BookService _service = new BookService();

    public void RequestBookUpdate(string author) {
        BookObserver observer = new BookObserver(Context.connectionId);
        IDisposable unsubscriber = _service.RequestBookUpdate(author, observer);
    }
}

BookService 返回转换为 Observable 的 LINQ 查询:

public IDisposable RequestBookUpdate(string author, BookObserver observer) {
    var query = from b in db.Book where b.Author.Contains(author) select b;
    IObservable<Book> observable = query.ToObservable();
    IDisposable unsubscriber = observable.Subscribe(observer);
    return unsubscriber;
}

BookObserver 只是将新检索到的项目发送回请求更新的特定客户端(由connectionId 标识):

// omissis

private static readonly IHubContext _context = GlobalHost.ConnectionManager.GetHubContext<BookHub>();
private readonly string _connectionId;

public BookObserver(string connectionId) {
    connectionId = _connectionId:
}

public void OnNext(Book value) {
    _context.Clients.Client(_connectionId).foundNewBook(value);
}

我不关心BookService 实例被销毁,但我希望BookObserver 保持活动状态,因此我只能在客户端断开连接时调用unsubscriber.Dispose()。这可能吗?

【问题讨论】:

  • 了解更多背景信息会很有见地。出于什么目的,您希望在连接结束之前保持 Observer 处于未处理状态?
  • 我想收听同一个 SQL 查询,并在提交一些新数据后立即获取更新。由于数据从单独的进程(我无法控制)提交到数据库中,唯一的解决方案是轮询(数据库和 Web 服务)或使用 Rx.NET 观察IQueryable并通过 WebSocket 发送更新。
  • 在 IQueryable 中创建查询并不能保证您将收到所有未来的元素。为此,您需要一个能够提供未来元素的查询提供程序。 db 似乎没有这个功能。
  • 移除元素目前不是我们的要求之一。
  • 您似乎认为.ToObservable() 会将单个数据库查询变成重复查询数据库以进行更改的查询。它没有。您需要根据Observable.Interval 创建一个查询,并自己进行重复的数据库调用以获取新记录,然后返回它们。您的 observables 签名最终应该类似于 IObservable&lt;IEnumerable&lt;Book&gt;&gt;,因为每次调用数据库都会返回零个或多个书籍。

标签: c# .net signalr system.reactive signalr-hub


【解决方案1】:

Observer 在收到对OnComplete 的调用时会被自动释放。这实际上是一个非常好的模式,因为这意味着您不必像这样手动处理Subscriptions:

Observable.Range(0, 100)
    .Subscribe(...);

Observable.Interval(TimeSpan.FromSeconds(1))
    .Take(10)
    .Subscribe();

因此,为了确保您的观察者在想要被释放之前不会被释放,您可以连接另一个空的、永远不会完成的、可观察到您的源。

IObservable<Book> observable = query.ToObservable()
    .Concat(Observable.Never<Book>());

但是,根据您尝试执行的操作,最好在其他地方(例如客户端)处理此问题。

【讨论】:

  • 关键是要进行类似连续查询(在这种情况下是可观察的),一旦某些数据发生变化就通知观察者。这样,观察者可以通过 SignalR 通知和更新客户端。如果仅返回“常规查询”数据时自动处置观察者,则 RX(和 WebSockets)的功能将丢失;你不同意吗?
  • 如果 Observable 正在完成,这意味着您的查询没有更多可用的元素。问题不在于 Observable 或 Observer,而在于生产者。在这种情况下,生产者是数据库,数据库只为您提供现有元素,而不是未来的元素。
  • 你确定吗?因为从常规 Rest 服务执行相同的方法并在数据库中添加一行会产生 OnNext 调用。
  • 你能在问题中提供一个比较的Rest代码示例吗?
  • 对同一个 BookService 的调用完全相同,但它是由 ApiController 生成的。当然,这只是为了测试,因为新结果无法推送给客户 :-)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多