Rust中Redis Stream响应解析报错:Bulk维度不匹配如何解决?
解决Rust解析Redis Stream响应类型不匹配问题
问题代码
pub fn consume(){ let redis_client = redis::Client::open("redis://default:123456@127.0.0.1:6379/1").expect("get redis client failed"); let mut con = redis_client .get_connection() .expect("get redis connection failed"); let stream_name = "texhub-server:proj:s-comp-queue"; let stream_id = "0"; loop { let result = con.xread(&[stream_name], &[stream_id]); if let Err(e) = result.as_ref() { println!("{}", e); return; } let unwarp: Vec<(String, String)> = result.unwrap(); } }
错误信息
Response was of incompatible type - TypeError: "Bulk response of wrong dimension" (response was [bulk(string-data('"texhub-server:proj:s-comp-queue"'), bulk(bulk(string-data('"1693881638829-0"'), bulk(string-data('"name"'), string-data('"test"'))), bulk(string-data('"1693881702534-0"'), bulk(string-data('"name"'), string-data('"test"'))), bulk(string-data('"1693891547874-0"'), bulk(string-data('"c98f73bc869143d084eb38a0fc38a8e7"'), string-data('""'))), bulk(string-data('"1693892469370-0"'), bulk(string-data('"c98f73bc869143d084eb38a0fc38a8e7"'), string-data('""'))), bulk(string-data('"1693892476263-0"'), bulk(string-data('"c98f73bc869143d084eb38a0fc38a8e7"'), string-data('""'))), bulk(string-data('"1693899506942-0"'), bulk(string-data('"c98f73bc869143d084eb38a0fc38a8e7"'), string-data('""')))))])
解决方案
Redis XREAD命令的响应是多层嵌套结构,对应Rust类型需要完全匹配:
- 最外层是数组,每个元素对应一个Stream的结果
- 每个Stream结果是(Stream名称,消息列表)的元组
- 每个消息是(消息ID,键值对扁平数组)的元组
方法一:手动匹配嵌套类型
将解析类型改为Vec<(String, Vec<(String, Vec<String>)>)>,自行将扁平键值对数组转为可用结构:
pub fn consume(){ let redis_client = redis::Client::open("redis://default:123456@127.0.0.1:6379/1").expect("get redis client failed"); let mut con = redis_client .get_connection() .expect("get redis connection failed"); let stream_name = "texhub-server:proj:s-comp-queue"; let mut stream_id = "0"; loop { let result: redis::RedisResult<Vec<(String, Vec<(String, Vec<String>)>)>> = con.xread(&[stream_name], &[stream_id]); if let Err(e) = result.as_ref() { println!("{}", e); return; } let reply = result.unwrap(); // 遍历处理每个Stream的消息 for (stream, messages) in reply { println!("Stream: {}", stream); for (msg_id, kv_pairs) in messages { println!("Message ID: {}", msg_id); // 将扁平数组转为HashMap let mut kv_map = std::collections::HashMap::new(); for chunk in kv_pairs.chunks(2) { if let [key, value] = chunk { kv_map.insert(key.clone(), value.clone()); } } println!("Content: {:?}", kv_map); // 更新下一次读取的ID,避免重复读取历史消息 stream_id = msg_id; } } } }
方法二:使用redis-rs专用类型(推荐)
redis-rs内置StreamReadReply和StreamEntry类型,自动解析响应结构,无需手动处理嵌套:
use redis::streams::{StreamReadReply, StreamEntry}; pub fn consume(){ let redis_client = redis::Client::open("redis://default:123456@127.0.0.1:6379/1").expect("get redis client failed"); let mut con = redis_client .get_connection() .expect("get redis connection failed"); let stream_name = "texhub-server:proj:s-comp-queue"; let mut stream_id = "0"; loop { let result: redis::RedisResult<StreamReadReply> = con.xread(&[stream_name], &[stream_id]); if let Err(e) = result.as_ref() { println!("{}", e); return; } let reply = result.unwrap(); // 遍历处理每个Stream的消息 for stream in reply.streams { println!("Stream: {}", stream.name); for entry in stream.entries { println!("Message ID: {}", entry.id); // entry.fields直接是HashMap<String, String> println!("Content: {:?}", entry.fields); // 更新下一次读取的ID stream_id = entry.id.as_str(); } } // 无新消息时短暂休眠,避免空循环占用CPU if reply.streams.iter().all(|s| s.entries.is_empty()) { std::thread::sleep(std::time::Duration::from_millis(100)); } } }
注意事项
- 每次读取后必须更新
stream_id为最后一条消息的ID,否则会重复读取历史数据 - 如需阻塞等待新消息,可使用
xread_options并设置block参数,替代轮询逻辑
内容的提问来源于stack exchange,提问作者Dolphin
相关产品推荐
相关产品推荐

