应用网络函数,或无需 DLL 的 MySQL:第 I 部分 - 连通器·进阶篇
(2/3)· 用原生网络函数打通 MT5 与数据库,告别动态库依赖的部署噩梦
从套接字字节流拼出 MySQL 报文
在 MT5 里直接对接 MySQL,第一步不是发查询,而是把 Socket 吐出来的零散字节拼回完整数据包。下面这段是接收循环的核心:先问 SocketIsReadable 有多少字节可读,有才进读取分支,否则直接跳过。 uint len=SocketIsReadable(m_socket); if(len) { //--- Read data from the socket to the buffer int rsp_len=SocketRead(m_socket,buf,len,m_timeout); m_rx_counter+= rsp_len; //--- Send the buffer for handling ENUM_TRANSACTION_STATE res = Incoming(buf,rsp_len); ... } //--- Handle received data ENUM_TRANSACTION_STATE Incoming(uchar &data[], uint len); SocketRead 的第三个参数是本次最多搬多少字节,m_timeout 控制阻塞上限;读到的长度累加到 m_rx_counter,这个计数器在排查「丢包还是慢」时很实用。 真正拼包逻辑在 Incoming() 里。状态枚举定义了四种结果:ERROR(-1)、IN_PROGRESS(0)、COMPLETE、SUBQUERY_COMPLETE。包头前 4 字节先收,凑满 4 字节才能用 reader.TotalLength(m_hdr) 算出整包长度,然后 ArrayResize 一次性开好缓冲。 ENUM_TRANSACTION_STATE CMySQLTransaction::Incoming(uchar &data[], uint len) { int ptr=0; // index of the current byte in the 'data' buffer ENUM_TRANSACTION_STATE result=MYSQL_TRANSACTION_IN_PROGRESS; // result of handling accepted data while(len>0) { if(m_packet.total_length==0) { //--- If the amount of data in the packet is unknown while(m_rcv_len<4 && len>0) { m_hdr[m_rcv_len] = data[ptr]; m_rcv_len++; ptr++; len--; } //--- Received the amount of data in the packet if(m_rcv_len==4) { //--- Reset error codes etc. m_packet.Reset(); m_packet.total_length = reader.TotalLength(m_hdr); m_packet.number = m_hdr[3]; //--- Length received, reset the counter of length bytes m_rcv_len = 0; //--- Highlight the buffer of a specified size if(ArrayResize(m_packet.data,m_packet.total_length)!=m_packet.total_length) return MYSQL_TRANSACTION_ERROR; // internal error } else // if the amount of data is still not accepted return MYSQL_TRANSACTION_IN_PROGRESS; } //--- Collect packet data while(len>0 && m_rcv_len<m_packet.total_length) { m_packet.data[m_rcv_len] = data[ptr]; m_rcv_len++; ptr++; len--; } //--- Make sure the package has been collected already 注意 m_hdr[3] 存的是包序号(sequence number),不是长度;长度由前 3 字节经 TotalLength 解出。若 ArrayResize 返回的数组大小不等于申请值,直接返回 ERROR,说明 MT5 终端内存分配失败,这种时候多半要降并发或关掉别的耗内存 EA。外汇与贵金属交易本就高杠杆高风险,用 socket 直连数据库做信号源时,断包重连逻辑没写好可能让 EA 在波动期失联。
class="type">uint len=SocketIsReadable(m_socket); if(len) { class=class="str">"cmt">//--- Read data from the socket to the buffer class="type">int rsp_len=SocketRead(m_socket,buf,len,m_timeout); m_rx_counter+= rsp_len; class=class="str">"cmt">//--- Send the buffer for handling ENUM_TRANSACTION_STATE res = Incoming(buf,rsp_len); ... } class=class="str">"cmt">//--- Handle received data ENUM_TRANSACTION_STATE Incoming(class="type">uchar &data[], class="type">uint len); enum ENUM_TRANSACTION_STATE { MYSQL_TRANSACTION_ERROR=-class="num">1, class=class="str">"cmt">// Error MYSQL_TRANSACTION_IN_PROGRESS=class="num">0, class=class="str">"cmt">// In progress MYSQL_TRANSACTION_COMPLETE, class=class="str">"cmt">// Fully completed MYSQL_TRANSACTION_SUBQUERY_COMPLETE class=class="str">"cmt">// Partially completed }; ENUM_TRANSACTION_STATE CMySQLTransaction::Incoming(class="type">uchar &data[], class="type">uint len) { class="type">int ptr=class="num">0; class=class="str">"cmt">// index of the current byte in the &class="macro">#x27;data&class="macro">#x27; buffer ENUM_TRANSACTION_STATE result=MYSQL_TRANSACTION_IN_PROGRESS; class=class="str">"cmt">// result of handling accepted data class="kw">while(len>class="num">0) { if(m_packet.total_length==class="num">0) { class=class="str">"cmt">//--- If the amount of data in the packet is unknown class="kw">while(m_rcv_len<class="num">4 && len>class="num">0) { m_hdr[m_rcv_len] = data[ptr]; m_rcv_len++; ptr++; len--; } class=class="str">"cmt">//--- Received the amount of data in the packet if(m_rcv_len==class="num">4) { class=class="str">"cmt">//--- Reset error codes etc. m_packet.Reset(); m_packet.total_length = reader.TotalLength(m_hdr); m_packet.number = m_hdr[class="num">3]; class=class="str">"cmt">//--- Length received, reset the counter of length bytes m_rcv_len = class="num">0; class=class="str">"cmt">//--- Highlight the buffer of a specified size if(ArrayResize(m_packet.data,m_packet.total_length)!=m_packet.total_length) class="kw">return MYSQL_TRANSACTION_ERROR; class=class="str">"cmt">// internal error } else class=class="str">"cmt">// if the amount of data is still not accepted class="kw">return MYSQL_TRANSACTION_IN_PROGRESS; } class=class="str">"cmt">//--- Collect packet data class="kw">while(len>class="num">0 && m_rcv_len<m_packet.total_length) { m_packet.data[m_rcv_len] = data[ptr]; m_rcv_len++; ptr++; len--; } class=class="str">"cmt">//--- Make sure the package has been collected already
「收包不全时如何中断事务」
在 MT5 里直接对接 MySQL 协议做底层通信时,必须判断已接收字节数是否覆盖整包。上面这段逻辑里,若 m_rcv_len 小于 m_packet.total_length,函数直接返回 MYSQL_TRANSACTION_IN_PROGRESS,把半包状态挂起,等下一次 OnChartEvent 或网络回调补齐。 判定通过后,代码将 m_rcv_len 与 m_packet.total_length 同时清零,意味着当前包已消费完,连接对象回到可接收下一包的状态。这个清零动作若漏掉,后续包长度累加会错乱,表现为偶发性的协议解析偏移。 紧随其后的 ENUM_PACKET_TYPE 枚举列出了六种包类型:NONE、DATA、EOF、OK、GREETING、ERROR。实际调试时,GREETING 只在握手后出现一次,而 DATA 与 EOF 常成对出现;若连续收到两个 EOF 而无 DATA,大概率是服务端提前断流,外汇或贵金属实盘环境下的网络闪断概率会明显抬高,需自行加重连逻辑。
if(m_rcv_len<m_packet.total_length) class="kw">return MYSQL_TRANSACTION_IN_PROGRESS; class=class="str">"cmt">//--- Handle received MySQL packet class=class="str">"cmt">//... class=class="str">"cmt">//--- m_rcv_len = class="num">0; m_packet.total_length = class="num">0; } class="kw">return result; } enum ENUM_PACKET_TYPE { MYSQL_PACKET_NONE=class="num">0, class=class="str">"cmt">// None MYSQL_PACKET_DATA, class=class="str">"cmt">// Data MYSQL_PACKET_EOF, class=class="str">"cmt">// End of file MYSQL_PACKET_OK, class=class="str">"cmt">// Ok MYSQL_PACKET_GREETING, class=class="str">"cmt">// Greeting MYSQL_PACKET_ERROR class=class="str">"cmt">// Error };
◍ 封装并投递 MySQL 报文的两个入口
向 socket 写数据比收数据简单:把 payload 塞进已准备好的发送缓冲,再调用一次 SocketSend() 即可。ping 走的是 0x0E 命令码,且序列号恒为 0,包头占 4 字节预留位,这点与查询不同。 query 的构造几乎一样,差别只在负载:命令码换成 0x03,后面紧接 SQL 字符串。两种发送在 SocketSend 返回值不等于待发长度时都直接返回 false,调用方应据此判断链路是否断掉。 无论 ping 还是 query,发完都要接此前写过的 ReceiveData() 接收逻辑。若该函数未报内部错误,业务才倾向视为成功,否则需重置接收缓冲重试。外汇与贵金属相关的行情存储若走这套 socket 通道,须留意网络抖动带来的重发风险。
class=class="str">"cmt">//+------------------------------------------------------------------+ class=class="str">"cmt">//| Form and send ping | class=class="str">"cmt">//+------------------------------------------------------------------+ class="type">bool CMySQLTransaction::ping(class="type">void) { if(reset_rbuf()==false) { SetUserError(MYSQL_ERR_INTERNAL_ERROR); class="kw">return false; } class=class="str">"cmt">//--- Prepare the output buffer m_tx_buf.Reset(); class=class="str">"cmt">//--- Reserve a place for the packet header m_tx_buf.Add(0x00,class="num">4); class=class="str">"cmt">//--- Place the command code m_tx_buf+=class="type">uchar(0x0E); class=class="str">"cmt">//--- Form a header m_tx_buf.AddHeader(class="num">0); class="type">uint len = m_tx_buf.Size(); class=class="str">"cmt">//--- Send a packet if(SocketSend(m_socket,m_tx_buf.Buf,len)!=len) class="kw">return false; m_tx_counter+= len; class="kw">return true; } class=class="str">"cmt">//+------------------------------------------------------------------+ class=class="str">"cmt">//| Form and send a query | class=class="str">"cmt">//+------------------------------------------------------------------+ class="type">bool CMySQLTransaction::query(class="type">class="kw">string s) { if(reset_rbuf()==false) { SetUserError(MYSQL_ERR_INTERNAL_ERROR); class="kw">return false; } class=class="str">"cmt">//--- Prepare the output buffer m_tx_buf.Reset(); class=class="str">"cmt">//--- Reserve a place for the packet header m_tx_buf.Add(0x00,class="num">4); class=class="str">"cmt">//--- Place the command code m_tx_buf+=class="type">uchar(0x03); class=class="str">"cmt">//--- Add the query class="type">class="kw">string m_tx_buf+=s; class=class="str">"cmt">//--- Form a header m_tx_buf.AddHeader(class="num">0); class="type">uint len = m_tx_buf.Size(); class=class="str">"cmt">//--- Send a packet if(SocketSend(m_socket,m_tx_buf.Buf,len)!=len) class="kw">return false; m_tx_counter+= len; class="kw">return true; }
拆解 CMySQLTransaction 的私有骨架
做 MT5 外接 MySQL 的业务封装时,CMySQLTransaction 是把连接、授权、收发包都收进一个类里的核心。它底下挂了几个私有成员:m_packet 是当前正在拆的 MySQL 数据包实例,m_tx_buf 是拼查询语句用的发送缓冲,m_auth 管密码加扰和握手参数,m_rbuf 接服务器回的 Ok 或 Data 包,reader 则专门解析包字节。 公开方法不用自己管 socket 的开关——调 CMySQLTransaction::Query() 时连接会按需建或关。持续连接模式下,首次 Query 建链,超过 m_keep_alive_tout(毫秒,0 表示不保活)就断;想保活就得在 OnTimer 里调它的 OnTimer,且计时器周期要小于 m_ping_period 和超时值,否则 ping 跟不上就掉线。 连接参数、账号、特殊客户端参数都得在 Query 之前设好。下面这段类声明的私有段,把网络层和时间戳全摊开了:m_timeout 是等 TCP 数据的毫秒数,m_rx_counter / m_tx_counter 是收发明文字节计数,m_dT 记最近查询时刻,m_last_ping_timestamp 记最近保活时刻——这些量在排查“为什么查询卡住”时直接能拿来打印。
class=class="str">"cmt">//+------------------------------------------------------------------+ class=class="str">"cmt">//| MySQL transaction class | class=class="str">"cmt">//+------------------------------------------------------------------+ class CMySQLTransaction { class="kw">private: class=class="str">"cmt">//--- Authorization data class="type">class="kw">string m_host; class=class="str">"cmt">// MySQL server IP address class="type">uint m_port; class=class="str">"cmt">// TCP port class="type">class="kw">string m_user; class=class="str">"cmt">// User name class="type">class="kw">string m_password; class=class="str">"cmt">// Password class=class="str">"cmt">//--- Timeouts class="type">uint m_timeout; class=class="str">"cmt">// timeout of waiting for TCP data(ms) class="type">uint m_timeout_conn; class=class="str">"cmt">// timeout of establishing a server connection class=class="str">"cmt">//--- Keep Alive class="type">uint m_keep_alive_tout; class=class="str">"cmt">// time(ms), after which the connection is closed; the value of class="num">0 - Keep Alive is not used class="type">uint m_ping_period; class=class="str">"cmt">// period of sending ping(in ms) in the Keep Alive mode class="type">bool m_ping_before_query; class=class="str">"cmt">// send &class="macro">#x27;ping&class="macro">#x27; before &class="macro">#x27;query&class="macro">#x27; (this is reasonable in case of large ping sending periods) class=class="str">"cmt">//--- Network class="type">int m_socket; class=class="str">"cmt">// socket handle class="type">class="kw">ulong m_rx_counter; class=class="str">"cmt">// counter of bytes received class="type">class="kw">ulong m_tx_counter; class=class="str">"cmt">// counter of bytes passed class=class="str">"cmt">//--- Timestamps class="type">class="kw">ulong m_dT; class=class="str">"cmt">// last query time class="type">uint m_last_resp_timestamp; class=class="str">"cmt">// last response time class="type">uint m_last_ping_timestamp; class=class="str">"cmt">// last ping time class=class="str">"cmt">//--- Server response CMySQLPacket m_packet; class=class="str">"cmt">// accepted packet class="type">uchar m_hdr[class="num">4]; class=class="str">"cmt">// packet header
「连接会话类的内部骨架与保活接口」
在 MT5 里直接对接 MySQL,核心不是发一条 SQL 就完事,而是先有一个事务类把 socket 收发、包解析和鉴权全揽下来。下面这段类声明揭示了它怎么组织缓冲区与回调。 m_rcv_len 记录已收包头的字节数,m_tx_buf 是出站传输缓冲,m_auth 专门处理握手期的授权请求;服务端回包放在 m_rbuf 动态数组,m_responses 统计回包条数。 接收侧靠 ReceiveData(ushort error_code) 等 socket 数据,Incoming() 按长度分派给 Ok / Greeting / Data / EOF / Error 五类处理器,每种返回 ENUM_TRANSACTION_STATE 让上层知道会话处于什么阶段。 对外接口里 Config() 设定 host、port、user、password 和 keep_alive_tout;KeepAliveTimeout() 与 PingPeriod() 控制保活,PingBeforeQuery() 决定每次查询前是否先 ping。OnTimer() 在开启保活后由定时器驱动。 想验证这套结构,把类声明贴进 MT5 头文件,建个 EA 调 Config("127.0.0.1",3306,"root","",10),再循环调 Responses() 看回包计数是否随 Query() 增长。外汇与贵金属行情存储走这个通道时,注意网络抖动可能引发会话断连,属高风险环节。
class="type">uint m_rcv_len; class=class="str">"cmt">// counter of packet header bytes class=class="str">"cmt">//--- Transfer buffer CData m_tx_buf; class=class="str">"cmt">//--- Authorization request class CMySQLLoginRequest m_auth; class=class="str">"cmt">//--- Server response buffer and its size CMySQLResponse m_rbuf[]; class="type">uint m_responses; class=class="str">"cmt">//--- Waiting and accepting data from the socket class="type">bool ReceiveData(class="type">class="kw">ushort error_code); class=class="str">"cmt">//--- Handle received data ENUM_TRANSACTION_STATE Incoming(class="type">uchar &data[], class="type">uint len); class=class="str">"cmt">//--- Packet handlers for each type ENUM_TRANSACTION_STATE PacketOkHandler(CMySQLPacket *p); ENUM_TRANSACTION_STATE PacketGreetingHandler(CMySQLPacket *p); ENUM_TRANSACTION_STATE PacketDataHandler(CMySQLPacket *p); ENUM_TRANSACTION_STATE PacketEOFHandler(CMySQLPacket *p); ENUM_TRANSACTION_STATE PacketErrorHandler(CMySQLPacket *p); class=class="str">"cmt">//--- Miscellaneous class="type">bool ping(class="type">void); class=class="str">"cmt">// send ping class="type">bool query(class="type">class="kw">string s); class=class="str">"cmt">// send a query class="type">bool reset_rbuf(class="type">void); class=class="str">"cmt">// initialize the server response buffer class="type">uint tick_diff(class="type">uint prev_ts); class=class="str">"cmt">// get the timestamp difference class=class="str">"cmt">//--- Parser class CMySQLPacketReader reader; class="kw">public: CMySQLTransaction(); ~CMySQLTransaction(); class=class="str">"cmt">//--- Set connection parameters class="type">bool Config(class="type">class="kw">string host,class="type">uint port,class="type">class="kw">string user,class="type">class="kw">string password,class="type">uint keep_alive_tout); class=class="str">"cmt">//--- Keep Alive mode class="type">void KeepAliveTimeout(class="type">uint tout); class=class="str">"cmt">// set timeout class="type">void PingPeriod(class="type">uint period) {m_ping_period=period;} class=class="str">"cmt">// set ping period in seconds class="type">void PingBeforeQuery(class="type">bool st) {m_ping_before_query=st;} class=class="str">"cmt">// enable/disable ping before a query class=class="str">"cmt">//--- Handle timer events(relevant when using Keep Alive) class="type">void OnTimer(class="type">void); class=class="str">"cmt">//--- Get the pointer to the class for working with authorization CMySQLLoginRequest *Handshake(class="type">void) {class="kw">return &m_auth;} class=class="str">"cmt">//--- Send a request class="type">bool Query(class="type">class="kw">string q); class=class="str">"cmt">//--- Get the number of server responses class="type">uint Responses(class="type">void) {class="kw">return m_responses;}
◍ 从响应指针摸到链路账单
把 MySQL 会话类收个尾,真正能派上用场的是那几个只读接口。Response(uint idx) 按索引拿某次服务器响应的指针,无参重载直接返回第 0 个,相当于默认看最新一笔回包。 GetServerError() 把 m_packet.error 这个结构体吐出来,连不上库或者语句报错时,错误码和报文都在里面,比只看 bool 返回值有用得多。 RequestDuration() 返回 m_dT,单位是微秒级 ulong,记的是上一次事务从发到收的耗时;RxBytesTotal() 和 TxBytesTotal() 分别给出累计收包和发包字节数,m_rx_counter 与 m_tx_counter 就是底层计数器。 ResetBytesCounters() 一句把两个计数器清零,跑长周期采样前调一下,后面读到的才是干净区间。外汇和贵金属行情峰值时 MySQL 往返可能突然拉长,这几个数配合日志能帮你判断是网络抖还是语句慢,这类外部依赖本身也有高波动风险。
class=class="str">"cmt">//--- Get the pointer to the server response by index CMySQLResponse *Response(class="type">uint idx); CMySQLResponse *Response(class="type">void) {class="kw">return Response(class="num">0);} class=class="str">"cmt">//--- Get the server error structure MySQLServerError GetServerError(class="type">void) {class="kw">return m_packet.error;} class=class="str">"cmt">//--- Options class="type">class="kw">ulong RequestDuration(class="type">void) {class="kw">return m_dT;} class=class="str">"cmt">// get the last transaction duration class="type">class="kw">ulong RxBytesTotal(class="type">void) {class="kw">return m_rx_counter;} class=class="str">"cmt">// get the number of received bytes class="type">class="kw">ulong TxBytesTotal(class="type">void) {class="kw">return m_tx_counter;} class=class="str">"cmt">// get the number of passed bytes class="type">void ResetBytesCounters(class="type">void) {m_rx_counter=class="num">0; m_tx_counter=class="num">0;} class=class="str">"cmt">// reset the counters of received and passed bytes };