使用notify-rs在Rust线程间共享状态的问题求助
Rust文件传感器计数同步问题
我是Rust新手,正在编写一个file_sensor:文件创建后启动计数器,如果一段时间内未收到第二个文件,传感器将以0退出码终止。下方代码能体现当前问题(省略了提及的post函数)。我已为此困扰数小时,试过Arc、Mutex甚至全局变量都没解决。当前用Ticktock-rs实现Timer,需要要么在EventKind::Create(CreateKind::File)的匹配代码块中获取heartbeat,要么让循环里的file_count显示正确值。目前代码可运行,但循环中的file_count始终为0。
相关代码:
use std::env; use std::path::Path; use std::{thread, time}; use std::process::ExitCode; use ticktock::Timer; use notify::{ Watcher, RecommendedWatcher, RecursiveMode, Result, event::{EventKind, CreateKind, ModifyKind, Event} }; // 补充原代码未定义的Args结构体 #[derive(Debug)] struct Args; impl Args { fn parse() -> Self { Args } } fn main() -> Result<()> { let now = time::Instant::now(); let mut heartbeat = Timer::apply( |_, count| { *count += 1; *count }, 0, ) .every(time::Duration::from_millis(500)) .start(now); let mut file_count = 0; let args = Args::parse(); let REQUEST_SENSOR_PATH = env::var("REQUEST_SENSOR_PATH").expect("$REQUEST_SENSOR_PATH is not set"); let mut watcher = notify::recommended_watcher(move|res: Result<Event>| { match res { Ok(event) => { match event.kind { EventKind::Create(CreateKind::File) => { file_count += 1; // 处理文件逻辑 } _ => { /* 处理其他变更 */ } } println!("{:?}", event); }, Err(e) => { println!("watch error: {:?}", e); ExitCode::from(101); }, } })?; watcher.watch(Path::new(&REQUEST_SENSOR_PATH), RecursiveMode::Recursive)?; loop { let now = time::Instant::now(); if let Some(n) = heartbeat.update(now){ println!("Heartbeat: {}, fileCount: {}", n, file_count); if n > 10 { heartbeat.set_value(0); // 文件到达时重置计时器的逻辑 } } } Ok(()) }
问题原因与解决方案
核心问题
notify的watcher闭包使用了move关键字,会将file_count的所有权转移到闭包所在的后台线程,主循环中的file_count是完全独立的另一个变量,因此永远显示初始值0。同时heartbeat也无法在闭包中直接修改,因为它只存在于主线程。
解决代码
用Arc<Mutex<T>>实现线程安全的共享可变状态,让主线程和watcher线程能安全访问同一变量:
use std::env; use std::path::Path; use std::sync::{Arc, Mutex}; use std::time; use std::process::exit; use ticktock::Timer; use notify::{ Watcher, RecommendedWatcher, RecursiveMode, Result, event::{EventKind, CreateKind, Event} }; #[derive(Debug)] struct Args; impl Args { fn parse() -> Self { Args } } fn main() -> Result<()> { let now = time::Instant::now(); // 用Arc<Mutex>包装Timer,支持多线程共享修改 let heartbeat = Arc::new(Mutex::new( Timer::apply( |_, count| { *count += 1; *count }, 0, ) .every(time::Duration::from_millis(500)) .start(now) )); // 同样包装file_count let file_count = Arc::new(Mutex::new(0)); let REQUEST_SENSOR_PATH = env::var("REQUEST_SENSOR_PATH").expect("$REQUEST_SENSOR_PATH is not set"); // 克隆Arc指针传给闭包(仅复制指针,底层数据共享) let hb_clone = Arc::clone(&heartbeat); let fc_clone = Arc::clone(&file_count); let mut watcher = notify::recommended_watcher(move|res: Result<Event>| { match res { Ok(event) => { match event.kind { EventKind::Create(CreateKind::File) => { // 锁住Mutex修改file_count let mut fc = fc_clone.lock().unwrap(); *fc += 1; // 重置计时器 let mut hb = hb_clone.lock().unwrap(); hb.set_value(0); } _ => { /* 处理其他事件 */ } } println!("{:?}", event); }, Err(e) => { println!("watch error: {:?}", e); exit(101); }, } })?; watcher.watch(Path::new(&REQUEST_SENSOR_PATH), RecursiveMode::Recursive)?; loop { let now = time::Instant::now(); // 锁住heartbeat更新并读取 let mut hb = heartbeat.lock().unwrap(); if let Some(n) = hb.update(now){ // 读取file_count let fc = file_count.lock().unwrap(); println!("Heartbeat: {}, fileCount: {}", n, *fc); // 达到阈值时退出程序 if n > 10 { println!("长时间无新文件,退出传感器"); exit(0); } } // 让出CPU时间片,避免空转占用资源 std::thread::sleep(time::Duration::from_millis(100)); } }
关键说明
Arc实现多线程共享所有权,Mutex保证同一时间只有一个线程能访问内部变量,避免数据竞争。- 闭包中使用
Arc::clone仅复制指针,不会复制底层数据,性能开销极小。 - 原代码中
ExitCode::from(101)在闭包里无效,因为闭包运行在独立线程,需用exit()直接终止程序。 - 循环中添加
sleep减少CPU占用,避免空转。
内容的提问来源于stack exchange,提问作者Happy Machine
相关产品推荐
相关产品推荐

