【问题标题】:Start multiple threads with Tokio使用 Tokio 启动多个线程
【发布时间】:2018-10-29 05:50:29
【问题描述】:

我正在尝试创建一个基本的 tcp 服务器:

  1. 服务器应该能够向所有连接的客户端广播消息流
  2. 服务器应该能够接收来自所有客户端的命令并处理它们

这是我在 main 函数中得到的:

let (server_tx, server_rx) = mpsc::unbounded();
let state = Arc::new(Mutex::new(Shared::new(server_tx)));

let addr = "127.0.0.1:6142".parse().unwrap();

let listener = TcpListener::bind(&addr).unwrap();

let server = listener.incoming().for_each(move |socket| {
    // Spawn a task to process the connection
    process(socket, state.clone());
    Ok(())
}).map_err(|err| {
    println!("accept error = {:?}", err);
});

println!("server running on localhost:6142");

let _messages = server_rx.for_each(|_| {
    // process messages here
    Ok(())
}).map_err(|err| {
    println!("message error = {:?}", err);
});

tokio::run(server);  

(playground)

我使用来自 tokio 存储库的 chat.rs 示例作为基础。
我正在通过传入的 tcp 消息向 server_tx 发送数据。
我遇到的问题是消费它们。
我正在使用server_rx.for_each(|_| {“消费”传入的消息流,现在,我如何告诉 tokio 运行它?

tokio::run 接受一个未来,但我有 2 个(可能更多)。如何组合它们以使它们并行运行?

【问题讨论】:

    标签: rust rust-tokio


    【解决方案1】:

    一起加入未来:

    let messages = server_rx.for_each(|_| {
        println!("Message broadcasted");
        Ok(())
    }).map_err(|err| {
        println!("accept error = {:?}", err);
    });
    
    tokio::run(server.join(messages).map(|_| ()));
    

    需要map() 组合器,因为Join Item 关联类型是一个元组((), ())tokio::run() 使用需要 Future::Item 类型为 () 的未来任务

    【讨论】:

      猜你喜欢
      • 2022-09-29
      • 1970-01-01
      • 2014-06-16
      • 2023-03-06
      • 2022-11-03
      • 2017-03-27
      • 2022-12-04
      • 1970-01-01
      • 2010-11-22
      相关资源
      最近更新 更多