Node.js基于Worker多线程优化交易路由价格抓取性能求助
用Node.js实现交易路由并发计算提升性能
你的问题核心是把串行执行的10条路由计算改成并行,从而把总耗时从20秒压缩到单条路由的耗时(约2秒)。下面提供两种方案,按需选择:
方案一:轻量并行(Promise.all,推荐IO密集场景)
如果你的calc函数主要是和交易所API交互(IO密集型,大部分时间在等待响应),用Promise.all就能实现并行,不需要Worker线程,代码改动最小:
async function func() { const start_time = performance.now(); // 同时发起所有路由的计算请求 const results = await Promise.all( routes.map(async (route) => { const result_amount = await calc(route, amount_wei); return { route, result_amount }; }) ); // 批量处理结果 results.forEach(({ route, result_amount }) => { if (result_amount[5] > amount_start * 1) { console.log(`Good Trade on route: ${route[0]}`); } }); console.log(`Total execution time: ${performance.now() - start_time}ms`); } async function main() { while (true) { await func(); // 可选:添加间隔避免频繁请求触发交易所限流 // await new Promise(resolve => setTimeout(resolve, 1000)); } } main();
方案二:Worker线程(适合CPU密集/避免阻塞主线程)
如果calc函数包含大量同步计算(CPU密集型),或者你担心并行请求阻塞主线程,可以用Node.js内置的worker_threads模块实现多线程并行:
第一步:创建Worker脚本(trade-worker.js)
把单条路由的计算逻辑独立到Worker中:
const { parentPort } = require('worker_threads'); // 这里复制你的calc函数实现,或者从公共模块导入 async function calc(route, amount_wei) { // 原有的QuickSwap + SushiSwap兑换逻辑 } // 监听主线程的任务请求 parentPort.on('message', async ({ route, amount_wei }) => { try { const result_amount = await calc(route, amount_wei); parentPort.postMessage({ success: true, route, result_amount }); } catch (err) { parentPort.postMessage({ success: false, route, error: err.message }); } });
第二步:修改主线程代码
创建多个Worker并行处理路由:
const { Worker } = require('worker_threads'); const path = require('path'); async function func() { const start_time = performance.now(); // 为每条路由创建Worker任务 const workerTasks = routes.map(route => { return new Promise((resolve, reject) => { const worker = new Worker(path.resolve(__dirname, 'trade-worker.js')); // 发送任务数据给Worker worker.postMessage({ route, amount_wei }); // 接收Worker返回结果 worker.on('message', msg => { worker.terminate(); msg.success ? resolve(msg) : reject(new Error(msg.error)); }); // 处理Worker错误 worker.on('error', err => { worker.terminate(); reject(err); }); }); }); try { const results = await Promise.all(workerTasks); results.forEach(msg => { if (msg.result_amount[5] > amount_start * 1) { console.log(`Good Trade on route: ${msg.route[0]}`); } }); console.log(`Total execution time: ${performance.now() - start_time}ms`); } catch (err) { console.error('Route processing failed:', err); } } async function main() { while (true) { await func(); // 可选:添加请求间隔 // await new Promise(resolve => setTimeout(resolve, 1000)); } } main();
注意事项
- 交易所API通常有请求频率限制,并行请求时要确保不超过限制,必要时添加请求间隔
- 如果用Worker线程,确保
calc依赖的库(比如Web3)能在Worker环境中正常运行 - 若需要频繁处理任务,可以实现Worker池复用线程,减少创建销毁的开销
内容的提问来源于stack exchange,提问作者Dani S
相关产品推荐
相关产品推荐

