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

Rust异步编程中如何实现类似Boost.Asio poll的非阻塞任务推进

在Tokio单线程环境下实现类似Boost.Asio poll()的非阻塞任务推进

Tokio没有直接提供和Boost.Asio poll()完全等价的API,但我们可以通过两种方式实现非阻塞推进所有就绪异步任务的效果,同时保持单线程运行,满足你的延迟要求。

方法一:手动轮询LocalSet(精确控制)

这种方式直接操作LocalSet的poll方法,手动推进所有就绪任务,完全不阻塞:

use tokio::time::Duration;
use futures::task::{Context, Poll};
use std::task::{Waker, RawWaker, RawWakerVTable};
use std::pin::Pin;

async fn some_async_function() {
    loop {
        println!("Before await");
        tokio::time::sleep(Duration::from_secs(2)).await;
        println!("After await");
    }
}

// 生成一个空的Waker(仅用于手动poll上下文,无需实际唤醒)
fn noop_waker() -> Waker {
    unsafe {
        Waker::from_raw(RawWaker::new(
            std::ptr::null(),
            &RawWakerVTable::new(
                |_| noop_raw_waker(),
                |_| {},
                |_| {},
                |_| {},
            ),
        ))
    }
}

fn noop_raw_waker() -> RawWaker {
    RawWaker::new(std::ptr::null(), &RawWakerVTable::new(
        |_| noop_raw_waker(),
        |_| {},
        |_| {},
        |_| {},
    ))
}

#[tokio::main(flavor = "current_thread")]
async fn main() {
    let mut local = tokio::task::LocalSet::new();
    local.spawn_local(some_async_function());
    local.spawn_local(some_async_function());

    let waker = noop_waker();
    let mut cx = Context::from_waker(&waker);

    loop {
        // 非阻塞推进所有就绪的异步任务
        match Pin::new(&mut local).poll(&mut cx) {
            Poll::Ready(_) => break, // 所有任务完成(本例中不会触发,因为任务是无限循环)
            Poll::Pending => {} // 无就绪任务,退出poll执行同步逻辑
        }

        // 在这里执行你的同步逻辑,比如处理消息队列
        println!("Running synchronous processing...");

        // 可选:添加小延迟避免CPU空转,根据延迟要求调整或移除
        std::thread::sleep(Duration::from_millis(50));
    }
}

方法二:使用yield_now简化实现(简洁高效)

利用Tokio调度器的yield_now().await,让出当前执行权,运行所有就绪任务后立即返回,实现类似poll()的效果:

use tokio::time::Duration;
use tokio::runtime::Handle;

async fn some_async_function() {
    loop {
        println!("Before await");
        tokio::time::sleep(Duration::from_secs(2)).await;
        println!("After await");
    }
}

#[tokio::main(flavor = "current_thread")]
async fn main() {
    let local = tokio::task::LocalSet::new();
    local.spawn_local(some_async_function());
    local.spawn_local(some_async_function());

    // 进入LocalSet上下文,确保本地任务能被调度
    let _enter = local.enter();
    let handle = Handle::current();

    loop {
        // 非阻塞推进所有就绪的异步任务
        handle.block_on(async {
            tokio::task::yield_now().await;
        });

        // 执行同步逻辑
        println!("Running synchronous processing...");

        // 可选:避免CPU空转
        std::thread::sleep(Duration::from_millis(50));
    }
}

关于LocalSet的用法

你原来的LocalSet用法是正确的:spawn_local用于创建!Send的本地任务,这类任务必须在LocalSet的上下文中才能被调度执行。两种实现方式都正确处理了LocalSet的上下文,确保任务能被正常推进。

注意事项

  1. 单线程保证:#[tokio::main(flavor = "current_thread")]确保Tokio运行时是单线程的,完全符合你的延迟要求。
  2. CPU占用:如果同步逻辑执行极快,循环会频繁触发导致CPU占用过高,可以添加小的同步延迟;如果延迟要求极高,可以移除延迟,根据实际场景调整。
  3. 定时器任务:Tokio的定时器在单线程环境下会被正确处理,两种方法都会检查定时器是否就绪,推进对应的休眠任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 09:54:57