【发布时间】: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