为 MetaTrader 5 开发 MQTT 客户端:TDD 方法 - 最终篇·进阶篇
(2/3)·从 SUBSCRIBE 服务到本地代理搭建,把协议层抽象落到可测代码里
「在 MT5 里拉起 MQTT 订阅服务」
客户端目前只订阅 QoS 0 和 QoS 1,但更新分时报价用 QoS 0 就够——每 500 毫秒推一次 tick,丢一帧立刻被下一帧顶掉,视觉上无感。若你更在意 RAM、CPU 和带宽,把代码里 UpdateRates 的注释行打开、printf tick 的那行注释掉,就能只刷报价不刷分时,不过市场报价窗口和图表可能缺数据。 启动参数用 input 接代理 host 和 port,实盘前必须换成你自己的 broker 地址。OnStart 里先 new 出 CConnect 和 CSubscribe 对象,Clean Start 设 true 表示不要 broker 会话保留,Keep Alive 放宽到 3600 秒免得开发期频繁 ping;客户端标识符绝不能和发布端重复,这里写的是 MT5_SUB。 所有 new 出来的类指针都要在出错或停止时 delete,否则 MT5 服务内存泄漏。SendConnect 发完 CONNECT 包后查 CONNACK 和返回码,SendSubscribe 发完查 SUBACK 及原因码,任一步返回非 0 就 CleanUp 并 return -1。 网络层有个坑:MQL5 的套接字函数在缓冲区空时可能返回 1 字节的虚假可读状态,导致偶发读取异常。AlgoBook 里提到这是 MQL5 包装行为与系统 API 的差异,附件注释代码里做了规避。日志专家页若打出 5270 错误,说明订阅端在空转自言自语,得先跑起本地 MQTT broker 才能联通。
class="macro">#class="kw">property service class=class="str">"cmt">//--- input parameters input class="type">class="kw">string host = "class="num">172.20.class="num">106.92"; input class="type">int port = class="num">80; class=class="str">"cmt">//--- global vars class="type">int skt; CConnect *conn; CSubscribe *sub; class="type">uchar conn_pkt[]; conn = new CConnect(host, port); conn.SetCleanStart(true); conn.SetKeepAlive(class="num">3600); conn.SetClientIdentifier("MT5_SUB"); conn.Build(conn_pkt); class="type">uchar sub_pkt[]; sub = new CSubscribe(); sub.SetTopicFilter("MySPX500"); sub.Build(sub_pkt); if(SendConnect(host, port, conn_pkt) == class="num">0) { Print("Client connected ", host); } if(!SendSubscribe(sub_pkt)) { class="kw">return -class="num">1; } if(rsp[class="num">0] >> class="num">4 != CONNACK) { Print("Not Connect acknowledgment"); CleanUp(); class="kw">return -class="num">1; } if(rsp[class="num">3] != MQTT_REASON_CODE_SUCCESS) class=class="str">"cmt">// Connect Return code(Connection accepted) { Print("Connection Refused"); CleanUp(); class="kw">return -class="num">1; } if(((rsp[class="num">0] >> class="num">4) & SUBACK) != SUBACK) { Print("Not Subscribe acknowledgment"); } else Print("Subscribed"); if(rsp[class="num">5] > class="num">2) class=class="str">"cmt">// Suback Reason Code(Granted QoS class="num">2) { Print("Subscription Refused with error code %d ", rsp[class="num">4]); CleanUp(); class="kw">return class="kw">false; } msg += CPublish().ReadMessageRawBytes(inpkt); class=class="str">"cmt">//printf("New quote arrived for MySPX500: %s", msg); class=class="str">"cmt">//UpdateRates(msg); printf("New tick arrived for MySPX500: %s", msg);
◍ 把字符串报价灌进自定义品种
实时接收来的行情多半是一串带分隔符的文本,要写进 MT5 的自定义品种缓冲,得先拆后填。下面这段函数就是把外部推来的 tick 文本转成 MqlTick 结构并追加到名为 MySPX500 的自定义品种里。 拆包时用 StringSplit 以 ASCII 45(即减号 '-')为界,把时间、买卖价、成交价、整数与真实成交量、毫秒时间和标记位依次落进数组。注意 new_ticks_arr 下标 0~7 必须和 MqlTick 成员顺序严格对齐,错位一个字段就会导致 CustomTicksAdd 返回 0。 最后一步 CustomTicksAdd("MySPX500", last_tick) 返回值小于 1 代表写入失败,此时用 Print 把 _LastError 打出来排查。外汇与贵金属自定义品种同样适用这套逻辑,但价差跳变频繁,写入失败概率可能偏高,属正常高风险现象,开 MT5 跑一遍便能验证。
class="type">void UpdateTicks(class="type">class="kw">string new_ticks) { class="type">class="kw">string new_ticks_arr[]; StringSplit(new_ticks, class="num">45, new_ticks_arr); class="type">MqlTick last_tick[class="num">1]; last_tick[class="num">0].time = StringToTime(new_ticks_arr[class="num">0]); last_tick[class="num">0].bid = StringToDouble(new_ticks_arr[class="num">1]); last_tick[class="num">0].ask = StringToDouble(new_ticks_arr[class="num">2]); last_tick[class="num">0].last = StringToDouble(new_ticks_arr[class="num">3]); last_tick[class="num">0].volume = StringToInteger(new_ticks_arr[class="num">4]); last_tick[class="num">0].time_msc = StringToInteger(new_ticks_arr[class="num">5]); last_tick[class="num">0].flags = (class="type">uint)StringToInteger(new_ticks_arr[class="num">6]); last_tick[class="num">0].volume_real = StringToDouble(new_ticks_arr[class="num">7]); if(CustomTicksAdd("MySPX500", last_tick) < class="num">1) { Print("Update ticks failed: ", _LastError); } }
在 WSL 里跑 Mosquitto 代理的坑与填法
想在本机同时调 MT5 终端和 MQTT 代理,用 Windows 的 WSL 是最轻量的方案。若只跑单个示例,回环地址加不同客户端标识符就能通;但要并行多个示例并搭开发环境,最好把代理单独丢一台机器,客户端/服务器分主机跑,连接、鉴权、网络问题会提早暴露。
默认装好的 Mosquitto 会作为 Ubuntu 服务随 WSL 启动,开发时反而碍事。停掉服务、用 mosquitto -v 前台拉起,日志直接进 STDOUT,比 tail 跟踪日志更全,少漏关键握手信息。
MT5 的联网白名单必须加 WSL 主机名(hostname 命令可查),否则连接会被拒。Mosquitto 新版本默认仅允许本机连,Windows 侧算“另一台机”,得在 mosquitto.conf 写一行 listener 1883 放通。
MT5 只放行 80/443,非 TLS 的 MQTT 走 1883,需用 redir 把 80 转 1883:redir --lport=80 --cport=1883。漏掉白名单或转发,专家日志会报 5273 无法收发套接字数据。整套若顺,10–15 分钟能跑通,订阅服务起来后日志会印出已订阅确认。
「把报价推到 Mosquitto 的发布端写法」
发布服务和订阅端骨架一致,但客户端标识符必须分开——曾有人直接复制粘贴订阅代码,忘了改 ID,Mosquitto 因两个相同 clientID 互相踢线,消息看似发出却没抵达,调试半天才发现不是代码错而是 ID 撞车。 发布端跑的是一个 do-while 循环,靠 IsStopped() 退出;QoS 0 下 broker 不会回 PUBACK,所以连包校验都省了,拼好包、连上、SendPublish 就完事。想推 K 线而不是逐笔tick,把 GetRates() 那行取消注释、GetLastTick() 注释掉即可,主题过滤器两边要对齐。 字符串转换在发布端同样躲不掉,至少用户属性没完全二进制化前,payload 还是得走文本。MT5 日志专家页里出现「MQTT 发布服务已启动并已连接」字样,说明 conn_pkt 已建好;再去 WSL 的 mosquitto_sub 订阅同主题,能收到就证明通道通了。 别把 clientID 当摆设 复制订阅端时顺手把 SetClientIdentifier 改成 "MT5_PUB" 这类独立串,否则 broker 可能吞消息且不报任何错,这种「非错误」最耗时间。
class="type">uchar conn_pkt[]; conn = new CConnect(host, port); conn.SetCleanStart(true); conn.SetKeepAlive(class="num">3600); conn.SetClientIdentifier("MT5_PUB"); conn.Build(conn_pkt); do { class="type">uchar pub_pkt[]; pub = new CPublish(); pub.SetTopicName("MySPX500"); class=class="str">"cmt">//class="type">class="kw">string payload = GetRates(); class="type">class="kw">string payload = GetLastTick(); pub.SetPayload(payload); pub.Build(pub_pkt); class="kw">delete(pub); class=class="str">"cmt">//ArrayPrint(pub_pkt); if(!SendPublish(pub_pkt)) { class="kw">return -class="num">1; CleanUp(); } ZeroMemory(pub_pkt); Sleep(class="num">500); } class="kw">while(!IsStopped()); class="type">class="kw">string GetLastTick() { class="type">MqlTick last_tick; if(SymbolInfoTick("class="macro">#USSPX500", last_tick)) { class="type">class="kw">string format = "%G-%G-%G-%d-%I64d-%d-%G"; class="type">class="kw">string out; out = TimeToString(last_tick.time, TIME_SECONDS); out += "-" + StringFormat(format, last_tick.bid, class=class="str">"cmt">//class="type">class="kw">double last_tick.ask, class=class="str">"cmt">//class="type">class="kw">double last_tick.last, class=class="str">"cmt">//class="type">class="kw">double last_tick.volume, class=class="str">"cmt">//class="type">class="kw">ulong last_tick.time_msc, class=class="str">"cmt">//class="type">long last_tick.flags, class=class="str">"cmt">//class="type">uint last_tick.volume_real);class=class="str">"cmt">//class="type">class="kw">double Print(last_tick.time, ": Bid = ", last_tick.bid, " Ask = ", last_tick.ask, " Last = ", last_tick.last, " Volume = ", last_tick.volume, " Time msc = ", last_tick.time_msc, " Flags = ", last_tick.flags,