【问题标题】:Redesigning and removing Task.Delay() from my C# function从我的 C# 函数中重新设计和删除 Task.Delay()
【发布时间】:2020-12-11 20:45:49
【问题描述】:

我的下面的代码在Task.Delay(500000) 时运行,但是当我的Task.Delay(5000) 时它没有给我任何结果,因为执行预期输出的持续时间非常短。我正在寻找一种重新设计代码的方法,我可以在没有Task.Delay() 的情况下处理这个问题,因为每次执行的响应时间可能会有所不同。我该怎么做?

注意:使用 Roald 在回答中建议的方法修改了代码。 Task.Delay() 上的早期查询无法为 Change Feed 异步处理。另一种方法是使用拉模型而不是推模型


    using Microsoft.Azure.Cosmos;
    using Microsoft.Azure.Documents.ChangeFeedProcessor;
    using System;
    using System.Collections.Generic;
    using System.Net;
    using System.Runtime.CompilerServices;
    using System.Threading;
    using System.Threading.Channels;
    using System.Threading.Tasks;
    
    namespace ConsoleApp1
    {
        public class ChangeFeedProcessorOptions
        {
            public int BufferCapacity { get; set; }
            public string ProcessorName { get; set; }
            public Container LeaseContainer { get; set; }
            public string InstanceName { get; set; }
            public DateTime StartTime { get; set; }
        }
        
        class Program
        {
    
            
    
            static async Task Main()
            {
                var client = new CosmosClient("AccountEndpoint = https://test.documents.azure.com:443/;AccountKey=oaEOA==;");
    
                var database = client.GetDatabase("testDatabase");
                var container = database.GetContainer("testContainer");
    
                var options = new ChangeFeedProcessorOptions
                {
                    BufferCapacity = 10,
                    InstanceName = "ChangeFeedInstanceName",
                    LeaseContainer = database.GetContainer("leases"),
                    ProcessorName = "ChangeFeedProcessorName",
                    StartTime = DateTime.Now.AddDays(-7).ToUniversalTime()
                };
    
                var count = 0;
                await foreach (var doc in container.GetChangeFeed<document>(options))
                {
                    Console.Write(doc, b: true);
    
                    count++;
    
                    if (count == 6)
                    {
                        break;
                    }
                }
            }     
    
    
            
    
        public static async IAsyncEnumerable<document> GetChangeFeed<document>(this Container self, ChangeFeedProcessorOptions options, [EnumeratorCancellation] CancellationToken cancellationToken = default)
        {
            var channel = Channel.CreateBounded<document>(new BoundedChannelOptions(options.BufferCapacity)
            {
                FullMode = BoundedChannelFullMode.Wait,
                SingleReader = true,
                SingleWriter = true
            });
    
            var processor = self
                .GetChangeFeedProcessorBuilder<document>(options.ProcessorName, async (items, cancellation) =>
                {
                    foreach (var item in items)
                    {
                        await channel.Writer.WriteAsync(item, cancellation);
                    }
                })
                .WithLeaseContainer(options.LeaseContainer)
                .WithInstanceName(options.InstanceName)
                .WithStartTime(options.StartTime)
                .Build();
    
            await processor.StartAsync();
            try
            {
                await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken))
                {
                    yield return item;
                }
            }
            finally
            {
                await processor.StopAsync();
            }
        }
    }
    }

【问题讨论】:

  • ProcessChanges 是不是一直在耗费时间?也许你可以用TaskCompletionSource 或SempahoreSlim 做点什么?
  • 行“await cfp.StartAsync();”应该阻塞直到收到所有数据(完整响应)并且不需要延迟。如果等待被删除,那么您将需要延迟。
  • ProcessChanges 确实需要时间,它是处理更改的委托。你能详细说明一下TaskCompletionSource
  • 它不适用于 Task.Delay() 删除
  • 只使用Task.Run 来block 等待结果有什么意义? Task.Delay() 无论如何都不需要。无论StartAsync 和StopAsync 做什么,在已经阻塞的呼叫中都不需要延迟。由于缺少相关代码,人们不得不猜测出了什么问题

标签: c# async-await


【解决方案1】:

我认为最灵活的方法是首先将更改提要转换为 IAsyncEnumerable,这样您就可以使用 linq 或一些直接的命令式代码来处理它。

您可以使用此扩展方法获取 IAsyncEnumerable

SDK 版本

public record ChangeFeedProcessorOptions
{
    public int BufferCapacity { get; init; }
    public string ProcessorName { get; init; }
    public Container LeaseContainer { get; init; }
    public string InstanceName { get; init; }
    public DateTime StartTime { get; init; }
}

