C#使用SemaphoreSlim進(jìn)行并發(fā)控制的最佳實(shí)踐
在現(xiàn)代異步編程中,高效處理I/O密集型操作是提升應(yīng)用性能的關(guān)鍵。然而,不加控制的并發(fā)往往會(huì)導(dǎo)致災(zāi)難性后果——下游服務(wù)過(guò)載、數(shù)據(jù)庫(kù)連接池耗盡、內(nèi)存暴漲。本文將深入探討C#中控制異步并發(fā)的標(biāo)準(zhǔn)解決方案:SemaphoreSlim,并提供生產(chǎn)級(jí)別的使用模式。
一、為什么需要控制異步并發(fā)?
假設(shè)我們需要處理1000個(gè)訂單,每個(gè)訂單需要調(diào)用一個(gè)外部支付接口:
// 危險(xiǎn)的反模式:瞬間發(fā)起1000個(gè)HTTP請(qǐng)求
public async Task ProcessOrdersDangerously(List<Order> orders)
{
var tasks = orders.Select(order => CallPaymentApiAsync(order));
await Task.WhenAll(tasks); // 瞬間并發(fā)過(guò)高!
}這種方式會(huì)同時(shí)發(fā)起1000個(gè)HTTP請(qǐng)求,可能導(dǎo)致:
- 目標(biāo)API服務(wù)器拒絕服務(wù)
- 本地網(wǎng)絡(luò)連接池耗盡
- 內(nèi)存使用量激增
- 整體性能反而下降
二、錯(cuò)誤解決方案辨析
在探索解決方案時(shí),開(kāi)發(fā)者常走入以下誤區(qū):
1. 誤用Parallel.ForEach
// 錯(cuò)誤:Parallel.ForEach用于CPU密集型同步操作
Parallel.ForEach(orders, async order =>
{
await CallPaymentApiAsync(order); // 實(shí)際上同步執(zhí)行
});
Parallel.ForEach 設(shè)計(jì)用于同步CPU密集型操作,將其用于異步I/O操作不僅無(wú)法有效控制并發(fā),還會(huì)造成線(xiàn)程池的浪費(fèi)。
2. 分批處理的問(wèn)題
// 次優(yōu)方案:雖能限制并發(fā),但效率低下
for (int i = 0; i < orders.Count; i += 10)
{
var batch = orders.Skip(i).Take(10);
await Task.WhenAll(batch.Select(CallPaymentApiAsync));
await Task.Delay(100); // 人工延遲降低效率
}
這種方法雖然限制了并發(fā)數(shù),但批次間的等待會(huì)導(dǎo)致總體處理時(shí)間延長(zhǎng),無(wú)法充分利用資源。
三、SemaphoreSlim:異步并發(fā)的標(biāo)準(zhǔn)解決方案
SemaphoreSlim 是.NET Framework 4.5引入的輕量級(jí)信號(hào)量,專(zhuān)為async/await設(shè)計(jì),是控制異步并發(fā)的事實(shí)標(biāo)準(zhǔn)。
核心工作機(jī)制
public class AsyncConcurrencyController
{
// 初始化信號(hào)量,設(shè)置最大并發(fā)數(shù)為5
private static readonly SemaphoreSlim _semaphore = new SemaphoreSlim(5, 5);
public async Task ProcessWithConcurrencyControl(List<Item> items)
{
var tasks = items.Select(async item =>
{
// 關(guān)鍵:異步等待信號(hào)量,不阻塞線(xiàn)程
await _semaphore.WaitAsync();
try
{
// 執(zhí)行受保護(hù)的異步操作
await ProcessItemAsync(item);
}
finally
{
// 關(guān)鍵:必須釋放信號(hào)量
_semaphore.Release();
}
});
await Task.WhenAll(tasks);
}
}
工作原理可視化:
初始狀態(tài): [√][√][√][√][√] [ ][ ][ ][ ][ ] ... (20個(gè)任務(wù))
↑ 5個(gè)并發(fā)槽可用
執(zhí)行過(guò)程:
1. 任務(wù)1-5立即獲取信號(hào)量并執(zhí)行
2. 任務(wù)6-20在WaitAsync()處等待
3. 任務(wù)1完成后釋放信號(hào)量
4. 任務(wù)6立即獲取釋放的信號(hào)量并開(kāi)始執(zhí)行
5. 如此循環(huán),始終保持最多5個(gè)并發(fā)
四、生產(chǎn)環(huán)境最佳實(shí)踐
1. 基礎(chǔ)封裝模式
public class ConcurrentExecutor
{
private readonly SemaphoreSlim _semaphore;
public ConcurrentExecutor(int maxConcurrency)
{
_semaphore = new SemaphoreSlim(maxConcurrency, maxConcurrency);
}
public async Task<TResult> ExecuteAsync<TResult>(
Func<Task<TResult>> operation,
CancellationToken cancellationToken = default)
{
await _semaphore.WaitAsync(cancellationToken);
try
{
return await operation();
}
finally
{
_semaphore.Release();
}
}
}
2. 帶超時(shí)控制的增強(qiáng)版本
public async Task<T> ExecuteWithTimeoutAsync<T>(
Func<Task<T>> operation,
TimeSpan timeout,
CancellationToken cancellationToken = default)
{
// 嘗試在指定時(shí)間內(nèi)獲取信號(hào)量
bool acquired = await _semaphore.WaitAsync(timeout, cancellationToken);
if (!acquired)
throw new TimeoutException($"無(wú)法在{timeout.TotalSeconds}秒內(nèi)獲取執(zhí)行許可");
try
{
return await operation();
}
finally
{
_semaphore.Release();
}
}
3. 批量處理與進(jìn)度報(bào)告
public async Task ProcessBatchWithProgressAsync<T>(
IEnumerable<T> items,
Func<T, Task> processor,
int maxConcurrency,
IProgress<int> progress = null,
CancellationToken cancellationToken = default)
{
var semaphore = new SemaphoreSlim(maxConcurrency, maxConcurrency);
int total = items.Count();
int completed = 0;
var tasks = items.Select(async item =>
{
await semaphore.WaitAsync(cancellationToken);
try
{
await processor(item);
}
finally
{
semaphore.Release();
Interlocked.Increment(ref completed);
progress?.Report((completed * 100) / total);
}
});
await Task.WhenAll(tasks);
}
五、高級(jí)應(yīng)用場(chǎng)景
1. 分層并發(fā)控制
// 場(chǎng)景:每個(gè)用戶(hù)最多5個(gè)并發(fā),全局最多50個(gè)并發(fā)
public class TieredConcurrencyController
{
private readonly SemaphoreSlim _globalSemaphore = new(50, 50);
private readonly ConcurrentDictionary<string, SemaphoreSlim> _userSemaphores = new();
public async Task ExecuteForUserAsync(string userId, Func<Task> operation)
{
// 獲取用戶(hù)級(jí)信號(hào)量(每個(gè)用戶(hù)獨(dú)立)
var userSemaphore = _userSemaphores.GetOrAdd(userId, _ => new SemaphoreSlim(5, 5));
// 先獲取全局許可
await _globalSemaphore.WaitAsync();
await userSemaphore.WaitAsync();
try
{
await operation();
}
finally
{
userSemaphore.Release();
_globalSemaphore.Release();
}
}
}
2. 與Polly結(jié)合實(shí)現(xiàn)彈性并發(fā)
public class ResilientConcurrentExecutor
{
private readonly SemaphoreSlim _semaphore;
private readonly AsyncPolicy _retryPolicy;
public async Task<T> ExecuteWithRetryAsync<T>(
Func<Task<T>> operation,
int maxConcurrency)
{
_semaphore = new SemaphoreSlim(maxConcurrency, maxConcurrency);
_retryPolicy = Policy
.Handle<HttpRequestException>()
.WaitAndRetryAsync(3, retryAttempt =>
TimeSpan.FromSeconds(Math.Pow(2, retryAttempt)));
await _semaphore.WaitAsync();
try
{
return await _retryPolicy.ExecuteAsync(operation);
}
finally
{
_semaphore.Release();
}
}
}
六、性能調(diào)優(yōu)與監(jiān)控
1. 動(dòng)態(tài)調(diào)整并發(fā)數(shù)
public class AdaptiveConcurrencyController
{
private SemaphoreSlim _semaphore;
private readonly int _initialConcurrency;
private readonly object _lock = new object();
public void AdjustConcurrencyBasedOnMetrics(
double successRate,
double avgLatency,
int errorCount)
{
lock (_lock)
{
int newLimit = CalculateOptimalConcurrency(
successRate, avgLatency, errorCount);
if (newLimit != _semaphore.CurrentCount)
{
var oldSemaphore = _semaphore;
_semaphore = new SemaphoreSlim(newLimit, newLimit);
// 遷移正在等待的任務(wù)到新信號(hào)量
MigrateWaiters(oldSemaphore, _semaphore);
}
}
}
}
2. 監(jiān)控信號(hào)量狀態(tài)
public class MonitoredSemaphoreSlim : SemaphoreSlim
{
public int CurrentWaitCount { get; private set; }
public TimeSpan AverageWaitTime { get; private set; }
public new async Task WaitAsync(CancellationToken cancellationToken)
{
var stopwatch = Stopwatch.StartNew();
CurrentWaitCount++;
try
{
await base.WaitAsync(cancellationToken);
}
finally
{
stopwatch.Stop();
CurrentWaitCount--;
UpdateAverageWaitTime(stopwatch.Elapsed);
}
}
}
七、注意事項(xiàng)與常見(jiàn)陷阱
- 避免信號(hào)量泄漏:務(wù)必在
finally塊中調(diào)用Release(),確保異常情況下也能釋放 - 不要過(guò)度限制:根據(jù)目標(biāo)服務(wù)的實(shí)際能力設(shè)置合理的并發(fā)數(shù)
- 區(qū)分資源類(lèi)型:
- CPU密集型:使用
Parallel.ForEach或TPL Dataflow - I/O密集型:使用
SemaphoreSlim+async/await
- CPU密集型:使用
- 考慮取消支持:始終傳遞
CancellationToken到WaitAsync()
八、總結(jié)
SemaphoreSlim 是C#異步編程中控制并發(fā)度的標(biāo)準(zhǔn)工具,它提供了輕量級(jí)、非阻塞的并發(fā)控制機(jī)制。通過(guò)正確使用WaitAsync()和Release()方法,配合try...finally確保資源釋放,可以構(gòu)建出高效、穩(wěn)定的異步處理系統(tǒng)。
核心建議:
- 對(duì)于HTTP API調(diào)用、數(shù)據(jù)庫(kù)訪問(wèn)等I/O操作,優(yōu)先使用
SemaphoreSlim - 設(shè)置并發(fā)數(shù)時(shí),考慮目標(biāo)服務(wù)的承受能力和網(wǎng)絡(luò)狀況
- 配合
CancellationToken實(shí)現(xiàn)優(yōu)雅的取消操作 - 在生產(chǎn)環(huán)境中添加適當(dāng)?shù)谋O(jiān)控和日志記錄
正確控制異步并發(fā)不僅能提升應(yīng)用性能,更是構(gòu)建穩(wěn)定、可擴(kuò)展分布式系統(tǒng)的基石。SemaphoreSlim以其簡(jiǎn)潔的API和可靠的行為,成為每個(gè).NET開(kāi)發(fā)者工具箱中不可或缺的工具。
以上就是C#使用SemaphoreSlim進(jìn)行并發(fā)控制的最佳實(shí)踐的詳細(xì)內(nèi)容,更多關(guān)于C# SemaphoreSlim并發(fā)控制的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
使用C#實(shí)現(xiàn)在Excel中創(chuàng)建和配置數(shù)據(jù)透視表
在數(shù)據(jù)分析和業(yè)務(wù)報(bào)告場(chǎng)景中,數(shù)據(jù)透視表(Pivot?Table)是一種強(qiáng)大的數(shù)據(jù)匯總工具,本文將介紹如何使用?C#?在?Excel?工作表中創(chuàng)建和配置數(shù)據(jù)透視表,感興趣的小伙伴可以了解下2026-03-03
C#編寫(xiě)一個(gè)控制臺(tái)程序的實(shí)現(xiàn)串口通信示例
本文主要介紹了C#編寫(xiě)一個(gè)控制臺(tái)程序的實(shí)現(xiàn)示例,實(shí)現(xiàn)串口通信功能,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2025-07-07
教你創(chuàng)建一個(gè)帶診斷工具的.NET鏡像
本文編寫(xiě)的初衷是因?yàn)樵谌豪镉泻芏嘈』锇橛龅缴a(chǎn)環(huán)境性能問(wèn)題的時(shí)候,.NET的runtime鏡像中沒(méi)有帶一些工具,安裝和使用起來(lái)很麻煩,所以分享一些我們公司內(nèi)部一些技巧,對(duì).NET鏡像帶診斷工具相關(guān)知識(shí)感興趣的朋友一起看看吧2022-07-07
Unity?UGUI的RawImage原始圖片組件使用示例詳解
這篇文章主要為大家介紹了Unity?UGUI的RawImage原始圖片組件使用示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-07-07

