【问题标题】:Confusion around CompletionQueue in an async C++ client异步 C++ 客户端中 CompletionQueue 的混淆
【发布时间】:2020-03-08 12:47:41
【问题描述】:

对于如何将CompletionQueue 用于异步 C++ 客户端,我有些困惑。我的服务器是 C#,所以我的问题完全是关于客户端以异步方式向服务器发送请求。

作为参考,以下是我设置客户端代码以发出异步请求的方式:

template<typename ResponseType, typename AsyncOpExecutor>
bool DoAsyncOp(const AsyncOpExecutor& op, const unsigned int deadlineMs, ResponseType& response)
{
    grpc::ClientContext ctx;
    grpc::CompletionQueue queue;

    const std::unique_ptr<grpc::ClientAsyncResponseReaderInterface<ResponseType>> asyncOpResponse = op(ctx, queue);

    grpc::Status status;
    int requestTag = 1;
    asyncOpResponse->Finish(&response, &status, (void*)requestTag);

    bool result = false;
    bool gotEvent = false;

    do
    {
        const std::chrono::time_point<std::chrono::system_clock> deadline = std::chrono::system_clock::now() + std::chrono::milliseconds(deadlineMs);

        void* got_tag;
        bool ok = false;
        const grpc::CompletionQueue::NextStatus nextStatus = queue.AsyncNext(&got_tag, &ok, deadline);

        switch (nextStatus)
        {
        case grpc::CompletionQueue::NextStatus::TIMEOUT:
            continue;

        case grpc::CompletionQueue::NextStatus::GOT_EVENT:
            assert(got_tag == (void*)requestTag);
                        // ok is always true even if I close the server while request is in progress.
            assert(ok);
            result = status.ok();
            gotEvent = true;
            break;

        // Given that I am creating a new CompletionQueue per request (not using a shared one), is this flag likely to occur?
        case grpc::CompletionQueue::NextStatus::SHUTDOWN:
            result = false;
            gotEvent = false;
            break;

        default:
            result = false;
            gotEvent = false;
            break;
        }
    } while (!gotEvent);

    return result;
}

我的第一个困惑是设置CompletionQueue 的最佳方式。 This answer 似乎暗示可以跨请求使用单个 CompletionQueue。如果多个线程使用同一个队列发出请求,这将如何表现。假设我将上面的代码更改为使用共享队列,而不是为每个请求创建一个新队列。

  • 一个线程如何知道它收到的响应是给它的,而不是给另一个线程的?

  • 我是否需要为每个线程分配一个唯一标签,并在每个线程上检查从队列中接收到的标签是否与我最初发送的标签匹配?

  • 如果线程 A 收到一个用于线程 B 的标签,这是否意味着线程 B 以后可以查询它的标签,还是因为线程 A 先看到它而导致该标签丢失?

    李>
  • 每个请求使用新队列而不是共享实际上存在一个主要问题吗?

我的第二个困惑是grpc::CompletionQueue::NextStatus::SHUTDOWN 结果。如果我对每个请求使用一个新队列,并且没有在队列上显式调用shutdown,那么这个结果是否可能发生?如果是,什么会触发它?我执行的一项测试是在请求进行时关闭服务器,但是我得到了grpc::CompletionQueue::NextStatus::GOT_EVENT 结果,状态设置为UNAVAILABLE,而不是得到关闭结果。

我最后的困惑在于ok 标志。我已经阅读了this answer,但是仍然不是很清楚。鉴于上面发布的用例和代码,如果我从队列中得到的结果是grpc::CompletionQueue::NextStatus::GOT_EVENT,ok 标志是否可以为假,如果是,什么会导致它为假?同样,这纯粹是围绕客户端,而不是 CompletionQueue 在服务器上的处理方式。

【问题讨论】:

    标签: grpc


    【解决方案1】:
    1. 使用从完成队列返回的标记来了解哪个完成队列操作启动了此操作。您必须确保没有两个 CQ 操作同时使用相同的标签。
    2. 飞行中的每个 CQ 操作同时需要一个不同的标签。无论这些操作是从同一个线程还是从不同线程启动的,这都适用。
    3. 每个标签仅来自 Next 函数一次。如果您有多个线程在同一个 CompletionQueue 上调用 Next,它们都应该能够处理(或传递)注册到该 CQ 的任何标签。
    4. 您可以为每个 RPC 使用一个新队列,就像同步 API 实际上在内部所做的那样。在这种情况下,它的效率较低,并且基本上会转移到同步 API,因为您将无法同时处理许多未完成的操作。

    您需要在某个时候调用 Shutdown,然后排空队列。这是 API 的标准。否则,特别是对于服务器 CQ,可能会出现泄漏。我的意思是调用 Next 直到它返回 false(或 AsyncNext 直到它返回 SHUTDOWN)。

    对于流式调用、服务器端调用等的 GOT_EVENT,ok 可以为 false。如果您只查看客户端一元调用(如您的示例),它不会为 false。错误的 ok 本质上意味着 RPC 的这一侧已损坏。相关文档位于 CQ 的头文件中。

    【讨论】:

    • 感谢您的回复,非常感谢。您能否详细说明第 3 点?在测试中,我生成了多个使用相同 CQ 的线程。我得到的结果是线程 A 能够获取线程 B 设置的标签,而线程 B 再也无法看到该标签。
    • 是的,一旦一个标签出现在一个线程中,它就永远不会被另一个线程看到。每个线程只能从 CQ 中出来一次,并且任何线程都可以从 CQ 中拉出任何标签,无论哪个线程将其放入 CQ。所以通常你会让每个线程检查标签,然后只根据标签来决定要做什么,而不是根据启动操作的线程。
    • 嗯,在这种情况下,共享 CQ 似乎是在自找麻烦,因为您必须协调许多方面。每个请求都有一个 CQ 似乎是最灵活和最安全的选择。我在这里专门谈论客户端。
    • 我们实际上推荐共享 CQ,因为它可以带来最高的效率(没有工作搁浅,单点轮询);您只需要使用标签指向一个维护该 RPC 当前状态的结构。如果每个请求都有一个单独的 CQ,那么您或多或少会回到同步 API 的设计点,并且可能应该只使用(更简单的)同步 API。
    • 我使用 CQ/Async 的唯一原因是,如果请求处理时间过长(即 TIMEOUT 而不是立即 GOT_EVENT),我不会阻塞 UI。在这种情况下,当我收到 TIMEOUT 时,我可以通知 UI 该请求仍处于待处理状态。尽管建议使用单个 CQ,但它会带来维护状态和确保不同请求不会相互交叉的问题。一个新的 CQ,虽然可能效率不高,但消除了这个问题,同时让我能够使用 AsyncNext。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-06-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-08-20
    • 2014-01-11
    • 1970-01-01
    相关资源
    最近更新 更多