【问题标题】:Receive concurrent asynchronous requests and process them one at a time接收并发的异步请求并一次处理一个
【发布时间】:2018-05-05 00:09:37
【问题描述】:

背景

我们有一个可以接收并发异步请求的服务操作,并且必须一次处理一个请求。

在以下示例中,UploadAndImport(...) 方法在多个线程上接收并发请求,但它对 ImportFile(...) 方法的调用必须一次发生一个。

外行描述

想象一个有许多工人(多线程)的仓库。人们(客户)可以同时(并发)向仓库发送多个包裹(请求)。当一个包裹进来时,工人从头到尾负责它,而放下包裹的人可以离开(即发即弃)。工人的工作是将每个包裹放入一个小滑槽,一次只能有一个工人将一个包裹放入滑槽,否则会出现混乱。如果投递包裹的人稍后签到(轮询端点),仓库应该能够报告包裹是否从滑槽下落。

问题

那么问题是如何编写一个服务操作...

  1. 可以接收并发的客户端请求,
  2. 在多个线程上接收和处理这些请求,
  3. 在接收请求的同一线程上处理请求,
  4. 一次处理一个请求,
  5. 是一种单向即发即弃的操作,并且
  6. 有一个单独的轮询端点,用于报告请求完成情况。

我们尝试了以下方法并且想知道两件事:

  1. 有没有我们没有考虑过的竞争条件?
  2. 是否有更规范的方式在 C#.NET 中使用面向服务的架构(我们碰巧使用 WCF)编写此场景?

示例:我们尝试了什么?

这是我们尝试过的服务代码。尽管感觉有点像 hack 或 kludge,但它确实有效。

static ImportFileInfo _inProgressRequest = null;

static readonly ConcurrentDictionary<Guid, ImportFileInfo> WaitingRequests = 
    new ConcurrentDictionary<Guid, ImportFileInfo>();

public void UploadAndImport(ImportFileInfo request)
{
    // Receive the incoming request
    WaitingRequests.TryAdd(request.OperationId, request);

    while (null != Interlocked.CompareExchange(ref _inProgressRequest, request, null))
    {
        // Wait for any previous processing to complete
        Thread.Sleep(500);
    }

    // Process the incoming request
    ImportFile(request);

    Interlocked.Exchange(ref _inProgressRequest, null);
    WaitingRequests.TryRemove(request.OperationId, out _);
}

public bool UploadAndImportIsComplete(Guid operationId) => 
    !WaitingRequests.ContainsKey(operationId);

这是示例客户端代码。

private static async Task UploadFile(FileInfo fileInfo, ImportFileInfo importFileInfo)
{
    using (var proxy = new Proxy())
    using (var stream = new FileStream(fileInfo.FullName, FileMode.Open, FileAccess.Read))
    {
        importFileInfo.FileByteStream = stream;
        proxy.UploadAndImport(importFileInfo);
    }

    await Task.Run(() => Poller.Poll(timeoutSeconds: 90, intervalSeconds: 1, func: () =>
    {
        using (var proxy = new Proxy())
        {
            return proxy.UploadAndImportIsComplete(importFileInfo.OperationId);
        }
    }));
}

很难在 Fiddle 中编写一个最小可行的示例,但 here is a start 给出了一种感觉并且可以编译。

与以前一样,上述内容似乎是一种 hack/kludge,我们正在询问其方法中的潜在缺陷以及更合适/规范的替代模式。

【问题讨论】:

  • 您说它必须在单独的线程上处理请求,但在原始线程上完成。处理和完成有什么区别?
  • 我的意思是接收请求的同一个线程也必须完成请求。
  • 请根据您的问题定义“完整”并告诉我它与“处理”有何不同。另外,如果完成发生在同一个线程上,那么有一个端点来轮询状态的目的是什么?
  • @JohnWu 我已经将“完成”和“处理”替换为“过程”一词,以使情况更清晰。
  • @JohnWu 进行轮询的原因是服务操作是单向的,即发即弃操作。

