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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 18:01:18