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

Rust消费Redis Stream时触发Interrupted system call (os error 4)错误求助

Redis Stream 消费报错 Interrupted system call (os error 4) 解决方法

我在Rust中尝试消费Redis Stream消息时,程序抛出Interrupted system call (os error 4)错误。已确认名为texhub-server:proj:s-comp-queue的Stream存在且不为空,也尝试将redis库升级至0.23.3并按照官方示例调整代码,但问题依旧。相关代码及配置如下:

初始消费代码

use redis::{Commands, PubSubCommands, ControlFlow};

fn main() {
    let redis_client = redis::Client::open("redis://default:123456@123456l:6379/1").expect("get redis client failed");
    let mut con = redis_client
        .get_connection()
        .expect("get redis connection failed");
    let mut count = 0;
    let consume_result = con.subscribe(&["texhub-server:proj:s-comp-queue"], |msg| {
        println!("{:?}", msg);
        count += 1;
        match count {
            10 => ControlFlow::Break(()),
            _ => ControlFlow::Continue,
        }
    });
    if let Err(err) = consume_result {
        println!("{}", err);
    }
}

推送Stream的代码

pub fn push_to_stream(
    stream_key: &str,
    params: &[(&str, &str)],
)  {
    let redis_client = redis::Client::open("redis://default:123456@123456l:6379/1").expect("get redis client failed");
    let mut con = redis_client
        .get_connection()
        .expect("get redis connection failed");
    let result = con.xadd::<&str, &str, &str, &str, String>(stream_key, "*", params);
    if let Err(err) = result {
        println!("{}", err);
    }
}

调整后的消费代码

use redis::{Commands, PubSubCommands, ControlFlow, RedisError};

fn main() {
    let redis_client = redis::Client::open("redis://default:123456@123456l:6379/1").expect("get redis client failed");
    let mut con = redis_client
        .get_connection()
        .expect("get redis connection failed");
    let mut pubsub = con.as_pubsub();
    let sub_result = pubsub.subscribe("texhub.compile_stream_redis_key");
    if let Err(sub_err) = sub_result {
        println!("subscribe error {}", sub_err);
    }
    loop {
        let msg = pubsub.get_message();
        if let Err(e) = msg {
            println!("subscribe message error {}", e);
            return;
        }
        let payload: Result<String, RedisError> = msg.as_ref().unwrap().get_payload();
        if let Err(e) = payload {
            println!("payload message error {}", e);
            return;
        }
        println!(
            "channel '{}': {}",
            msg.unwrap().get_channel_name(),
            payload.unwrap()
        );
    }
}

Cargo.toml配置

[package]
name = "rust-learn"
version = "0.1.0"
edition = "2018"

[dependencies]
rust_wheel = { git = "https://github.com/jiangxiaoqiang/rust_wheel.git", branch = "diesel2.0", features = ["model","common","rwconfig"]}
redis = "0.21.3"

问题原因及解决方法

核心错误:用Pub/Sub API消费Redis Stream

Redis的**Stream(流)和Pub/Sub(发布订阅)**是完全独立的数据结构,操作API不通用。你一直用Pub/Sub的subscribe、get_message方法消费Stream,这是根本错误,必然导致异常。

正确的Stream消费方式

需要使用xread、xreadgroup这类Stream专属命令,以下是修复后的消费代码:

use redis::{Commands, RedisResult};

fn main() -> RedisResult<()> {
    let redis_client = redis::Client::open("redis://default:123456@123456l:6379/1")?;
    let mut con = redis_client.get_connection()?;

    let stream_key = "texhub-server:proj:s-comp-queue";
    let mut last_id = "$"; // "$"表示从最新消息开始消费

    loop {
        // 阻塞读取Stream,每次最多取1条消息
        let messages: Vec<(String, Vec<(String, Vec<(String, String)>)>)> = 
            con.xread_options(&[stream_key], &[last_id], 1, None, true)?;

        for (key, entries) in messages {
            println!("Stream {} 的消息:", key);
            for (id, fields) in entries {
                println!("ID: {}, 内容: {:?}", id, fields);
                last_id = &id; // 更新消费位置,避免重复读取
            }
        }
    }
}

额外注意事项

  1. 版本适配:确保redis库版本(0.23.3)支持xread_options命令,若有差异可参考对应版本的官方文档调整。
  2. 错误处理:示例用?简化错误处理,可根据业务需求替换为expect或自定义错误逻辑。
  3. 消费组(可选):如果需要多消费者协同工作,建议使用xreadgroup创建消费组,实现消息的负载均衡与消费确认机制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 09:43:11