Rust应用中reqwest POST请求随机冻结问题排查求助
问题:Rust日志监控应用随机冻结在reqwest的send()方法处
我接触Rust开发仅2周,开发了一个监控日志文件并将批量数据发送至Elasticsearch数据库的应用。运行一段时间后,应用会随机冻结(CPU占用100%),且无错误提示。经CLion调试定位,冻结发生在reqwest请求的send()方法处。
日志文件每分钟新增约625行,冻结通常在处理5500-25000行后随机发生。怀疑问题与LogWatcher、reqwest、block_on或异步代码混合使用有关。
相关代码:
use std::{env}; use std::io::{stdout, Write}; use std::path::Path; use std::time::Duration; use logwatcher::{LogWatcher, LogWatcherAction}; use serde_json::{json, Value}; use serde_json::Value::Null; use tokio; #[tokio::main] async fn main() { let mut log_watcher = LogWatcher::register("/var/log/test.log").unwrap(); let mut counter = 0; let BULK_SIZE = 500; log_watcher.watch(&mut move |line: String| { // 日志文件新增行时触发 counter += 1; if counter >= BULK_SIZE { futures::executor::block_on(async { // 因LogWatcher非异步需同步阻塞执行 let _response = reqwest::Client::new() .post("http://127.0.0.1/test.php") // 测试用地址,连DB时也会出现问题 .header("Content-Type", "application/json") .body("{\"test\": true}") .timeout(Duration::from_secs(30)) .send() // 冻结位置 .await; if _response.is_ok(){ println!("Ok"); } }); counter = 0; } LogWatcherAction::None }); }
问题分析
- 核心问题:同步回调中阻塞异步运行时
LogWatcher的watch方法运行在同步线程中,回调触发时用futures::executor::block_on强制阻塞执行异步HTTP请求,会抢占Tokio运行时的线程资源。当日志高频新增、多次触发批量发送时,会导致运行时线程池饥饿甚至死锁,最终引发CPU占用100%、应用冻结。 - 额外隐患:重复创建reqwest Client
每次批量发送都新建reqwest::Client会浪费连接池资源,频繁创建销毁连接也可能加剧网络层面的不稳定。
修复方案
1. 用异步通道解耦同步回调与异步任务
通过Tokio的MPSC通道,让同步回调只负责提交任务,实际的HTTP请求由后台异步任务处理,彻底避免阻塞:
use std::{time::Duration}; use logwatcher::{LogWatcher, LogWatcherAction}; use tokio::sync::mpsc; use reqwest::Client; #[tokio::main] async fn main() { // 创建异步通道,缓冲队列设为100,避免突发日志堆积 let (tx, mut rx) = mpsc::channel(100); // 启动后台异步任务,专门处理批量发送逻辑 tokio::spawn(async move { // 仅初始化一次reqwest Client,复用连接池 let client = Client::new(); let target_url = "http://127.0.0.1/test.php"; // 监听通道,收到信号就执行发送 while let Some(_) = rx.recv().await { let response = client .post(target_url) .header("Content-Type", "application/json") .body("{\"test\": true}") .timeout(Duration::from_secs(30)) .send() .await; match response { Ok(_) => println!("批量发送成功"), Err(e) => eprintln!("发送失败: {}", e), } } }); let mut log_watcher = LogWatcher::register("/var/log/test.log").unwrap(); let mut counter = 0; let BULK_SIZE = 500; log_watcher.watch(&mut move |_line: String| { counter += 1; if counter >= BULK_SIZE { // 在同步回调中用blocking_send提交任务,不阻塞回调线程 if let Err(e) = tx.blocking_send(()) { eprintln!("任务提交失败: {}", e); } counter = 0; } LogWatcherAction::None }); }
2. 关键优化说明
- 通道解耦:同步回调只做计数和任务提交,不处理耗时的网络请求,避免阻塞LogWatcher的工作线程和Tokio运行时。
- 复用Client:
reqwest::Client自带连接池,复用能大幅提升网络请求效率,减少资源消耗。 - 错误处理:新增了通道发送和请求失败的错误打印,方便后续排查问题。
3. 扩展建议
如果需要发送实际的日志内容,可修改通道传递批量日志数据,示例如下:
// 修改通道类型为Vec<String> let (tx, mut rx) = mpsc::channel(100); // 回调中缓冲日志行 let mut buffer = Vec::with_capacity(BULK_SIZE); log_watcher.watch(&mut move |line: String| { buffer.push(line); if buffer.len() >= BULK_SIZE { // 发送缓冲的日志数据 let _ = tx.blocking_send(buffer.drain(..).collect()); } LogWatcherAction::None }); // 后台任务处理批量数据 while let Some(lines) = rx.recv().await { // 将lines转为Elasticsearch批量格式 let bulk_body = lines.into_iter() .map(|line| json!({"index": {}}).to_string() + "\n" + &line + "\n") .collect::<String>(); let response = client .post(target_url) .header("Content-Type", "application/x-ndjson") .body(bulk_body) .send() .await; // ... 错误处理 }
内容的提问来源于stack exchange,提问作者Typewar
相关产品推荐
相关产品推荐

