如何在单个Deno MainWorker上高效并发执行多脚本调用?
在单个Deno MainWorker上并发执行脚本的问题与优化需求
需求与概念逻辑
我尝试在单个Deno MainWorker上并发执行同一脚本的多次调用,并等待异步脚本的执行结果。概念上期望实现的run_worker函数逻辑如下:
type Tx = Sender<(String, Sender<String>)>; type Rx = Receiver<(String, Sender<String>)>; struct Runner { worker: MainWorker, futures: FuturesUnordered<Pin<Box<dyn Future<Output=(String, Result<Global<Value>, Error>)>>>>, response_futures: FuturesUnordered<Pin<Box<dyn Future<Output=(String, Result<(), SendError<String>>)>>>>, result_senders: HashMap<String, Sender<String>>, } impl Runner { fn new() ... async fn run_worker(&mut self, rx: &mut Rx, main_module: ModuleSpecifier, user_module: ModuleSpecifier) { self.worker.execute_main_module(&main_module).await.unwrap(); self.worker.preload_side_module(&user_module).await.unwrap(); loop { tokio::select! { msg = rx.recv() => { if let Some((id, sender)) = msg { let global = self.worker.js_runtime.execute_script("test", "mod.entry()").unwrap(); self.result_senders.insert(id, sender); self.futures.push(Box::pin(async { let resolved = self.worker.js_runtime.resolve_value(global).await; return (id, resolved); })); } }, script_result = self.futures.next() => { if let Some((id, out)) = script_result { self.response_futures.push(Box::pin(async { let value = deserialize_value(out.unwrap(), &mut self.worker); let res = self.result_senders.remove(&id).unwrap().send(value).await; return (id.clone(), res); })); } }, // also handle response_futures here else => break, } } } }
临时适配方案及存在的问题
由于Worker无法被多次可变借用,上述实现无法运行。我将Worker包装为RefCell,并自定义BorrowingFuture结构及其poll方法来适配:
struct BorrowingFuture { worker: RefCell<MainWorker>, global: Global<Value>, id: String }
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> { match Pin::new(&mut Box::pin(self.worker.borrow_mut().js_runtime.resolve_value(self.global.clone()))).poll(cx) { Poll::Ready(result) => Poll::Ready((self.id.clone(), result)), Poll::Pending => { cx.waker().clone().wake_by_ref(); Poll::Pending } } }
修改后的调用逻辑变为:
self.futures.push(Box::pin(BorrowingFuture{worker: self.worker, global: global.clone(), id: id.clone()}));
响应Future也需做类似修改,但该实现存在以下问题:
- 每次
poll都创建新Future,虽能运行但会带来性能损耗; - 响应Future也存在相同问题,每次
poll调用send逻辑不合理; - 因无法知晓脚本何时完成,每次
poll都调用waker.wake_by_ref,导致Future被高频轮询,性能低下。
注:我当前实现未使用select!,而是通过枚举作为多Future的输出类型,放入单个FuturesUnordered中匹配处理。
疑问
是否有更优的实现方式?或者MainWorker本就不支持此类使用场景?
完整main函数代码
#[tokio::main] async fn main() { let main_module = deno_runtime::deno_core::resolve_url(MAIN_MODULE_SPECIFIER).unwrap(); let user_module = deno_runtime::deno_core::resolve_url(USER_MODULE_SPECIFIER).unwrap(); let (tx, mut rx) = channel(1); let (result_tx, mut result_rx) = channel(1); let handle = thread::spawn(move || { let runtime = tokio::runtime::Builder::new_multi_thread().enable_all().build().unwrap(); let mut runner = Runner::new(); runtime.block_on(runner.run_worker(&mut rx, main_module, user_module)); }); tx.send(("test input".to_string(), result_tx)).await.unwrap(); let result = result_rx.recv().await.unwrap(); println!("result from worker {}", result); handle.join().unwrap(); }
内容的提问来源于stack exchange,提问作者proteus
相关产品推荐
相关产品推荐

