【问题标题】:Rx.Net - Publish method missing first few items when subscribing to Cold ObservableRx.Net - 订阅 Cold Observable 时,发布方法缺少前几项
【发布时间】:2020-10-06 11:39:46
【问题描述】:

Akavache 启发,我正在尝试创建一个为我提供IObservable<IArticle> 的解决方案。该方法基本上首先尝试获取数据库中存在的所有文章,然后尝试从 Web 服务获取更新的文章,并在从 Web 服务获取最新文章时尝试将它们保存回数据库。

由于 web 服务本质上是一个冷的 observable,我不想订阅两次,所以我使用Publish 连接到它。我的理解是我使用的是Publish 方法的正确版本,但是,很多时候该方法往往会错过GetNewsArticles 的前几篇文章。这是通过 UI 以及在下面的调用中添加的 Trace 调用观察到的。

除了解决问题之外,如果还了解如何调试/测试这段代码(除了引入 DI 注入NewsService),那就太好了。

public IObservable<IArticle> GetContents(string newsUrl, IScheduler scheduler)
{
    var newsService = new NewsService(new HttpClient());
    scheduler = scheduler ?? TaskPoolScheduler.Default;

    var fetchObject = newsService
        .GetNewsArticles(newsUrl,scheduler)
        .Do(x => Trace.WriteLine($"Parsing Articles {x.Title}"));

    return fetchObject.Publish(fetchSubject =>
    {
        var updateObs = fetchSubject
            .Do( x =>                         
            {
                // Save to database, all sync calls
            })
            .Where(x => false)
            .Catch(Observable.Empty<Article>());

        var dbArticleObs = Observable.Create<IArticle>(o =>
        {
            return scheduler.ScheduleAsync(async (ctrl, ct) =>
            {
                using (var session = dataBase.GetSession())
                {
                    var articles = await session.GetArticlesAsync(newsUrl, ct);
                    foreach (var article in articles)
                    {
                        o.OnNext(article);
                    }
                }
                o.OnCompleted();
            });
        });

        return
            dbArticleObs                // First get all the articles from dataBase cache
                .Concat(fetchSubject    // Get the latest articles from web service 
                    .Catch(Observable.Empty<Article>())
                    .Merge(updateObs))  // Update the database with latest articles
                .Do(x => Trace.WriteLine($"Displaying {x.Title}"));
    });
}

更新 - 添加了 GetArticles

public IObservable<IContent> GetArticles(string feedUrl, IScheduler scheduler)
{
    return Observable.Create<IContent>(o =>
    {
        scheduler = scheduler ?? DefaultScheduler.Instance;
        scheduler.ScheduleAsync(async (ctrl, ct) =>
        {
            try
            {
                using (var inputStream = await Client.GetStreamAsync(feedUrl))
                {
                    var settings = new XmlReaderSettings
                    {
                        IgnoreComments = true,
                        IgnoreProcessingInstructions = true,
                        IgnoreWhitespace = true,
                        Async = true
                    };

                    //var parsingState = ParsingState.Channel;
                    Article article = null;
                    Feed feed = null;

                    using (var reader = XmlReader.Create(inputStream, settings))
                    {
                        while (await reader.ReadAsync())
                        {
                            ct.ThrowIfCancellationRequested();
                            if (reader.IsStartElement())
                            {
                                switch (reader.LocalName)
                                {
                                    ...
                                    // parsing logic goes here
                                    ...
                                }
                            }
                            else if (reader.LocalName == "item" &&
                                     reader.NodeType == XmlNodeType.EndElement)
                            {
                                o.OnNext(article);
                            }
                        }
                    }

                    o.OnCompleted();
                }
            }
            catch (Exception e)
            {
                o.OnError(e);
            }

        });
        return Disposable.Empty;
    });
}
更新 2

在这里分享source code的链接。

【问题讨论】:

  • 一般来说,当您对Observable.Create 执行return Disposable.Empty; 时,您将创建一个行为不佳的查询。您能否提供一个完整的定义(minimal reproducible example),我看看是否可以在没有Observable.Create 的情况下重构它?
  • 最小的可重现示例意味着引入完整的新闻服务和数据库。我会尝试使用控制台应用程序创建一个,并在下周(希望下周初)将其发布到 github。
  • 我不知道是否需要所有这些内容的完整副本。我将使用完整的GetContents 代码开始一个新项目,然后只包含编译所需的NewsServiceDatabase 的部分。它真的不必是可运行的,只是可编译的。那么我很可能可以重构代码。
  • 我添加了指向源代码的链接。 original_source 分支有一个编译和工作代码。另外,我根据您的建议做了一些改进;他们在主分支。
  • @Enigmativity 只是想知道您是否有机会查看github.com/harvinders/RxTest/tree/original_issue 的代码或github.com/harvinders/RxTest/tree/main 的改进版本

标签: system.reactive


【解决方案1】:

我不喜欢您的代码的一些地方。我假设NewsServiceIDisposable,因为它需要HttpClient(这是一次性的)。您没有进行适当的清理。

此外,您还没有提供完整的方法 - 因为您已尝试将其缩减为问题 - 但这使得很难推断如何重写代码。

也就是说,在我看来非常可怕的一件事是Observable.Create。你能试试这个代码,看看它是否对你有用吗?

    var dbArticleObs =
        Observable
            .Using(
                () => dataBase.GetSession(),
                session =>
                    from articles in Observable.FromAsync(ct => session.GetArticlesAsync(newsUrl, ct))
                    from article in articles
                    select article);

现在,如果是这样,请尝试重写 fetchObject 以在新建 `NewService 时使用相同的 Observable.Using

无论如何,如果您能在问题中提供GetContentsNewsServicedataBase 代码的完整实现,那就太好了。

【讨论】:

  • 是的,您对HttpClient 的看法是完全正确的。我在问题末尾添加了NewsService 的代码。您能否解释一下为什么Observable.Create 的实现是错误的,您的建议是如何解决的?
猜你喜欢
  • 1970-01-01
  • 2019-04-28
  • 1970-01-01
  • 2016-04-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-06-19
  • 2015-04-09
相关资源
最近更新 更多