如何识别并行异步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
相关产品推荐
相关产品推荐

