【发布时间】:2020-11-14 11:19:28
【问题描述】:
我正在使用 MassTransit(与 RabbitMQ)在 C# 控制台应用程序(实际上作为 Windows 服务托管,使用 TopShelf)中与 NHibernate 一起使用消息队列。
我们的应用程序包含的流程:
- 监视扫描文档的文件共享(使用 FileSystemWatcher)
- 对文件进行一些处理(将其移动到新的文件位置,将记录插入到我们的数据库中)
- 对文档执行 OCR 并阅读它以查看它是否包含某些单词
- 修改我们的数据库记录以反映第 3 步的结果。
从高级技术实现来看,这项工作的步骤如下:
- 在 FileSystemWatcher 处理其 Created 事件的线程中进行初始处理。完成后,将消息发布到我们的队列以执行 OCR 和单词检查。
- MassTransit 处理消息,创建一个新的生命周期范围,并实例化一个消费者来处理它
- 消费者通过调用 IOCRService 实现来执行 OCR。完成后,在同一个消费者中,我们(从数据库中)获取我们想要搜索的词,然后读取文档文本以找到这些词。
- 将响应消息发送回 MassTransit/RabbitMQ,此消息的使用者会根据是否找到任何单词来修改我们的数据库条目。
我遇到的问题出在上述第 3 步中。这大概是我们旧代码的样子:
public class Consumer
{
public Task Handle(message)
{
_Ocr.DoOcr(message); //Performed in-process
var response = DoDirtyWordCheck(message);
_Publisher.Publish(response);
return Task.CompletedTask;
}
private CheckResponse DoDirtyWordCheck(Message message)
{
var wordsToFind = _DB.FindWords();
var response = _Checker.SearchForWords(message);
}
}
我对这个流程所做的重大改变是,OCR 步骤被移出进程,并放入微服务中,并通过使用 HttpClient 调用微服务来调用。新代码大体相同:
public class Consumer
{
public async Task Handle(message)
{
await _Ocr.DoOcr(message); //calls out to micro-service and awaits the result; method on interface //changed from void return to Task return
var response = DoDirtyWordCheck(message);
_Publisher.Publish(response);
}
private CheckResponse DoDirtyWordCheck(Message message)
{
var wordsToFind = _DB.FindWords(); //Fails here
var response = _Checker.SearchForWords(message);
}
}
然而,我发现,现在调用 _DB.FindWords() 时经常会失败。好吧,事实证明,这个调用发生在与我的生命周期范围开始的线程不同的线程上,这与对await _OCR.DoOCR(); 的调用是同一个线程@ 为什么这是个问题?由于 NHibernate Sessions 不是线程安全的,我们的 DB 层(非常复杂)确保操作只能在创建它的线程上执行。
现在,我以前对 async/await 的理解是,额外的线程不会有任何“技巧”,这样我就不必担心代码必须是线程安全的才能await它.
但是,在深入了解 async/await 并对其工作原理有所了解之后,似乎确实正在在 ThreadPool 线程上完成了一些工作(无论这是实际的等待的工作或等待之后的继续,我仍然不确定),这与控制台应用程序中没有SynchronizationContext这一事实有关,它决定了这个过程是如何完成的;而在 WPF 应用程序中,在 UI 线程上等待的工作必然会在同一个 UI 线程上继续(这是我期望在所有上下文中展示的那种行为)。
所以,这给我带来了我的终极问题:我如何确保在我调用 await 之后需要继续的代码在同一个线程上继续?
我了解我的上述代码流可以通过某种方式进行重组以防止这种情况发生,例如,我可以将 OCR 操作和脏字检查操作分成两个单独的使用者,以便每个操作都处于其不同的上下文/生命周期中范围,可能还有其他任何数量的东西。然而,任何需要这种重组的解决方案对我来说都是有问题的,因为它似乎表明 async/await 是一个比乍一看更容易泄漏的抽象。
看起来这个代码,可能是可以在任何上下文中运行的库代码,应该依赖于与线程模型有关的任何事情,只是因为这个链中的一个调用是现在等待。似乎它应该“正常工作”。在我的应用程序中,与此用例的生命周期和范围以及线程有关的所有内容都预计在消费者级别是原子的,并且我们希望通过 MassTransit 内置的处理来处理(并且已经处理) ,以及我们围绕它的代码结构。
现在,在我看来,我可以设置的SynchronizationContext 中存在一个可能的解决方案(或TaskSheduler?),这是在 WPF 应用程序中处理此类工作的原因,或者WinForms 应用程序,但在控制台应用程序中根本不使用(默认情况下)。不幸的是,在控制台应用程序中做任何事情似乎不再是一个常见的用例,更不用说可能提出这个要求的事情了,所以我真的找不到任何预先存在且经过良好测试的解决方案来做这种排序的事情。此外,由于任何意想不到的影响,我对 async-await 实现的深层内部结构还不够了解,因此我很乐意手动滚动自己的解决方案。
那么,有人可以帮我解决这个问题吗?我的基本假设是否有任何问题让我以错误的方式思考这个问题?还是程序本身的整体设计/结构真的有问题?
【问题讨论】:
标签: nhibernate async-await threadpool message-queue masstransit