exchange.Go
多线程异步支持函数,可将所有受支持函数的操作转换为异步并发执行。
exchange.Go(method)
exchange.Go(method, ...args)示例
-
exchange.Go()函数的使用范例。判断undefined时需使用typeof(xx) === "undefined",因为null == undefined在 JavaScript 中是成立的。javascriptfunction main(){ // 以下四种操作为并发多线程异步执行,不会耗时,会立即返回 var a = exchange.Go("GetTicker") var b = exchange.Go("GetDepth") var c = exchange.Go("Buy", 1000, 0.1) var d = exchange.Go("GetRecords", PERIOD_H1) // 调用 wait 方法等待异步获取 ticker 的结果 var ticker = a.wait() // 返回深度数据,如果获取失败也有可能返回 null var depth = b.wait() // 返回订单号,限定 1 秒超时,超时返回 undefined,若上次 wait 超时,此对象可继续调用 wait 等待 var orderId = c.wait(1000) if(typeof(orderId) == "undefined") { // 超时,重新获取 orderId = c.wait() } var records = d.wait() }pythondef main(): a = exchange.Go("GetTicker") b = exchange.Go("GetDepth") c = exchange.Go("Buy", 1000, 0.1) d = exchange.Go("GetRecords", PERIOD_H1) ticker, ok = a.wait() depth, ok = b.wait() orderId, ok = c.wait(1000) if ok == False: orderId, ok = c.wait() records, ok = d.wait()rustfn main() { // Rust 中 exchange.Go 为类型化写法:使用 Go:: 方法 token 指定并发函数,无参数传 (),有参数传元组 // 以下四种操作为并发多线程异步执行,不会耗时,会立即返回 let a = exchange.Go(Go::GetTicker, ()); let b = exchange.Go(Go::GetDepth, ()); // Rust 中没有 Buy 的 token,等价于 CreateOrder,第一个参数 "" 表示当前交易对 let c = exchange.Go(Go::CreateOrder, ("", "buy", 1000, 0.1)); let d = exchange.Go(Go::GetRecords, (PERIOD_H1,)); // 调用 wait 方法等待异步获取 ticker 的结果,wait(0) 会阻塞直到并发线程运行完毕(对应 JS 的无参 wait()) let ticker = a.wait(0); // 返回深度数据,如果获取失败也有可能返回 Err let depth = b.wait(0); // 返回订单号,限定 1 秒超时,超时返回 Err,若上次 wait 超时,此对象可继续调用 wait 等待 // 注意:Err 也可能是下单本身失败(与超时无法区分),此时再次 wait 会返回 Err 并记录出错信息 let mut orderId = c.wait(1000); if orderId.is_err() { // 超时,重新获取 orderId = c.wait(0); } let records = d.wait(0); }c++void main() { auto a = exchange.Go("GetTicker"); auto b = exchange.Go("GetDepth"); auto c = exchange.Go("Buy", 1000, 0.1); auto d = exchange.Go("GetRecords", PERIOD_H1); Ticker ticker; Depth depth; Records records; TId orderId; a.wait(ticker); b.wait(depth); if(!c.wait(orderId, 300)) { c.wait(orderId); } d.wait(records); } -
对已释放的并发对象调用其
wait()方法会报错:javascriptfunction main() { var d = exchange.Go("GetRecords", PERIOD_H1) // 等待 K 线数据返回结果 var records = d.wait() // 此处对已经 wait 过且已结束的异步操作再次调用 wait,将返回 null,并记录错误信息 var ret = d.wait() }pythondef main(): d = exchange.Go("GetRecords", PERIOD_H1) records, ok = d.wait() ret, ok = d.wait()rustfn main() { // Rust 中 exchange.Go 采用类型化写法:通过 Go:: 方法 token 指定并发函数 let d = exchange.Go(Go::GetRecords, (PERIOD_H1,)); // 等待 K 线数据返回结果,wait(0) 会阻塞直到运行完毕(对应 JS 的无参 wait()) let records = d.wait(0); // 此处对已经 wait 过且已结束的异步操作再次调用 wait,将返回 Err,并记录错误信息 let ret = d.wait(0); }c++void main() { auto d = exchange.Go("GetRecords", PERIOD_H1); Records records; d.wait(records); Records ret; d.wait(ret); } -
并发获取多个交易所行情:
javascriptfunction main() { while(true) { var beginTS = new Date().getTime() var arrRoutine = [] var arrTicker = [] var arrName = [] for(var i = 0; i < exchanges.length; i++) { arrRoutine.push(exchanges[i].Go("GetTicker")) arrName.push(exchanges[i].GetName()) } for(var i = 0; i < arrRoutine.length; i++) { arrTicker.push(arrRoutine[i].wait()) } var endTS = new Date().getTime() var tbl = { type: "table", title: "行情", cols: ["索引", "名称", "最新成交价"], rows: [] } for(var i = 0; i < arrTicker.length; i++) { tbl.rows.push([i, arrName[i], arrTicker[i].Last]) } LogStatus(_D(), "Total time for concurrent ticker retrieval:", endTS - beginTS, "ms", "\n", "`" + JSON.stringify(tbl) + "`") Sleep(500) } }pythonimport time import json def main(): while True: beginTS = time.time() arrRoutine = [] arrTicker = [] arrName = [] for i in range(len(exchanges)): arrRoutine.append(exchanges[i].Go("GetTicker")) arrName.append(exchanges[i].GetName()) for i in range(len(exchanges)): ticker, ok = arrRoutine[i].wait() arrTicker.append(ticker) endTS = time.time() tbl = { "type": "table", "title": "行情", "cols": ["索引", "名称", "最新成交价"], "rows": [] } for i in range(len(arrTicker)): tbl["rows"].append([i, arrName[i], arrTicker[i]["Last"]]) LogStatus(_D(), "Total time for concurrent ticker retrieval:", endTS - beginTS, "seconds", "\n", "`" + json.dumps(tbl) + "`") Sleep(500)rustfn main() { loop { let beginTS = UnixNano() / 1000000; let mut arrRoutine = Vec::new(); let mut arrTicker = Vec::new(); let mut arrName = Vec::new(); for e in exchanges.iter() { // Rust中exchange.Go为类型化写法,token为Go::GetTicker arrRoutine.push(e.Go(Go::GetTicker, ())); arrName.push(e.GetName()); } // 失败时记None占位,保持与arrName按索引对齐 for r in arrRoutine.iter() { arrTicker.push(r.wait(0).ok()); } let endTS = UnixNano() / 1000000; // Rust无JSON序列化,用format!拼接表格的JSON文本 let mut rows = String::new(); for i in 0..arrTicker.len() { if let Some(ticker) = &arrTicker[i] { if !rows.is_empty() { rows.push(','); } rows += &format!(r#"[{}, "{}", {}]"#, i, arrName[i], ticker.Last); } } let tbl = format!(r#"{{"type": "table", "title": "行情", "cols": ["索引", "名称", "最新成交价"], "rows": [{}]}}"#, rows); LogStatus!(_D(None), "Total time for concurrent ticker retrieval:", endTS - beginTS, "ms", "\n", format!("`{}`", tbl)); Sleep(500); } }c++void main() { while(true) { int length = exchanges.size(); auto beginTS = UnixNano() / 1000000; vector<Ticker> arrTicker(length); vector<string> arrName(length); // 注意,添加几个交易所对象,这里要执行几次exchanges[n].Go 函数,这个例子是需要添加四个交易所对象,具体可以修改 auto r0 = exchanges[0].Go("GetTicker"); auto r1 = exchanges[1].Go("GetTicker"); auto r2 = exchanges[2].Go("GetTicker"); auto r3 = exchanges[3].Go("GetTicker"); vector<GoObj*> arrRoutine = {&r0, &r1, &r2, &r3}; for(int i = 0; i < length; i++) { arrName[i] = exchanges[i].GetName(); } for(int i = 0; i < length; i++) { Ticker ticker; arrRoutine[i]->wait(ticker); arrTicker[i] = ticker; } auto endTS = UnixNano() / 1000000; json tbl = R"({ "type": "table", "title": "行情", "cols": ["索引", "名称", "最新成交价"], "rows": [] })"_json; for(int i = 0; i < length; i++) { json arr = R"(["", "", ""])"_json; arr[0] = str_format("%d", i); arr[1] = arrName[i]; arr[2] = str_format("%f", arrTicker[i].Last); tbl["rows"].push_back(arr); } LogStatus(_D(), "Total time for concurrent ticker retrieval:", str_format("%d", endTS - beginTS), "ms", "\n", "`" + tbl.dump() + "`"); Sleep(500); } } -
并发调用
exchange.IO("api", ...)函数:javascriptfunction main() { /* 测试OKX期货下单接口 POST /api/v5/trade/order */ var beginTS = new Date().getTime() var param = {"instId":"BTC-USDT-SWAP","tdMode":"cross","side":"buy","ordType":"limit","px":"16000","sz":"1","posSide":"long"} var ret1 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", JSON.stringify(param)) var ret2 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", JSON.stringify(param)) var ret3 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", JSON.stringify(param)) var id1 = ret1.wait() var id2 = ret2.wait() var id3 = ret3.wait() var endTS = new Date().getTime() Log("id1:", id1) Log("id2:", id2) Log("id3:", id3) Log("Concurrent order time:", endTS - beginTS, "ms") }pythonimport time import json def main(): beginTS = time.time() param = {"instId":"BTC-USDT-SWAP","tdMode":"cross","side":"buy","ordType":"limit","px":"16000","sz":"1","posSide":"long"} ret1 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", json.dumps(param)) ret2 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", json.dumps(param)) ret3 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", json.dumps(param)) id1, ok1 = ret1.wait() id2, ok2 = ret2.wait() id3, ok3 = ret3.wait() endTS = time.time() Log("id1:", id1) Log("id2:", id2) Log("id3:", id3) Log("Concurrent order time:", endTS - beginTS, "seconds")rustfn main() { /* 测试OKX期货下单接口 POST /api/v5/trade/order */ let beginTS = UnixNano() / 1000000; // Rust不支持JSON序列化,直接使用原始字符串构造参数 let param = r#"{"instId":"BTC-USDT-SWAP","tdMode":"cross","side":"buy","ordType":"limit","px":"16000","sz":"1","posSide":"long"}"#; // 在Rust中,exchange.Go采用类型化写法,token为Go::IO,参数以元组形式传入 let ret1 = exchange.Go(Go::IO, ("api", "POST", "/api/v5/trade/order", "", param)); let ret2 = exchange.Go(Go::IO, ("api", "POST", "/api/v5/trade/order", "", param)); let ret3 = exchange.Go(Go::IO, ("api", "POST", "/api/v5/trade/order", "", param)); let id1 = ret1.wait(0); let id2 = ret2.wait(0); let id3 = ret3.wait(0); let endTS = UnixNano() / 1000000; Log!("id1:", id1); Log!("id2:", id2); Log!("id3:", id3); Log!("Concurrent order time:", endTS - beginTS, "ms"); }c++void main() { auto beginTS = UnixNano() / 1000000; json param = R"({"instId":"BTC-USDT-SWAP","tdMode":"cross","side":"buy","ordType":"limit","px":"16000","sz":"1","posSide":"long"})"_json; auto ret1 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", param.dump()); auto ret2 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", param.dump()); auto ret3 = exchange.Go("IO", "api", "POST", "/api/v5/trade/order", "", param.dump()); json id1 = R"({})"_json; json id2 = R"({})"_json; json id3 = R"({})"_json; ret1.wait(id1); ret2.wait(id2); ret3.wait(id3); auto endTS = UnixNano() / 1000000; Log("id1:", id1); Log("id2:", id2); Log("id3:", id3); Log("Concurrent order time:", endTS - beginTS, "ms"); } -
自动释放机制的测试
javascriptfunction main() { var counter = 0 var arr = [] // 用于测试持续引用并发对象的相关变量 var symbols = ["BTC_USDT", "ETH_USDT", "SOL_USDT", "LTC_USDT", "EOS_USDT"] while (true) { var arrRoutine = [] for (var symbol of symbols) { var r = exchange.Go("GetTicker", symbol) arrRoutine.push(r) // 记录并发对象,用于调用 r.wait() 函数获取结果,每轮循环都会清空 // arr.push(r) // 若使用这句代码,运行时会持续引用并发对象,导致其无法自动释放;当并发数超过 2000 时,会报错:```InternalError: too many routine wait, max is 2000```。 counter++ } // 遍历 arrRoutine 并调用 r.wait() 获取结果 LogStatus(_D(), "routine number:", counter) Sleep(50) } }rustfn main() { let mut counter = 0; let mut arr: Vec<TypedRoutine<Go::GetTicker>> = Vec::new(); // 用于测试持续引用并发对象的相关变量 let symbols = ["BTC_USDT", "ETH_USDT", "SOL_USDT", "LTC_USDT", "EOS_USDT"]; loop { let mut arrRoutine = Vec::new(); for symbol in symbols { // Rust 中 exchange.Go 为类型化写法,token 为 Go::GetTicker let r = exchange.Go(Go::GetTicker, (symbol,)); arrRoutine.push(r); // 记录并发对象,用于调用 r.wait(0) 函数获取结果,每轮循环都会清空 // arr.push(r); // 若使用这句代码,运行时会持续引用并发对象,导致其无法自动释放;当并发数超过 2000 时,会报错:InternalError: too many routine wait, max is 2000。 counter += 1; } // 遍历 arrRoutine 并调用 r.wait(0) 获取结果 LogStatus!(_D(None), "routine number:", counter); Sleep(50); } }
返回值
| 类型 | 描述 |
object |
|
参数
| 名称 | 类型 | 必填 | 描述 |
method | string | 是 |
|
arg | string / number / bool / object / array / function / any (平台支持的任意类型) | 否 | 并发执行函数的参数,参数 |
参考
Mail_Go HttpQuery_Go EventLoop exchange.IO(API 限流控制)
备注
该函数仅在实盘运行时创建多线程执行任务,回测不支持多线程并发执行任务(回测中可用,但仍为顺序执行)。
exchange.Go()函数返回对象后,可通过该对象调用其wait()函数获取线程返回的数据。当并发的多线程任务执行完毕且相关变量不再被引用时,系统底层会自动处理资源回收。
wait()方法支持超时参数:
1、不设置超时参数,即wait(),或将超时参数设置为0,即wait(0)。此时wait()函数会阻塞等待,直到并发线程运行完毕,并返回并发线程的执行结果。
2、将超时参数设置为-1,即wait(-1)。此时wait()函数会立即返回,不同编程语言的返回值有所不同,具体可参考本小节的调用示例。
3、设置具体的超时参数,即wait(300),此时wait()函数最多等待300毫秒后返回。
虽然系统底层具有自动回收机制,但如果持续引用相关变量,并发线程将不会被释放。当并发线程数量超过2000个时会报错:"too many routine wait, max is 2000"。
支持的函数:GetTicker,GetDepth,GetTrades,GetRecords,GetAccount,GetOrders,GetOrder,CancelOrder,Buy,Sell,GetPositions,IO等。这些函数并发调用时均基于当前exchange交易所对象执行。
Python语言与JavaScript语言的区别在于,Python语言中并发对象的wait()函数返回两个值,第一个为异步API调用返回的结果,第二个表示异步调用是否完成。
python
def main():
d = exchange.Go("GetRecords", PERIOD_D1)
# ok是一定返回True的, 除非策略被停止
ret, ok = d.wait()
# 如果等待超时, 或者wait了一个已经结束的实例,ok返回False
ret, ok = d.wait(100)