【问题标题】:Return IAsyncEnumerable from grpc with timeout try-catch使用超时 try-catch 从 grpc 返回 IAsyncEnumerable
【发布时间】:2022-11-29 01:20:31
【问题描述】:

我有一个 gRPC 客户端,我想有一个方法来简化它的使用。该方法应返回从 gRPC 服务器流式传输的项目的 IAsyncEnumerable。 我有一个指定的流式传输超时不超过。如果发生超时,我只想带走到目前为止我设法获取的所有项目。

这是我尝试做的:

    public async IAsyncEnumerable<Item> Search(
        SearchParameters parameters, 
        CancellationToken cancellationToken, 
        IDictionary<string, string> headers = null)
    {
        try
        {
            await _client.Search(
                    MapInput(parameters),
                    cancellationToken: cancellationToken,
                    deadline: DateTime.UtcNow.Add(_configuration.Timeout),
                    headers: MapHeaders(headers))
                .ResponseStream.ForEachAsync(item =>
                {
                    yield return MapSingleItem(item); // compilation error
                });
        }
        catch (RpcException ex) when (ex.StatusCode == StatusCode.DeadlineExceeded)
        {
            _logger.LogWarning("Steam finished due to timeout, a limited number of items has been returned");
        }
    }

从逻辑上讲,这应该有效。但是,yield 关键字在 lambda 中不受支持,因此无法编译。还有其他方法可以写吗?

【问题讨论】:

  • 你在使用grpc库吗?
  • 是的,我正在使用 Grapc.Tools 生成客户端。我代码中的_client变量是Grpc.Tools自动生成客户端的结果

标签: c# .net grpc iasyncenumerable


【解决方案1】:

您需要一个中间缓冲区来保存项目,因为 IAsyncEnumerable&lt;Item&gt; 的消费者可以按照自己的节奏枚举它。 Channel&lt;T&gt; 类是用于此目的的优秀异步缓冲区。

您可能要考虑的另一件事是,如果消费者过早地放弃了对 IAsyncEnumerable&lt;Item&gt; 的枚举,会发生什么情况,无论是通过 breaking 或 returning 故意放弃,还是因为它遇到异常而不情愿。你需要注意这种情况,最好的方法是取消iteratorfinally块中的linkedCancellationTokenSource

把所有东西放在一起:

public async IAsyncEnumerable<Item> Search(
    SearchParameters parameters, 
    [EnumeratorCancellation] CancellationToken cancellationToken = default,
    IDictionary<string, string> headers = null)
{
    Channel<Item> channel = Channel.CreateUnbounded<Item>();
    using var linkedCTS = CancellationTokenSource
        .CreateLinkedTokenSource(cancellationToken);

    Task producer = Task.Run(async () =>
    {
        try
        {
            await _client.Search(
                    MapInput(parameters),
                    cancellationToken: linkedCTS.Token,
                    deadline: DateTime.UtcNow.Add(_configuration.Timeout),
                    headers: MapHeaders(headers))
                .ResponseStream.ForEachAsync(item =>
                {
                    channel.Writer.TryWrite(item);
                }).ConfigureAwait(false);
            channel.Writer.Complete();
        }
        catch (Exception ex) { channel.Writer.Complete(ex); }
    });

    try
    {
        await foreach (var item in channel.Reader.ReadAllAsync()
            .ConfigureAwait(false))
        {
            yield return item;
        }
    }
    finally
    {
        linkedCTS.Cancel();
        producer.GetAwaiter().GetResult();
    }
}

当令牌被取消时,生成的 IAsyncEnumerable&lt;Item&gt; 很可能会以 OperationCanceledException 完成。如果您希望您的令牌具有停止语义,您应该首先将其重命名为stoppingToken,然后相应地处理producer任务中的OperationCanceledException异常。

【讨论】:

  • 似乎过于复杂。 Rx.net 可以做得更容易......
  • 谢谢@Theodor。我不明白的一件事是最后一个 finally 块。你说我们应该注意消费者过早完成枚举的情况,这就是为什么需要取消令牌的原因。但是,我不明白它在这段代码中是如何工作的。我们如何发现枚举完成?我的意思是,如果某人只从 IAsyncEnjmerable 中获取 3 个元素,并且没有获取更多元素,那么 finally 块如何触发?另外,我想知道你为什么对生产者使用 GetAwaiter().GetResult()?
【解决方案2】:

使用 Rx.net,您可以使用 .Debounce 运算符和 .TakeUntil 运算符来执行此操作。

var inputObservable = input .ToObservable()
      .Publish()
      .RefCount();

var timeout = inputObs
     .Throttle(TimeSpan.FromSeconds(10));
var outputObs = inputObservable
    .TakeUntil(timeout);
  

return outputObs
     .ToAsyncEnumerable()
     .ToListAsync();

【讨论】:

    猜你喜欢
    • 2012-05-18
    • 1970-01-01
    • 2018-11-05
    • 2015-07-07
    • 2014-04-27
    • 2015-10-19
    • 2016-02-01
    • 1970-01-01
    • 2015-10-24
    相关资源
    最近更新 更多