You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 03:55:28