Rust中批量数据超时发送逻辑异常,如何修正?
问题分析与代码修正
需求说明
从Flume的Receiver<RequestMessage>接收HTTP请求数据,缓存到Vec中,满足以下任一条件时将缓存数据发送至新通道:
- 缓存长度≥3
- 缓存长度为1或2,且超过3秒未收到新数据
原代码核心问题
原代码的超时逻辑完全错误:
- 仅在收到新数据后才触发超时等待,无法处理"长时间无新数据"的场景
- 用
spawn_blocking执行std::thread::sleep(4秒),外层套timeout(3秒),导致is_ok()永远为false,超时发送逻辑根本不会触发 - 未处理通道关闭后剩余缓存数据的发送
修正后的代码
use flume::{Receiver, Sender}; use std::sync::Arc; use tokio::sync::RwLock; use tokio::time::{sleep, Duration}; use tokenizers::Tokenizer; // 假设定义了以下类型(根据实际场景调整) #[derive(Debug)] struct RequestMessage { message: String, // 其他业务字段... } #[derive(Debug)] struct SessionIdsWithTokens { // 目标通道所需字段... } async fn process_flume1_data( rx: Receiver<RequestMessage>, tx2: Sender<SessionIdsWithTokens>, request_vec: Arc<RwLock<Vec<RequestMessage>>>, mut tokenizer: Tokenizer, ) { // 初始化超时定时器,3秒后触发 let mut timeout_timer = sleep(Duration::from_secs(3)); loop { tokio::select! { // 监听新数据到达 Ok(request_message) = rx.recv_async() => { println!("Received data from Flume channel1: {:?}", request_message.message); let mut guard = request_vec.write().await; guard.push(request_message); // 收到数据后重置超时定时器 timeout_timer = sleep(Duration::from_secs(3)); // 缓存满3条,立即发送 if guard.len() >= 3 { send_to_channel(&tx2, guard, &mut tokenizer).await; } }, // 超时触发(3秒无新数据) _ = &mut timeout_timer => { let mut guard = request_vec.write().await; // 缓存有1-2条数据时发送 if !guard.is_empty() && guard.len() <= 2 { send_to_channel(&tx2, guard, &mut tokenizer).await; } // 重置定时器,继续等待 timeout_timer = sleep(Duration::from_secs(3)); }, // 原通道关闭,退出前发送剩余数据 else => { let mut guard = request_vec.write().await; if !guard.is_empty() { send_to_channel(&tx2, guard, &mut tokenizer).await; } break; } } } } // 封装发送逻辑,避免重复代码 async fn send_to_channel( tx2: &Sender<SessionIdsWithTokens>, mut guard: tokio::sync::RwLockWriteGuard<'_, Vec<RequestMessage>>, tokenizer: &mut Tokenizer, ) { // 这里实现RequestMessage到SessionIdsWithTokens的转换逻辑 // 示例: let session_data = SessionIdsWithTokens { // 填充转换后的数据... }; // 发送到新通道,处理发送失败场景 if let Err(e) = tx2.send_async(session_data).await { eprintln!("Failed to send to channel2: {}", e); } // 发送完成后清空缓存 guard.clear(); }
代码说明
- 核心逻辑:用
tokio::select!同时监听三个事件:新数据到达、超时触发、原通道关闭,确保所有场景都被覆盖 - 超时控制:每次收到新数据时重置定时器,保证超时只会在"连续3秒无新数据"时触发
- 缓存管理:发送完成后立即清空缓存,避免重复发送
- 边界处理:原通道关闭时自动发送剩余缓存数据,防止数据丢失
- 代码复用:将发送逻辑封装为独立函数,减少冗余代码
内容的提问来源于stack exchange,提问作者hoson
相关产品推荐
相关产品推荐

