gpt4 book ai didi

c# - 支持多线程的异步任务队列

转载 作者:太空狗 更新时间:2023-10-29 17:49:15 25 4
gpt4 key购买 nike

我需要实现一个库来请求 vk.com API。问题是 API 每秒仅支持 3 个请求。我想要 API 异步。

重要提示:API 应支持多线程安全访问。

我的想法是实现一些称为 throttler 的类,它允许不超过 3 个请求/秒并延迟其他请求。

接下来是界面:

public interface IThrottler : IDisposable
{
Task<TResult> Throttle<TResult>(Func<Task<TResult>> task);
}

用法就像

var audio = await throttler.Throttle(() => api.MyAudio());
var messages = await throttler.Throttle(() => api.ReadMessages());
var audioLyrics = await throttler.Throttle(() => api.AudioLyrics(audioId));
/// Here should be delay because 3 requests executed
var photo = await throttler.Throttle(() => api.MyPhoto());

如何实现节流?

目前我将其实现为由后台线程处理的队列。

public Task<TResult> Throttle<TResult>(Func<Task<TResult>> task)
{
/// TaskRequest has method Run() to run task
/// TaskRequest uses TaskCompletionSource to provide new task
/// which is resolved when queue processed til this element.
var request = new TaskRequest<TResult>(task);

requestQueue.Enqueue(request);

return request.ResultTask;
}

这是处理队列的后台线程循环的缩短代码:

private void ProcessQueue(object state)
{
while (true)
{
IRequest request;
while (requestQueue.TryDequeue(out request))
{
/// Delay method calculates actual delay value and calls Thread.Sleep()
Delay();
request.Run();
}

}
}

是否可以在没有后台线程的情况下实现这个?

最佳答案

因此,我们将从一个更简单的问题的解决方案开始,即创建一个最多同时处理 N 个任务的队列,而不是限制为每秒启动 N 个任务,并以此为基础:

public class TaskQueue
{
private SemaphoreSlim semaphore;
public TaskQueue()
{
semaphore = new SemaphoreSlim(1);
}
public TaskQueue(int concurrentRequests)
{
semaphore = new SemaphoreSlim(concurrentRequests);
}

public async Task<T> Enqueue<T>(Func<Task<T>> taskGenerator)
{
await semaphore.WaitAsync();
try
{
return await taskGenerator();
}
finally
{
semaphore.Release();
}
}
public async Task Enqueue(Func<Task> taskGenerator)
{
await semaphore.WaitAsync();
try
{
await taskGenerator();
}
finally
{
semaphore.Release();
}
}
}

我们还将使用以下辅助方法将 TaskCompletionSource 的结果与“任务”相匹配:

public static void Match<T>(this TaskCompletionSource<T> tcs, Task<T> task)
{
task.ContinueWith(t =>
{
switch (t.Status)
{
case TaskStatus.Canceled:
tcs.SetCanceled();
break;
case TaskStatus.Faulted:
tcs.SetException(t.Exception.InnerExceptions);
break;
case TaskStatus.RanToCompletion:
tcs.SetResult(t.Result);
break;
}

});
}

public static void Match<T>(this TaskCompletionSource<T> tcs, Task task)
{
Match(tcs, task.ContinueWith(t => default(T)));
}

现在对于我们的实际解决方案,我们可以做的是每次我们需要执行节流操作时,我们创建一个 TaskCompletionSource,然后进入我们的 TaskQueue 并添加一个项目启动任务,匹配 TCS 到它的结果,不等待它,然后延迟任务队列 1 秒。然后任务队列将不允许任务开始,直到过去一秒不再有 N 个任务开始,而操作本身的结果与创建 Task 相同:

public class Throttler
{
private TaskQueue queue;
public Throttler(int requestsPerSecond)
{
queue = new TaskQueue(requestsPerSecond);
}
public Task<T> Enqueue<T>(Func<Task<T>> taskGenerator)
{
TaskCompletionSource<T> tcs = new TaskCompletionSource<T>();
var unused = queue.Enqueue(() =>
{
tcs.Match(taskGenerator());
return Task.Delay(TimeSpan.FromSeconds(1));
});
return tcs.Task;
}
public Task Enqueue<T>(Func<Task> taskGenerator)
{
TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>();
var unused = queue.Enqueue(() =>
{
tcs.Match(taskGenerator());
return Task.Delay(TimeSpan.FromSeconds(1));
});
return tcs.Task;
}
}

关于c# - 支持多线程的异步任务队列,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/34315589/

25 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com