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; // 更新消费位置,避免重复读取 } } } }
额外注意事项
- 版本适配:确保redis库版本(0.23.3)支持
xread_options命令,若有差异可参考对应版本的官方文档调整。 - 错误处理:示例用
?简化错误处理,可根据业务需求替换为expect或自定义错误逻辑。 - 消费组(可选):如果需要多消费者协同工作,建议使用
xreadgroup创建消费组,实现消息的负载均衡与消费确认机制。
内容的提问来源于stack exchange,提问作者Dolphin
相关产品推荐
相关产品推荐