public static async IAsyncEnumerable<T> GetChangeFeed<T>(this Container self, ChangeFeedProcessorOptions options, [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    var channel = Channel.CreateBounded<T>(new BoundedChannelOptions(options.BufferCapacity)
    {
        FullMode = BoundedChannelFullMode.Wait,
        SingleReader = true,
        SingleWriter = true
    });

    var processor = self
        .GetChangeFeedProcessorBuilder<T>(options.ProcessorName, async (items, cancellation) =>
        {
            foreach (var item in items)
            {
                await channel.Writer.WriteAsync(item, cancellation);
            }
        })
        .WithLeaseContainer(options.LeaseContainer)
        .WithInstanceName(options.InstanceName)
        .WithStartTime(options.StartTime)
        .Build();

    await processor.StartAsync();
    try
    {
        await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken))
        {
            yield return item;
        }
    }
    finally
    {
        await processor.StopAsync();
    }
}

SDK 版本 > 3.15.0

public record ChangeFeedProcessorOptions
{
    public DateTime StartTime { get; init; }
    public TimeSpan PollInterval { get; init; }
}

public static async IAsyncEnumerable<T> GetChangeFeed<T>(this Container self, ChangeFeedProcessorOptions options, [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    var iterator = self.GetChangeFeedIterator<T>(ChangeFeedStartFrom.Time(options.StartTime));

    while (iterator.HasMoreResults)
    {
        FeedResponse<T> items;

        try
        {
            items = await iterator.ReadNextAsync(cancellationToken);
        }
        catch (CosmosException ex) when (ex.StatusCode == HttpStatusCode.NotModified)
        {
            // No changes
            await Task.Delay(options.PollInterval, cancellationToken);
            continue;
        }

        foreach (var item in items)
        {
            yield return item;
        }
    }
}

然后像这样使用它:

static async Task Main()
{
    var client = new CosmosClient("AccountEndpoint = https://test.documents.azure.com:443/;AccountKey=oaEOA==;");

    var database = client.GetDatabase("testDatabase");
    var container = database.GetContainer("testContainer");

    var options = new ChangeFeedProcessorOptions
    {
        BufferCapacity = 10,
        InstanceName = ChangeFeedInstanceName,
        LeaseContainer = database.GetContainer("leases"),
        ProcessorName = ChangeFeedProcessorName,
        StartTime = DateTime.Now.Subtract(MaxAge).ToUniversalTime()
    };
    
    var count = 0;
    await foreach (var doc in container.GetChangeFeed<Recording>(options))
    {
        WriteObject(doc, b: true);
        
        count++;
        
        if (count == 6)
        {
            break;
        }
    }
}

如果您添加System.Linq.Async,甚至更好

await foreach (var doc in container.GetChangeFeed<Recording>(options).Take(6))
{
    WriteObject(doc, b: true);
}

对于不同的实现,您还可以查看here,它使用两个信号量而不是通道来实现相同的结果。

powershell cmdlet 中的异步代码

您遇到的问题是由于在 cmdlet 中调用 WriteObject、WriteVerbose、WriteWarning 等需要来自主线程。
为了解决这个问题,您需要在ProcessRecord 中运行一个消息泵,并在您需要调用这些方法中的任何一个时使用它来回发到主线程,这正是您在 WinForm 或 WPF 中使用 Dispatcher 必须执行的操作。 处理这个问题的库是PowerShellAsync

使用您的代码将成为的库

[Cmdlet(VerbsCommon.Get, "ChangedRecording")]
[OutputType(typeof(Recording))]
public class SyncRecording : AsyncCmdlet
{
    // ...
    
    protected override async Task ProcessRecordAsync()
    {
        var container = ...;
        await foreach (var doc in container.GetChangeFeed<Recording>(options).Take(6))
        {
            WriteObject(doc, b: true);
        }
    }
}

【讨论】:

  • 嗨 Roald,经过几天的尝试,我终于转向了您的方法。您上面建议的代码几乎没有错误,例如公共记录 ChangeFeedProcessorOptions 什么程序集引用记录?
  • @AnkitKumar 记录和仅初始化属性是 c#9 的特性,你可以在这里阅读它们 devblogs.microsoft.com/dotnet/c-9-0-on-the-record ,如果你不能使用 c#9,只需将 'record' 替换为 'class' 和'init' 和 'set'。
  • @AnkitKumar 是的,我测试了它.. 你需要阅读那些错误信息!编译你的代码它说“[CS1106] 扩展方法必须在非泛型静态类中定义”。 GetChangeFeed 是一种扩展方法,因此必须在静态类中声明,但您将其复制到非静态的 Program 中。扩展方法中的更多信息:docs.microsoft.com/en-us/dotnet/csharp/programming-guide/…
  • 现在留下一个错误 container.GetChangeFeedBuilder(options)) ,它说容器不包含 GetChangeFeedBuilder 的定义。你用的是什么包
  • 是的,没错,GetChangeFeedProcessorBuilder 可用,但由于未获得 GetAsyncEnumerator 的定义,它会给出错误。如果您测试了代码,您是如何获得 GetChangedFeedBuilder 的??
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-02-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-12-28
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多