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

Rust多线程请求乱序:无sleep时批量执行读写的原因问询

问题:Rust多线程请求处理出现批量执行而非按顺序执行的现象

我构建了一个基础的Rust多线程应用,主线程随机生成并发送读写请求,由独立的读线程、写线程分别处理。代码可正常运行,但在process_result的match分支后不添加500ms sleep时,原本随机生成的请求会出现批量执行读操作后批量执行写操作(或反之)的情况;添加sleep后则能按请求队列顺序执行。想理解该现象的原因,猜测可能与unbounded无界通道的处理顺序有关。

代码实现

use crossbeam_channel::unbounded; // Import from crossbeam_channel
use rand::prelude::*;
use std::sync::{mpsc, Arc, Mutex};
use std::thread::scope;
use std::time::Duration;

#[derive(Debug, Default)]
enum OpType {
    #[default]
    Read,
    Write(u8),
}

impl OpType {
    fn rand_write() -> Self {
        let mut rng = rand::thread_rng();
        OpType::Write(rng.gen_range(0..100))
    }
}

pub fn test_thread3() {
    let data = Arc::new(Mutex::new(0u8));

    let (read_tx, read_rx) = unbounded::<OpType>();
    let (write_tx, write_rx) = unbounded::<OpType>();
    let (conn_tx, conn_rx) = mpsc::channel::<OpType>();

    // Read thread process
    let process_read = || {
        println!("Starting Read Thread");
        let data = data.clone();
        while let Ok(OpType::Read) = read_rx.recv_timeout(Duration::from_millis(100)) {
            if let Ok(guard) = data.lock() {
                println!("Data = {:?}", *guard);
            }
        }
    };

    //write thread process
    let process_write = || {
        let data = data.clone();
        println!("Starting Write Thread");
        while let Ok(OpType::Write(b)) = write_rx.recv_timeout(Duration::from_millis(100)) {
            if let Ok(mut guard) = data.lock() {
                println!("Writing {:?}", b);
                *guard = b;
            }
        }
    };

    // Main thread process
    let process_result = move || {
        println!("Main Thread Started");
        while let Ok(req) = conn_rx.recv_timeout(Duration::from_millis(50)) {
            match req {
                OpType::Read => {
                    let res = read_tx.send(req);
                    if res.is_err() {
                        println!("Panicked when reading");
                    }
                }
                OpType::Write(_) => {
                    let res = write_tx.send(req);
                    if res.is_err() {
                        println!("Panicked when writing");
                    }
                }
            }
        }
    };

    let mut user_request: Vec<OpType> = Vec::new();
    for _ in 0..=2 {
        user_request.push(OpType::default());
        user_request.push(OpType::rand_write());
    }
    user_request.shuffle(&mut rand::thread_rng());
    println!("{:?}", &user_request);

    for connection_request in user_request {
        conn_tx
            .send(connection_request)
            .expect("Error Sending Request To The ConnectionPool");
    }

    _ = scope(|s| {
        s.spawn(process_result);
        s.spawn(process_read);
        s.spawn(process_write);
    });
}

fn main() {
    test_thread3();
}

示例结果

当前无sleep时的执行结果(批量处理同类型请求):

[Write(92), Read, Write(52), Write(47), Read, Read]
Main Thread Started
Starting Write Thread
Starting Read Thread
Writing 92
Writing 52
Writing 47
Data = 47
Data = 47
Data = 47

预期结果:请求按队列顺序交替执行(如先写92,再读,再写52,再写47,再读,再读)


原因分析

  1. 无界通道的快速缓存:process_result线程处理conn_rx请求的速度极快,无界通道(unbounded)会瞬间将所有同类型请求缓存起来。比如当队列里连续几个写请求时,write_tx会一次性把这些请求全部发送到通道,写线程会连续处理完所有缓存的写请求,期间读线程没有任务可处理。
  2. 线程调度的特性:操作系统的线程调度器会倾向于让正在执行任务的线程持续运行,直到它进入等待状态(比如通道为空)。当写线程开始处理批量写请求时,调度器会优先分配CPU时间给它,直到通道里的写请求全部处理完毕,才会切换到读线程处理读请求。
  3. sleep的作用:添加sleep后,process_result发送请求的速度被放慢,给了读/写线程足够的时间交替处理请求。每次发送一个请求后,主线程休眠,调度器会切换到对应的处理线程执行,因此看起来是按队列顺序执行。

解决方案

如果需要让请求严格按发送顺序执行,不能依赖sleep或线程调度的不确定性,可采用以下方式:

方式1:使用线程让步替代固定sleep

在发送每个请求后调用std::thread::yield_now(),主动让调度器切换到其他线程,这样无需固定休眠时间,更高效:

// 修改process_result中的match分支
match req {
    OpType::Read => {
        let res = read_tx.send(req);
        if res.is_err() {
            println!("Panicked when reading");
        }
        std::thread::yield_now(); // 主动让出CPU
    }
    OpType::Write(_) => {
        let res = write_tx.send(req);
        if res.is_err() {
            println!("Panicked when writing");
        }
        std::thread::yield_now(); // 主动让出CPU
    }
}

方式2:保证请求的串行处理(严格顺序)

如果必须严格按请求顺序执行,可以让一个线程处理所有请求,或者通过同步机制控制请求的处理顺序。比如只用一个通道,让单个线程按顺序处理读写请求:

// 合并读写线程为一个处理线程
let process_all = || {
    println!("Starting Processing Thread");
    let data = data.clone();
    while let Ok(req) = all_rx.recv_timeout(Duration::from_millis(100)) {
        match req {
            OpType::Read => {
                if let Ok(guard) = data.lock() {
                    println!("Data = {:?}", *guard);
                }
            }
            OpType::Write(b) => {
                if let Ok(mut guard) = data.lock() {
                    println!("Writing {:?}", b);
                    *guard = b;
                }
            }
        }
    }
};

这种方式会失去多线程并行的优势,但能严格保证执行顺序。


内容的提问来源于stack exchange,提问作者Apurba Pokharel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:44:51