标签: c# .net multithreading asynchronous concurrency


【解决方案1】:

在线程数限制的情况下使用生产者-消费者模式来管道请求的简单解决方案。

您仍然需要实现一个简单的进度报告器或事件。我建议用微软的SignalR 库提供的异步通信来代替昂贵的轮询方法。它使用 WebSocket 来启用异步行为。客户端和服务器可以在集线器上注册它们的回调。使用 RPC,客户端现在可以调用服务器端方法,反之亦然。您将使用集线器(客户端)向客户端发布进度。以我的经验,SignalR 使用起来非常简单,并且有很好的文档记录。它有一个适用于所有著名服务器端语言(例如 Java)的库。

在我看来,轮询与“一劳永逸”完全相反。您不能忘记,因为您必须根据间隔检查某些内容。基于事件的通信,如 SignalR,是即发即忘的,因为你开火并会收到提醒(因为你忘记了)。 “事件端”将调用您的回调,而不是等待您自己做!

要求 5 被忽略,因为我没有得到任何理由。等待线程完成将消除火灾并忘记字符。

private BlockingCollection<ImportFileInfo> requestQueue = new BlockingCollection<ImportFileInfo>();
private bool isServiceEnabled;
private readonly int maxNumberOfThreads = 8;
private Semaphore semaphore = new Semaphore(numberOfThreads);
private readonly object syncLock = new object();

public void UploadAndImport(ImportFileInfo request) 
{            
  // Start the request handler background loop
  if (!this.isServiceEnabled)
  {
    this.requestQueue?.Dispose();
    this.requestQueue = new BlockingCollection<ImportFileInfo>();

    // Fire and forget (requirement 4)
    Task.Run(() => HandleRequests());
    this.isServiceEnabled = true;
  }

  // Cache multiple incoming client requests (requirement 1) (and enable throttling)
  this.requestQueue.Add(request);
}

private void HandleRequests()
{
  while (!this.requestQueue.IsCompleted)
  {
    // Wait while thread limit is exceeded (some throttling)
    this.semaphore.WaitOne();

    // Process the incoming requests in a dedicated thread (requirement 2) until the BlockingCollection is marked completed.
    Task.Run(() => ProcessRequest());
  }

  // Reset the request handler after BlockingCollection was marked completed
  this.isServiceEnabled = false;
  this.requestQueue.Dispose();
}

private void ProcessRequest()
{
  ImportFileInfo request = this.requestQueue.Take();
  UploadFile(request);

  // You updated your question saying the method "ImportFile()" requires synchronization.
  // This a bottleneck and will significantly drop performance, when this method is long running. 
  lock (this.syncLock)
  {
    ImportFile(request);
   }

  this.semaphore.Release();
}

备注:

  • BlockingCollection 是 IDisposable
  • TODO:您必须通过将 BlockingCollection 标记为已完成来“关闭”它: “BlockingCollection.CompleteAdding()”,否则它将不确定地循环等待进一步的请求。也许您为客户端引入了一个额外的请求方法来取消和/或更新进程并将添加到 BlockingCollection 中标记为已完成。或者在将其标记为已完成之前等待空闲时间的计时器。或者让您的请求处理程序线程阻塞或旋转。
  • 如果您需要取消支持,请将 Take() 和 Add(...) 替换为 TryTake(...) 和 TryAdd(...)
  • 代码未经测试
  • 您的“ImportFile()”方法是多线程环境中的瓶颈。我建议让它线程安全。对于需要同步的 I/O,我会将数据缓存在 BlockingCollection 中,然后将它们一一写入 I/O。

