【问题标题】:How to use Rust shiplift from a hyper server如何从超级服务器使用 Rust 升船
【发布时间】:2019-05-20 18:31:36
【问题描述】:

我正在尝试编写一个简单的 Rust 程序,它使用 shiplift 读取 Docker 统计数据,并使用 rust-prometheus 将它们公开为 Prometheus 指标。

shiplift stats 示例自行正确运行,我正在尝试将其集成到服务器中

fn handle(_req: Request<Body>) -> Response<Body> {
    let docker = Docker::new();
    let containers = docker.containers();
    let id = "my-id";
    let stats = containers
        .get(&id)
        .stats().take(1).wait();
    for s in stats {
        println!("{:?}", s);
    }
    // ...
}

// in main
let make_service = || {
    service_fn_ok(handle)
};

let server = Server::bind(&addr)
    .serve(make_service);

但似乎流永远挂起(我无法产生任何错误消息)。

我还在shiplift example 中尝试了相同的重构(使用takewait 而不是tokio::run),但在这种情况下我收到错误executor failed to spawn task: tokio::spawn failed (is a tokio runtime running this future?)shiplift 是否需要 tokio

编辑: 如果我理解正确的话,我的尝试是行不通的,因为wait 将阻止tokio 执行程序,而stats 将永远不会产生结果。

【问题讨论】:

    标签: rust rust-tokio


    【解决方案1】:

    shiplift 的 API 是异步的,这意味着 wait() 和其他函数返回 Future,而不是阻塞主线程直到结果准备好。 Future 在传递给执行程序之前实际上不会执行任何 I/O。您需要将Future 传递给tokio::run,如您链接到的示例中所示。您应该阅读tokio docs 以更好地了解如何在 rust 中编写异步代码。

    【讨论】:

    • 看起来我无法从hyper 服务调用tokio::run,它以attempted to run an executor while another executor is already running 失败。
    • 尝试hyper::rt::spawn 而不是tokio::runhyper 初始化它自己的执行程序,但 tokio::run 尝试初始化一个新的执行程序,这就是您看到错误消息的原因。 spawn 将在现有 executor 上运行 future。
    • hyper::rt::spawn 可能是正确的方法。另外,我已经意识到我的错误,我发布的代码示例完全错过了。我会更新它以使其更清晰
    【解决方案2】:

    我对@9​​87654321@ 工作原理的理解有很多错误。基本上:

    • 如果一个服务应该处理future,不要使用service_fn_ok来创建它(它用于同步服务):使用service_fn
    • 不要使用wait:所有期货都使用同一个执行程序,执行将永远挂起(文档中有警告,但是哦……);
    • 正如 ecstaticm0rse 所指出的,hyper::rt::spawn 可用于异步读取统计信息,而不是在服务中执行此操作

    tokio 是否需要某种方式的升船?

    是的。它使用hyper,如果默认的tokio 执行器不可用,则抛出executor failed to spawn task(无论如何处理期货几乎总是需要执行器)。

    这是我最终得到的最小版本(tokio 0.1.20 和 hyper 0.12):

    use std::net::SocketAddr;
    use std::time::{Duration, Instant};
    
    use tokio::prelude::*;
    use tokio::timer::Interval;
    
    use hyper::{
        Body, Response, service::service_fn_ok,
        Server, rt::{spawn, run}
    };
    
    fn init_background_task(swarm_name: String) -> impl Future<Item = (), Error = ()> {
        Interval::new(Instant::now(), Duration::from_secs(1))
            .map_err(|e| panic!(e))
            .for_each(move |_instant| {
                futures::future::ok(())  // unimplemented: call shiplift here
            })
    }
    
    fn init_server(address: SocketAddr) -> impl Future<Item = (), Error = ()> {
        let service = move || {
            service_fn_ok(|_request| Response::new(Body::from("unimplemented")))
        };
        Server::bind(&address)
            .serve(service)
            .map_err(|e| panic!("Server error: {}", e))
    }
    
    
    fn main() {
        let background_task = init_background_task("swarm_name".to_string());
        let server = init_server(([127, 0, 0, 1], 9898).into());
    
        run(hyper::rt::lazy(move || {
            spawn(background_task);
            spawn(server);
            Ok(())
        }));
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2012-02-29
      • 2016-07-17
      • 1970-01-01
      • 1970-01-01
      • 2014-04-27
      • 2016-05-04
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多