Tokio 如何用 Semaphore 的 acquire 与 acquire_many 控制并发并处理 AcquireError
【免费下载链接】tokioA runtime for writing reliable asynchronous applications with Rust. Provides I/O, networking, scheduling, timers, ...项目地址: https://gitcode.com/GitHub_Trending/to/tokio
在实际的异步服务里,常见的需求是:同时打开的文件数、同时发出的请求数、同时处理的连接数不能无限增长。Tokio 在tokio::sync模块中提供的Semaphore就是为此设计的——模块文档把它概括为 "Limits the amount of concurrency":信号量持有一组 permit(许可),任务进入临界区前必须先拿到 permit。本文基于仓库内的 API 文档与测试代码(semaphore.rs、batch_semaphore.rs、sync_semaphore.rs),演示如何用acquire与acquire_many限制并发,并在信号量被关闭时正确识别和处理AcquireError。
前提说明:Semaphore位于tokio::sync,属于syncfeature 控制的同步原语;其阻塞变体(如blocking_acquire)在源码中同样以#[cfg(feature = "sync")]标注。模块文档还指出,这些同步原语是 runtime-agnostic 的,可以跨 Tokio 运行时实例使用;仅在 Tokio 运行时内使用时才会参与协作式调度。下文示例沿用仓库文档示例的运行方式(#[tokio::main])。
创建信号量并理解 acquire 的行为
Semaphore::new(permits)创建初始 permit 数为permits的信号量;如果permits超过MAX_PERMITS(文档定义为usize::MAX >> 3),会 panic。acquire的语义是:
- 有剩余 permit 时,立即返回一个
SemaphorePermit; - 没有剩余 permit 时,异步等待,直到有 permit 被释放后按队列顺序分发给调用者;
- 信号量已关闭(closed)时,返回
AcquireError。
下面的代码来自 semaphore.rs 中acquire方法上的官方文档示例,它同时演示了如何用available_permits()验证剩余 permit 数:
use tokio::sync::Semaphore; #[tokio::main(flavor = "current_thread")] async fn main() { let semaphore = Semaphore::new(2); let permit_1 = semaphore.acquire().await.unwrap(); assert_eq!(semaphore.available_permits(), 1); let permit_2 = semaphore.acquire().await.unwrap(); assert_eq!(semaphore.available_permits(), 0); drop(permit_1); assert_eq!(semaphore.available_permits(), 1); }assert_eq!就是文档给出的验证方式:每次获取后剩余 permit 数应当等于初始值减去已持有数,释放(drop)后回涨。permit 的归还发生在它被 drop 时——SemaphorePermit的Drop实现会调用sem.add_permits(self.permits)把 permit 还给信号量。文档示例(如限制文件句柄数的例子)还展示了惯用写法:用let _permit = PERMITS.acquire().await.unwrap();让 permit 的作用域覆盖整个文件操作,函数返回时自动归还。
用 acquire_many 一次获取多个 permit
acquire_many(n: u32)一次获取n个 permit,返回单个持有n个 permit 的SemaphorePermit;失败条件与acquire相同(仅关闭时返回AcquireError)。文档示例:
use tokio::sync::Semaphore; #[tokio::main(flavor = "current_thread")] async fn main() { let semaphore = Semaphore::new(5); let permit = semaphore.acquire_many(3).await.unwrap(); assert_eq!(semaphore.available_permits(), 2); }使用acquire_many必须注意文档明确写出的公平性语义:permit 按请求顺序发放,如果队列前排的acquire_many请求的 permit 数超过当前可用数,它会阻塞后面的普通acquire完成,即使信号量里的 permit 足够满足后者。仓库测试 semaphore_batch.rs 中的poll_acquire_many_unavailable验证了这一点:一个请求 5 个 permit 的等待者排在前面时,后面请求 3 个的等待者即使总量足够也不会被唤醒,直到前者被满足。设计批量获取时要考虑这一队头阻塞效应。
两个方法都声明了相同的取消安全性(Cancel safety):acquire/acquire_many内部用队列公平分配 permit,取消调用会使你失去在队列中的位置,而不是保留位置。
跨任务共享:Arc 与 acquire_owned
如果并发任务由tokio::spawn产生,每个任务都要引用同一个信号量,文档给出的做法是把信号量包在Arc里再 clone 给各任务。文档示例中限制并发请求数的完整模式(取自 semaphore.rs 的模块文档):
use std::sync::Arc; use tokio::sync::Semaphore; #[tokio::main(flavor = "current_thread")] async fn main() { // Define maximum number of parallel requests. let semaphore = Arc::new(Semaphore::new(5)); // Spawn many tasks that will send requests. let mut jhs = Vec::new(); for task_id in 0..50 { let semaphore = semaphore.clone(); let jh = tokio::spawn(async move { // Acquire permit before sending request. let _permit = semaphore.acquire().await.unwrap(); // Send the request. let response = send_request(task_id).await; // Drop the permit after the request has been sent. drop(_permit); response }); jhs.push(jh); } // Collect responses from tasks. let mut responses = Vec::new(); for jh in jhs { let response = jh.await.unwrap(); responses.push(response); } } async fn send_request(task_id: usize) { // 你的实际请求逻辑 }这里的关键判断是:permit 必须在任务边界内保持存活。acquire返回的SemaphorePermit<'a>借用信号量,不能跨spawn移动;当 permit 需要在 spawn 之前获取并移入任务时,要用acquire_owned(信号量必须是Arc<Semaphore>),它返回可跨任务移动的OwnedSemaphorePermit。文档中限制入站连接的示例正是这种用法:在 accept 循环中先acquire_owned().await拿到 permit,再 spawn 任务,任务内先 drop socket、最后 drop permit,从而保证同一时刻至多 N 个连接在被处理。
处理 AcquireError:它只在信号量关闭时出现
AcquireError的定义在 batch_semaphore.rs 中,文档写得很直接:
An
acquireoperation can only fail if the semaphore has been closed.
也就是说,permit 不足不会让acquire/acquire_many报错,而是让它们异步等待;唯一的错误路径是信号量已被close()。这一点决定了错误处理写法:拿到Err(AcquireError)时,说明系统进入了关闭状态,应当停止重试、走优雅退出路径,而不是当作"暂时没资源"来处理。
close()的文档说明:阻止信号量再发出新的 permit,并通知所有等待中的任务。官方示例展示了完整流程(等待者被唤醒后收到错误,之后try_acquire返回TryAcquireError::Closed):
use tokio::sync::Semaphore; use std::sync::Arc; use tokio::sync::TryAcquireError; #[tokio::main(flavor = "current_thread")] async fn main() { let semaphore = Arc::new(Semaphore::new(1)); let semaphore2 = semaphore.clone(); tokio::spawn(async move { let permit = semaphore.acquire_many(2).await; assert!(permit.is_err()); println!("waiter received error"); }); println!("closing semaphore"); semaphore2.close(); // Cannot obtain more permits assert_eq!(semaphore2.try_acquire().err(), Some(TryAcquireError::Closed)) }这段代码同时给出了验证方法:close()之后,try_acquire()的错误应是TryAcquireError::Closed,而之前挂起的acquire_many(2)会返回错误(示例中的断言permit.is_err()即验证点)。注意非零 permit 的关闭状态是允许的:仓库测试 sync_semaphore.rs 的add_permits_closed验证了 close 之后add_permits仍会增加available_permits()的计数,且is_closed()保持为true。
在写错误处理时,注意区分两个错误类型:
| 方法 | 返回值 | 可能的错误 |
|---|---|---|
acquire/acquire_many(含_owned变体) | Result<..., AcquireError> | 仅AcquireError(信号量已关闭) |
try_acquire/try_acquire_many | Result<..., TryAcquireError> | Closed(已关闭)或NoPermits(无可用 permit) |
TryAcquireError的定义同样在 batch_semaphore.rs 中。如果你需要"拿不到就走别的路"的逻辑,用try_acquire系列并匹配NoPermits;而await acquire()拿到Err就只有一种含义——被关闭。配合is_closed()布尔方法可以在进入等待前主动检查。
相关限制与注意事项
文档中还有几条会直接影响使用方式的约束:
- 内存顺序保证:文档的 Memory ordering 一节说明,acquire(含
try系列与_owned变体)、释放 permit(drop、add_permits、forget_permits)和close都是AcqRel操作且全序排列。因此一个任务写数据后释放 permit,之后拿到该 permit 的任务保证能看到这些数据——用信号量在共享状态中传递数据是安全的。 - 不想归还 permit 时用
forget:SemaphorePermit::forget()会把 permit 标记为"遗忘",drop 时不再归还,可用于永久减少信号量的容量(文档中的 token bucket 限流示例正是靠add_permits+forget实现令牌不回收)。如果你的业务里出现"permit 数量莫名只减不增",先检查是否误用了forget。 - 阻塞变体的使用条件:
blocking_acquire/blocking_acquire_many是acquire/acquire_many的阻塞等价物,文档指定它们用于同步代码(如spawn_blocking内部或与同步代码共享信号量时),并在异步执行上下文中调用会 panic。 - 静态实例用
const_new:Semaphore::const_new(permits)允许写static信号量(文档示例中用static PERMITS: Semaphore = Semaphore::const_new(100)限制文件句柄数)。文档同时提醒:开启 unstable tracing 时,const_new创建的信号量不会被 instrument,不会出现在 tokio-console 中,需要被观测时应改用Semaphore::new。
验证清单
把上述步骤落到代码后,可以直接对照文档给出的检查点:
- 每次
acquire/acquire_many成功后,available_permits()应等于初始值减去本次获取数(文档示例均以assert_eq!验证); - permit drop 后计数回涨;若计数不回涨,检查是否调用了
forget; close()之后,新的acquire调用返回Err,挂起中的等待者被唤醒并返回Err,try_acquire返回TryAcquireError::Closed;- 并发上界由初始 permit 数决定:文档中"限制同时打开文件数"与"限制并发请求数"两个示例分别以 100 和 5 个 permit 作为上限。
仓库中还有一套可参考的行为验证代码:tests/sync_semaphore.rs(try_acquire、forget、add_permits等)和 sync/tests/semaphore_batch.rs(队列顺序、等待者唤醒时机)。当你对"permit 何时被唤醒、按什么顺序发放"存疑时,这两个测试文件是最贴近实现行为的参照。
【免费下载链接】tokioA runtime for writing reliable asynchronous applications with Rust. Provides I/O, networking, scheduling, timers, ...项目地址: https://gitcode.com/GitHub_Trending/to/tokio
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考