news 2026/9/12 17:44:00

Tokio 如何用 Semaphore 的 acquire 与 acquire_many 控制并发并处理 AcquireError

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Tokio 如何用 Semaphore 的 acquire 与 acquire_many 控制并发并处理 AcquireError

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),演示如何用acquireacquire_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 时——SemaphorePermitDrop实现会调用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 中,文档写得很直接:

Anacquireoperation 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_manyResult<..., TryAcquireError>Closed(已关闭)或NoPermits(无可用 permit)

TryAcquireError的定义同样在 batch_semaphore.rs 中。如果你需要"拿不到就走别的路"的逻辑,用try_acquire系列并匹配NoPermits;而await acquire()拿到Err就只有一种含义——被关闭。配合is_closed()布尔方法可以在进入等待前主动检查。

相关限制与注意事项

文档中还有几条会直接影响使用方式的约束:

  1. 内存顺序保证:文档的 Memory ordering 一节说明,acquire(含try系列与_owned变体)、释放 permit(drop、add_permitsforget_permits)和close都是AcqRel操作且全序排列。因此一个任务写数据后释放 permit,之后拿到该 permit 的任务保证能看到这些数据——用信号量在共享状态中传递数据是安全的。
  2. 不想归还 permit 时用forgetSemaphorePermit::forget()会把 permit 标记为"遗忘",drop 时不再归还,可用于永久减少信号量的容量(文档中的 token bucket 限流示例正是靠add_permits+forget实现令牌不回收)。如果你的业务里出现"permit 数量莫名只减不增",先检查是否误用了forget
  3. 阻塞变体的使用条件blocking_acquire/blocking_acquire_manyacquire/acquire_many的阻塞等价物,文档指定它们用于同步代码(如spawn_blocking内部或与同步代码共享信号量时),并在异步执行上下文中调用会 panic
  4. 静态实例用const_newSemaphore::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,挂起中的等待者被唤醒并返回Errtry_acquire返回TryAcquireError::Closed
  • 并发上界由初始 permit 数决定:文档中"限制同时打开文件数"与"限制并发请求数"两个示例分别以 100 和 5 个 permit 作为上限。

仓库中还有一套可参考的行为验证代码:tests/sync_semaphore.rs(try_acquireforgetadd_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 17:42:26

spaCy 如何通过 entry points 机制注册自定义组件与语言类?

spaCy 如何通过 entry points 机制注册自定义组件与语言类&#xff1f; 【免费下载链接】spaCy &#x1f4ab; Industrial-strength Natural Language Processing (NLP) in Python 项目地址: https://gitcode.com/GitHub_Trending/sp/spaCy 写出自定义 pipeline 组件或自…

作者头像 李华
网站建设 2026/9/12 17:38:54

2026年7月九江市新房价格深度分析报告

一、报告背景与数据说明本报告基于2026年7月九江市新房实际成交案例&#xff0c;结合区域分布、楼盘定位、户型结构与成交价格等维度&#xff0c;对当前九江新房市场进行深度剖析。报告数据来源于公开成交备案信息与典型楼盘样本&#xff0c;旨在为购房者、投资者及行业研究者提…

作者头像 李华
网站建设 2026/9/12 17:38:36

微信小程序消息怎么通过API发送?个人微信API接口功能应用

小程序卡片消息是聊天和 App 服务之间的桥——用户在对话里点一下卡片&#xff0c;直接落到小程序的具体页面。它的工程要点在"卡片参数怎么拼"和"落地页怎么接"。 一、小程序卡片的结构 发送小程序消息的参数组合比链接消息复杂&#xff1a;小程序的 ap…

作者头像 李华
网站建设 2026/9/12 17:38:27

SpringBoot汽车销售管理系统:实时库存与电子合同实践

1. 项目背景与核心价值汽车销售行业正经历从传统线下模式向数字化管理的转型浪潮。去年我参与改造某4S店管理系统时&#xff0c;亲眼目睹了手工台账导致的库存混乱——销售员A刚签下一台宝马3系的订单&#xff0c;销售员B却在同一时间把同一辆车卖给了另一位客户。这种尴尬局面…

作者头像 李华