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

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类型需要完全匹配:

  1. 最外层是数组,每个元素对应一个Stream的结果
  2. 每个Stream结果是(Stream名称,消息列表)的元组
  3. 每个消息是(消息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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 10:07:47