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

Rust基于TCP流实现Pub/Sub多读写端的问题与解决方案

TCP发布订阅模式套接字API开发需求

TCP流默认代表本地套接字与单个远程套接字之间的一对一连接,本次需要开发一套API实现标准Pub/Sub发布订阅效果:

  • 用户向单个本地写套接字写入的数据,可被多个远程读套接字同时读取
  • 支持创建多个绑定到同一个远程写套接字的本地读套接字,所有读套接字都能收到写入端发送的全量数据
  • 读写端在独立进程中创建时,上述逻辑必须保持一致

测试用例

以下测试用例定义了API的预期使用方式与行为:

#[test]
fn single_writer_single_reader() {
    //reader 1 should print "hello world"
    create_reader(|s| {println!("reader 1 {}", s)});
    let mut writer = create_writer();
    writer.write_all(b"hello world").unwrap();
    thread::sleep(Duration::from_millis(1000));
}

#[test]
fn multi_writer_single_reader() {
    //reader 1 should print "hello" twice
    create_reader(|s| {println!("reader 1 {}", s)});
    let mut writer1 = create_writer();
    let mut writer2 = create_writer();
    writer1.write_all(b"hello").unwrap();
    writer2.write_all(b"hello").unwrap();
    thread::sleep(Duration::from_millis(1000));
}

#[test]
fn single_writer_multi_reader() {
    //each reader should print "hello world"
    create_reader(|s| {println!("reader 1 {}", s)});
    create_reader(|s| {println!("reader 2 {}", s)});
    let mut writer = create_writer();
    writer.write_all(b"hello world").unwrap();
    thread::sleep(Duration::from_millis(1000));
}

现有实现逻辑

读端实现

create_reader接口会启动新线程,运行绑定在WRITER_LISTENER_ADDR地址的服务(即writer监听器),持续接收所有传入的写端连接,为每个连接创建独立的读线程。读线程会先等待reader监听器接收当前reader的注册连接,之后持续从对应写端连接读取数据,每收到一条消息就调用传入的回调函数处理,对应实现代码如下:

pub fn create_reader(callback: fn(String)) {
    std::thread::spawn(move || {listen_for_writers_and_make_readers(callback)});
}

fn listen_for_writers_and_make_readers(callback: fn(String)) {
    //listen for remote sockets that are connecting to WRITER_LISTENER_ADDR.
    let writer_listener = TcpListener::bind(WRITER_LISTENER_ADDR).expect("writer_listener: failed to bind");
    println!("writer_listener: started on {}", WRITER_LISTENER_ADDR);
    loop {
        let (stream, _) = writer_listener.accept().unwrap();
        println!("writer_listener: New writer connected");
        let _reader_thread = std::thread::spawn( move || {read_loop(stream, callback)});
    }

}

fn read_loop(stream: TcpStream, callback: fn(String)) {
    let mut incoming_stream = stream; //stream from the writer
    //wait for the reader listener to accept this reader (this is to signal the reader listener that a new reader is created)
    TcpStream::connect(READER_LISTENER_ADDR).expect("reader: failed to connect to reader listener");
    println!("reader: connected to reader listener");
    println!("reader: starting read loop");
    loop {
        let mut buf = [0; 1024];
        match incoming_stream.read(&mut buf) {
            Ok(0) => {
                println!("reader: writer disconnected");
                break;
            },
            Ok(_) => {
                let msg = String::from_utf8(buf.to_vec()).unwrap();
                println!("reader: calling callback");
                callback(msg);
            }
            Err(e) => {
                println!("{}", e);
                break;
            }
        }
    }
}

写端实现

create_writer接口会创建一个TcpStream实例(即writer对象),等待writer监听器接收该连接后,启动新线程运行绑定在READER_LISTENER_ADDR地址的服务(即reader监听器),持续接收所有传入的读端注册连接,每接收一个连接就克隆一次writer实例。原设计假设对原始writer调用write()方法时,数据会同步到所有克隆的writer实例,进而发送给所有关联的reader端,对应实现代码如下:

fn create_writer() -> TcpStream {
    //wait for the writer listener to accept this writer
    let writer = TcpStream::connect(WRITER_LISTENER_ADDR).expect("writer: failed to connect to writer listener");
    println!("writer: connected to the writer listener");
    //create a thread that listens for readers and clones the writer
    let writer_clone = writer.try_clone().unwrap();
    std::thread::spawn( move || {listen_for_readers_and_clone_writer(&writer_clone)});
    return writer
}

fn listen_for_readers_and_clone_writer(writer: &TcpStream) {
    //listen for remote sockets that are connecting to READER_LISTENER_ADDR
    let reader_listener = TcpListener::bind(READER_LISTENER_ADDR).expect("reader_listener: failed to bind");
    println!("reader_listener: started on {}", READER_LISTENER_ADDR);
    let mut writers = vec![];
    loop {
        let (_stream, _) = reader_listener.accept().unwrap();
        println!("reader_listener: New reader connected");
        writers.push(writer.try_clone().unwrap());
    }
}

遇到的问题

当前single_writer_single_reader单写单读测试可正常通过,但其余两个测试均抛出地址占用错误:

  1. 多写单读multi_writer_single_reader测试报错:reader_listener: failed to bind: Os { code: 10048, kind: AddrInUse, message: "Only one usage of each socket address (protocol/network address/port) is normally permitted." }
  2. 单写多读single_writer_multi_reader测试报错:writer_listener: failed to bind: Os { code: 10048, kind: AddrInUse, message: "Only one usage of each socket address (protocol/network address/port) is normally permitted." }

核心疑问:

  • 是否支持多个TcpListener实例绑定到同一个套接字地址?
  • 如果不支持,该发布订阅场景的正确实现方案是什么?

补充调整尝试

曾尝试通过给监听器分配可选地址段的方式修复地址占用问题,代码如下:

#[derive(Clone, Copy)]
struct AddrRange([&'static str; 10]);

static WRITER_LISTENER_ADDR_RANGE : AddrRange = AddrRange([
    "127.0.0.1:80",
    "127.0.0.1:81",
    "127.0.0.1:82",
    "127.0.0.1:83",
    "127.0.0.1:84",
    "127.0.0.1:85",
    "127.0.0.1:86",
    "127.0.0.1:87",
    "127.0.0.1:88",
    "127.0.0.1:89"
]);
static READER_LISTENER_ADDR_RANGE : AddrRange = AddrRange([
    "127.0.0.1:90",
    "127.0.0.1:91",
    "127.0.0.1:92",
    "127.0.0.1:93",
    "127.0.0.1:94",
    "127.0.0.1:95",
    "127.0.0.1:96",
    "127.0.0.1:97",
    "127.0.0.1:98",
    "127.0.0.1:99"
]);

修改后多写单读测试可正常通过,但单写多读测试中,克隆writer实例时并不会如预期创建新的网络连接,新的疑问:如何让writer_listener感知到writer被克隆的事件,进而启动新的read_loop线程处理对应连接?


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 13:27:25