【讨论】:

  • 感谢您为我指明了生产者-消费者的方向,举了一个例子,并提供了一些陷阱。这为我提供了另一种调查方法。
  • @Shaun Luttin 好的,我想你编辑了你的问题。让我更新我的示例,通过引入类似于链式 BlockingCoillections 的管道来同步 ImportFIle() 上的调用。但是 UploadFile() 打算并行运行?
  • 是的。很抱歉编辑我的问题。如果这看起来合适,我可以将其回滚。您询问UploadFile() 是否应该是并行的。我对这个问题进行了更多思考,归结为:UploadAndImport 有两个部分。第一部分可以并行,第二部分必须同步。
  • 如果您有兴趣,我已经将我原来的 .NET Fiddle 分叉了一个我认为更清晰的版本:dotnetfiddle.net/VUzdXR 内存中的队列仍然具有@JohnWu 提到的脆弱性如果我们收到大量请求,就会发生。
  • @ShaunLuttin 我对您的环境一无所知。如果您期望“大量请求”,那么您将遇到比内存等硬件资源更严重的问题。您的响应时间太长,让您的客户等待。您必须提前知道预期的负载并相应地设计您的系统。当处理每个请求的执行时间太长时,这里就会出现瓶颈。然后,您可以并行托管多个服务,并可以使用反向代理来平衡负载。或者你可以买一个更大的CPU。水平或垂直缩放。
【解决方案2】:

问题在于您的总带宽非常小——一次只能运行一个作业——并且您想要处理并行请求。这意味着排队时间可能会有很大差异。在内存中实现作业队列可能不是最佳选择,因为它会使您的系统更加脆弱,并且随着业务的增长更难以横向扩展。

一种传统的、可扩展的架构方式是:

  • 接受请求的 HTTP 服务,负载平衡/冗余,没有会话状态。
  • 一个 SQL Server 数据库,用于将请求保留在队列中,并返回一个持久的唯一作业 ID。
  • 用于处理队列、一次一个作业并将作业标记为完成的 Windows 服务。服务的工作进程可能是单线程的。

此解决方案要求您选择 Web 服务器。一个常见的选择是运行 ASP.NET 的 IIS。在该平台上,每个请求都保证以单线程方式处理(即您无需过多担心竞争条件),但由于名为thread agility 的功能,请求可能以不同的线程结束,但在原始同步上下文中,这意味着您可能永远不会注意到,除非您正在调试和检查线程 ID。

【讨论】:

  • 在我们的场景中,该服务是一个多线程 Windows 服务,它通过 TCP 接收来自同一台机器上的客户端的请求。
  • 感谢您提及内存队列的脆弱性。
【解决方案3】:

鉴于我们系统的约束上下文,这是我们最终使用的实现:

static ImportFileInfo _importInProgressItem = null;

static readonly ConcurrentQueue<ImportFileInfo> ImportQueue = 
    new ConcurrentQueue<ImportFileInfo>();

public void UploadAndImport(ImportFileInfo request) {
    UploadFile(request);
    ImportFileSynchronized(request);
}

// Synchronize the file import, 
// because the database allows a user to perform only one write at a time.
private void ImportFileSynchronized(ImportFileInfo request) {
    ImportQueue.Enqueue(request);
    do {
        ImportQueue.TryPeek(out var next);
        if (null != Interlocked.CompareExchange(ref _importInProgressItem, next, null)) {
            // Queue processing is already under way in another thread.
            return;
        }

        ImportFile(next);
        ImportQueue.TryDequeue(out _);
        Interlocked.Exchange(ref _importInProgressItem, null);
    }
    while (ImportQueue.Any());
}

public bool UploadAndImportIsComplete(Guid operationId) =>
    ImportQueue.All(waiting => waiting.OperationId != operationId);

此解决方案适用于我们预期的负载。该负载最多涉及大约 15-20 个并发 PDF 文件上传。多达 15-20 个文件的批次往往会同时到达,然后安静几个小时,直到下一批到达。

欢迎批评和反馈。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-12-20
    • 1970-01-01
    • 2020-07-12
    • 2020-11-24
    • 2017-08-20
    • 1970-01-01
    • 2019-05-13
    相关资源
    最近更新 更多