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

如何识别并行异步Future以避免重复执行FFI事件循环?

处理FFI异步Future的并行执行优化问题

问题背景

我实现了一个持续检查FFI对象状态并返回Poll结果的Future,但只有调用一次FFI事件循环后,该对象的状态才会发生变化。现有代码如下:

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
   run_event_loop_once();

   if get_ffi_object().check_status_for_id(self.id) {
     return Poll::Ready(Ok(()));
   }

   return Poll::Pending;
}

但当用户并行执行两个这类Future时,异步执行器会在每个Future的poll方法中各运行一次事件循环,这不仅毫无意义,还会消耗IO循环性能。我希望在tokio/async-std的IO循环中只执行一次该事件循环,同时无法修改第三方提供的FFI对象和事件循环。

预期实现逻辑如下:

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
   
   if is_first_parallel_future() {
      run_event_loop_once();
   }

   if get_ffi_object().check_status_for_id(self.id) {
     return Poll::Ready(Ok(()));
   }

   return Poll::Pending;
}

核心痛点:如何识别并行执行的Future中的"第一个"?普通计数器无法满足需求,因为每个Future都会递增计数,无法判断是否是当前轮次的第一个。

解决方案

方案1:线程本地存储标记执行状态

利用异步执行器的任务调度特性,同一线程内的任务会被依次poll。通过线程本地存储(TLS)标记当前线程是否已在本次调度轮次中执行过事件循环:

use std::cell::Cell;

thread_local! {
    static EVENT_LOOP_RUN_THIS_ROUND: Cell<bool> = Cell::new(false);
}

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
    EVENT_LOOP_RUN_THIS_ROUND.with(|flag| {
        if !flag.get() {
            run_event_loop_once();
            flag.set(true);
        }
    });

    let ready = get_ffi_object().check_status_for_id(self.id);
    if ready {
        Poll::Ready(Ok(()))
    } else {
        EVENT_LOOP_RUN_THIS_ROUND.with(|flag| flag.set(false));
        cx.waker().wake_by_ref();
        Poll::Pending
    }
}

方案2:全局原子计数器控制

通过原子计数器跟踪当前活跃的Pending任务数量,只有当计数器从0变为1时,才执行事件循环:

use std::sync::atomic::{AtomicUsize, Ordering};

static ACTIVE_TASKS: AtomicUsize = AtomicUsize::new(0);

fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
    let prev_count = ACTIVE_TASKS.fetch_add(1, Ordering::SeqCst);
    
    if prev_count == 0 {
        run_event_loop_once();
    }

    let ready = get_ffi_object().check_status_for_id(self.id);
    let result = if ready {
        Poll::Ready(Ok(()))
    } else {
        cx.waker().wake_by_ref();
        Poll::Pending
    };

    ACTIVE_TASKS.fetch_sub(1, Ordering::SeqCst);
    result
}

方案3:后台任务独立执行事件循环(推荐)

创建长期运行的后台任务专门负责执行FFI事件循环,所有业务Future只负责检查状态并等待通知,从根本上避免重复执行问题:

use tokio::sync::broadcast;
use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};

lazy_static::lazy_static! {
    static ref STATUS_UPDATED: broadcast::Sender<()> = {
        let (tx, _) = broadcast::channel(1);
        tokio::spawn(async {
            loop {
                run_event_loop_once();
                let _ = STATUS_UPDATED.send(());
                tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
            }
        });
        tx
    };
}

struct FfiFuture {
    id: u64,
    receiver: broadcast::Receiver<()>,
}

impl FfiFuture {
    fn new(id: u64) -> Self {
        Self {
            id,
            receiver: STATUS_UPDATED.subscribe(),
        }
    }
}

impl Future for FfiFuture {
    type Output = Result<(), ()>;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        if get_ffi_object().check_status_for_id(self.id) {
            return Poll::Ready(Ok(()));
        }

        let this = self.get_mut();
        match this.receiver.poll_recv(cx) {
            Poll::Ready(_) => {
                if get_ffi_object().check_status_for_id(this.id) {
                    Poll::Ready(Ok(()))
                } else {
                    cx.waker().wake_by_ref();
                    Poll::Pending
                }
            }
            Poll::Pending => Poll::Pending,
        }
    }
}

方案对比

  • 方案1:轻量级实现,适合单线程执行器或同一线程内的并行任务,缺点是跨线程并行任务仍可能重复执行事件循环。
  • 方案2:全局控制逻辑,跨线程场景也能保证同一时间只有一次事件循环执行,缺点是原子操作存在轻微性能开销。
  • 方案3:解耦度最高,彻底分离事件循环与业务Future,性能最优,适合长期运行的异步场景,需要依赖执行器的后台任务支持。

内容的提问来源于stack exchange,提问作者Ramesh Kithsiri HettiArachchi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 21:36:25