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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 10:45:47