JavaScript多线程
JavaScript策略可以用threading对象创建真正并行执行的线程,并用消息、共享字典、锁等对象在线程之间通信。本页说明怎么选用、怎么组织线程代码;各函数的参数和返回值见语法手册Threads。
先选对工具
| 需求 | 推荐 | 适用语言 |
|---|---|---|
| 同时发出几个API请求(如同时取多个交易所的行情),等结果回来 | exchange.Go,配合EventLoop等待完成事件 | 所有语言 |
| 长时间在后台运行的任务:独立的行情采集、风控巡检、耗时计算 | threading.Thread | 仅JavaScript |
| 在策略内提供HTTP、WebSocket或TCP服务 | threading.Serve | 仅JavaScript |
只是并发几个请求时,exchange.Go()更简单,也不需要处理线程间的数据传递。本页的threading对象只适用于JavaScript策略,Python、Rust策略请使用exchange.Go()。
回测系统中可以调用这些函数,但线程实际是按顺序执行的,只用于保证代码在回测中能运行。
线程运行在隔离的环境中
传给threading.Thread()的函数在一个独立的JavaScript环境中执行,这是写线程代码时最需要注意的一点:
- 线程函数不能引用外部的变量和闭包,也不能调用策略里自定义的其它函数。需要的数据通过
threading.Thread(func, arg1, arg2, ...)的参数传入。 - 普通对象、数组作为参数时是深拷贝:线程里修改它不影响其它线程。需要多个线程看到同一份数据时,使用
threading.Dict()创建的字典。 - 函数也可以作为参数传入;
threading.Thread()还支持传入函数源码字符串,用于在线程中加载外部库。 - 线程里可以直接调用平台的API函数,如
exchange.GetTicker()、Log()。 - 线程函数的返回值通过
join()取回:t.join().ret。
线程之间怎样交换数据
| 方式 | 用法 | 说明 |
|---|---|---|
| 消息 | t.postMessage(msg)发给线程t;线程内用threading.currentThread().peekMessage(timeout)读取自己收到的消息;子线程用threading.mainThread().postMessage(msg)发回主线程 | 每个线程有自己的收件箱,按顺序读取。peekMessage(-1)不阻塞,没有消息时返回空值 |
| 共享字典 | var d = threading.Dict(),作为参数传入线程后各线程d.get(key)、d.set(key, value) | 适合保存「最新状态」,例如最新行情、运行标志 |
| 线程数据 | t.setData(key, value)、t.getData(key) | 挂在某个线程对象上的键值,线程结束(join()、terminate())后失效 |
| 同步对象 | threading.Lock()、threading.Event()、threading.Condition() | 作为参数传入线程,用于互斥访问和等待通知 |
线程收到消息时也会产生事件,可以用线程对象的eventLoop统一等待消息和其它事件。
线程的生命周期
t.join()等待线程结束并取回返回值,可以设置超时;t.terminate()强制结束线程。- 线程结束且不再被引用时,资源会自动回收,不必为了释放资源调用
join()。持续引用、无法回收的线程累计超过2000个时会报错。 threading.pending()返回正在运行的线程数(包括主线程)。- 实盘停止时所有线程一起结束。
peekMessage()、join()、锁和事件的等待都会被停止打断。
在策略内提供服务
threading.Serve(地址, 处理函数, ...参数)在策略进程内启动HTTP(含WebSocket)或TCP服务,每个请求或连接在独立的线程中调用处理函数,返回Server对象(addr()取实际监听地址,close()关闭)。处理函数与线程函数一样运行在隔离环境中,需要的数据通过参数传入,常用threading.Dict()与主线程共享状态。地址写法、ctx对象的方法见Serve。
旧的全局函数__Serve()仍可使用,它只返回监听地址字符串,新代码请使用threading.Serve()。
示例
示例
-
多个线程并行计算,主线程汇总结果
每个线程拉取一个交易对的K线并计算均线,结果通过返回值交给主线程。注意交易对通过参数传入,线程函数里没有引用外部变量。
javascriptfunction main() { var symbols = ["BTC_USDT", "ETH_USDT", "SOL_USDT"] var threads = [] for (var i = 0; i < symbols.length; i++) { threads.push(threading.Thread(function(symbol, period) { // 在线程中运行:只能使用参数和平台API var records = exchange.GetRecords(symbol, period) if (!records || records.length < 20) { return null } var ma = TA.MA(records, 20) return {symbol: symbol, close: records[records.length - 1].Close, ma20: ma[ma.length - 1]} }, symbols[i], PERIOD_H1)) } for (var i = 0; i < threads.length; i++) { var r = threads[i].join().ret if (r) { Log(r.symbol, "收盘价:", r.close, "MA20:", r.ma20) } } } -
后台线程采集行情,主线程读取与下发指令
后台线程把最新价写进共享字典,并把异常通过消息报告给主线程;主线程通过消息通知后台线程退出。
javascriptfunction main() { var shared = threading.Dict() var worker = threading.Thread(function(dict, symbol) { while (true) { // 读取主线程发来的指令,-1 表示不阻塞 var cmd = threading.currentThread().peekMessage(-1) if (cmd == "stop") { break } var ticker = exchange.GetTicker(symbol) if (ticker) { dict.set("last", ticker.Last) dict.set("time", ticker.Time) } else { threading.mainThread().postMessage("行情获取失败: " + GetLastError()) } Sleep(1000) } return "worker exited" }, shared, "BTC_USDT") for (var i = 0; i < 10; i++) { // 最多等 1 秒后台线程的消息 var msg = threading.currentThread().peekMessage(1000) if (msg) { Log("后台线程报告:", msg) } LogStatus("最新价:", shared.get("last"), "时间:", _D(shared.get("time"))) } worker.postMessage("stop") Log(worker.join().ret) } -
用threading.Serve提供状态查询接口
主线程把状态写进共享字典,HTTP处理函数从参数取到同一个字典并返回JSON。
javascriptfunction main() { var state = threading.Dict() var server = threading.Serve("http://127.0.0.1:8088", function(ctx, st) { if (ctx.path() == "/status") { ctx.setHeader("Content-Type", "application/json") ctx.write(JSON.stringify({last: st.get("last"), updated: st.get("updated")})) } else { ctx.setStatus(404) } }, state) Log("服务地址:", server.addr()) while (true) { var ticker = exchange.GetTicker("BTC_USDT") if (ticker) { state.set("last", ticker.Last) state.set("updated", _D()) } Sleep(3000) } }
参考