【发布时间】: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代码开始一个新项目,然后只包含编译所需的NewsService和Database的部分。它真的不必是可运行的,只是可编译的。那么我很可能可以重构代码。 -
我添加了指向源代码的链接。 original_source 分支有一个编译和工作代码。另外,我根据您的建议做了一些改进;他们在主分支。
-
@Enigmativity 只是想知道您是否有机会查看github.com/harvinders/RxTest/tree/original_issue 的代码或github.com/harvinders/RxTest/tree/main 的改进版本
标签: system.reactive