打开导航 打开导航 Open menu
量化交易系统开发(C++)/ 期权策略与风控 / 行业研究分析
量化交易系统開發(C++)/ 期权策略与风控 / 行业研究分析
C++ trading systems / options strategy & risk controls / equity research
adrian@adrianxv.cn
文章目录

交易系统:OrderServer 与 Market Data Publisher 交易系统:OrderServer 与 Market Data Publisher 交易系统:OrderServer 与 Market Data Publisher

本章覆盖 Order Server 和 Market Data Publisher 两个与客户端相互通信的组件

撰写编码协议

这部分先讲本书的交易系统当中,MarketDataPublisher 组件与客户端通信使用的编码协议,然后讲 OrderServer 组件与客户端通信使用的编码协议,最后讲在生产中使用的 Single Binary Encoding 协议。

Market Data Protocol

撮合引擎内部使用的用来传递市场更新的结构体是 MEMarketUpdate,这是内部格式;而此处我们需要定义 MarketDataPublisher 是面向广大市场参与者的公共格式。其实就是在 MEMarketUpdate 的基础上增加了序列号 seq_num_,客户端之后会用这个序列号来检测 UDP 丢包、进行恢复。

#pragma pack(push, 1)
struct MDPMarketUpdate {
private:
    size_t seq_num_;
    MEMarketUpdate me_market_update_;
};
#pragma pack(pop)

为了满足 snapshot stream 的需要,我们在 MarketUpdateType 枚举值(对应 me_market_update_.type_ 这个对象的可能取值)当中新增了以下几种消息类型:

  1. CLEAR:清空客户端本地订单簿
  2. SNAPSHOT_START:快照开始
  3. SNAPSHOT_END:快照结束

另外,新增以下一条 LFQueue,用于在 MarketDataConsumerSnapshotSynthesizer 之间通信:

typedef Common::LFQueue<Exchange::MDPMarketUpdate> MDPMarketUpdateLFQueue;

Order Protocol

#pragma pack(push, 1)
struct OMClientRequest {
    size_t seq_num_;
    MEClientRequest me_client_request_;
};

struct OMClientResponse {
    size_t seq_num_;
    MEClientResponse me_client_response_;
};
#pragma pack(pop)

这里的 OM 指的是 Order Management

本书当中使用 #pragma pack(1) 定义的的 C++ 结构体可以做到类似 Single Binary Encoding 的编码,因为字段按照固定顺序排列,最终

这里没有一个独立的序列化器把结构体字段转成二进制,而是依赖编译器根据结构体定义和 #pragma pack 来决定内存布局。这个对象在内存中的连续字节就是随后被发送的二进制。在 TCPSocket::send() 当中,memcpy 把内存中的字节复制到 send_buffer_,随后 OrderServer 线程通过 epoll 等待 socket 可写,然后 socket 把内存缓冲区中的字节发送到对端。实际上,由于服务器与客户端都是 C++,对端收到消息后把二进制 reinterpret_cast 成 struct 就可以了。

所以本书的实现可以看成是一个建议的 Binary Protocol 实现。下面我们看一下真正的 SBE 的实现。

Simple Binary Encoding (SBE)

SBE 协议是开源的,其 Github 仓库链接。它可以用来编码任何结构化信息,甚至包括图片、视频、PDF 等,但是 SBE 的典型用途是低延迟结构化信息,例如订单、行情、执行回报。所以如果在本书的交易系统当中使用 SBE,那么 Market Data Publisher, Order Server/Gateway 都可以用 SBE 来编码。

交易所与客户端之间使用 SBE 来通信的流程大致是:

  1. SBE 社区提供一套标准,定义 SBE XML schema 的语法以及允许的结构
  2. 交易所根据这套标准,编写一套包含业务语义的 SBE XML schema,当中可能包含 OrderId, price 等字段
  3. SBE 社区提供 SBE tools,能够根据规范的 SBE XML Schema 生成一系列 codec (是 coder-decoder 的简写,编码译码器)
  4. 交易所和客户端使用 SBE tools 根据交易所定义好的 schema 生成 codec。这一系列的 codec 可能包含 OrderEncoder, OrderDecoder, MarketDataEncoder, MarketDataDecoder
  5. 应用层使用这套 Encoder, Decoder 来进行通信时,只需要知道它们提供哪些 API,无需了解底层机制。具体有哪些 API,完全由 schema 决定,它们反映业务语义,比如:
OrderEncoder encoder;
encoder.wrapForEncode(buffer, 0, sizeof(buffer))
       .side(Side::BUY)
       .qty(100)
       .price(12345);

OrderDecoder decoder;
decoder.wrapForDecode(buffer, messageOffset, actingBlockLength, actingVersion, bufferLength);
auto side = decoder.side();
auto qty = decoder.qty();
auto price = decoder.price();

OrderEncoder 只负责把编码好的二进制消息放入内存缓冲区,后续由 OrderGateway 线程的 epoll 唤醒 socket 发送消息,这些就不是 encoder 的职责了。

我们可以根据 XML schema 和 SBE message header 推断最终 encoder 放到内存缓冲区当中的二进制信息的大概样貌

15 00 01 00 07 00 00 00  // 这是 header,包含 blockLength、templateId、schemaId、version
2A 00 00 00 00 00 00 00  // 这是 orderId = 42
01                       // 这是 side = BUY = 1
64 00 00 00              // 这是 qty = 100
39 30 00 00 00 00 00 00  // 这是 price = 12345

据我目前所知,中国的交易所并不使用 SBE,但是 SBE 是理解交易系统编码协议的很好的切入口。

OrderServer

alt text

OrderServer 的职责是:

  1. 创建 TCP 服务端并接收客户连接(监听 socket 的职责)
  2. 从 TCP 缓冲区按固定大小解析 OMClientRequest(这个 TCP缓冲区指的是 Server 端的缓冲区?)
  3. 检查客户端 client_id, socket 映射和请求序号(什么是 socket 映射和请求序号?)
  4. 把请求交给 OrderServer 内部的 FIFOSequencer 排序
  5. 把排序后的内部 MEClientRequest 放入 Matching Engine 的无锁队列

如果是返回私有订单回调,流程则是:

  1. 从 Matching Engine 接收 MEClientResponse
  2. 根据 client_id 找到对应 TCP 连接,封装成 OMClientResponse 并返回客户端

OrderServer 的成员对象

class OrderServer {
private:
    const std::string iface_; // OrderServer 监听所使用的本机网络接口
    const int port_ = 0;      // OrderServer 对外监听的本地接口。

    ClientResponseLFQueue *outgoing_responses = nullptr; // OrderServer 与 MatchingEngine 通信的队列。这里只是指针,不拥有

    volatile bool run_ = false;

    std::string time_str_;
    Logger logger_;

    std::array<size_t, ME_MAX_NUM_CLIENTS> cid_next_outgoing_seq_num_;
    std::array<size_t, ME_MAX_NUM_CLIENTS> cid_next_exp_seq_num_;

    std::array<Common::TCPSocket*, ME_MAX_NUM_CLIENTS> cid_tcp_socket_;
    Common::TCPServer tcp_server_;

    FIFOSequencer fifo_sequencer_;
};

iface_ 是网络接口名称。OrderServer 需要通过操作系统 API 找到接口名称对应的本机 IP 地址,然后将 socket 绑定到该接口。port_ 本身表示 OrderServer 机器上的一个端口。 OrderServer 的监听 socket 和很多连接 socket 绑定的本机 IP 地址和本机端口都是一致的;不同连接 socket 的对端 IP 地址和端口号各异。

cid_next_outgoing_seq_num_cid_next_exp_seq_num_ 分别用来保存期望下一个 ClientResponse 和 ClientRequest 应当携带的 seq_num_ 序列号。

cid_tcp_socket_ 保存指向每一个与客户端通信使用的连接 socket 的指针。此处的三个 std::array 都以 client_id 作为下标。client_id 实质上是 0 到 ME_MAX_NUM_CLIENTS - 1 范围内的一个整数。

tcp_server_ 是 OrderSever 内部的子组件,负责管理一切 socket 事务

注意此处 OrderServer 并没有一个专门的子组件来负责 ClientRequest 或是 ClientResponse 的编码与解码,因为本书实现的编码协议很简单且服务器与客户端都是 C++,所以发送方直接发送内存中对象的二进制字段,接收方直接接收后使用 reinterpret_cast 即可。若是使用真实的 SBE 协议,则通常有独立的 OrderEncoder / OrderDecoder。

OrderServer::OrderServer()

OrderServer::OrderServer(ClientReqestLFQueue* client_requests, ClientResponseLFQueue* client_responses_, const std::string &iface, int port) 
    : iface_(iface), port_(port), outgoing_responses_(client_responses), logger_("exchange_order_server.log"), tcp_server_(logger_), fifo_sequencer_(client_requests, &logger_) 
{
    cid_next_outgoing_seq_num_.fill(1); // 把 std::array 当中的所有元素都填上 1
    cid_next_exp_seq_num_.fill(1);
    cid_tcp_socket_.fill(nullptr);

    tcp_server_.recv_callback_ = [this](auto socket, auto rx_time) { recvCallcack(socket, rx_time); };
    tcp_server_.recv_finished_callback_ = [this]() { recvFinishedCallback(); };
}

tcp_server_ 本身是通用网络组件,不知道收到数据后如何处理,这里把 recvCallback() 作为函数对象(lambda)赋值给 tcp_server_.recv_finished_callback_

OrderServer::start()

MatchingEngine::start() 一样,保持不阻塞,把工作循环放在 run() 当中

auto OrderServer::start() {
    run_ = true;
    tcp_server_.listen(iface_, port_);
    Common::creatAndStartThread(-1, "Exchange/OrderServer", [this]() { run(); });
}

OrderServer::stop()

就这么简单

auto OrderServer::stop() {
    run_ = false;
}

OrderServer::run()

auto OrderServer::run() noexcept {
    logger_.log("%:% %() %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_));

    while (run_) {
        tcp_server_.poll();

        tcp_server_.sendAndRecv();

        // 遍历无锁队列当中每一条私有回调消息
        for (auto client_response = outgoing_responses_->getNextToRead(); outgoing_responses_->size() && client_response; client_response = outgoing_responses_->getNextToRead()) {

            // 获得 client_id
            client_id = client_responses->client_id_;

            // 获得这条回调消息应有的序列号,这里用了引用
            auto& next_outgoing_seq_num = cid_next_outgoing_seq_num_[client_id];

            logger_.log("%:% %() % Processing cid:% seq:% %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), client_id, next_outgoing_seq_num, client_response->toString());

            // 确保该客户端对应的连接 socket 存在
            ASSERT(cid_tcp_socket_[client_id] != nullptr, "Dont have a TCPSocket for ClientId:" + std::to_string(client_id));

            cid_tcp_socket_[client_id]->send(&next_outgoing_seq_num, sizeof(next_outgoing_seq_num));
            cid_tcp_socket_[client_id]->send(client_response, sizeof(client_response));

            outgoing_responses_->updateReadIndex();

            ++next_outgoing_seq_num;
        }
    }
}

OrderServer::recvCallback()

最终注册到 OrderServer::tcp_server_ 所持有的连接 socket 的 TCPSocket::recv_callback_ 的回调函数。负责从这个连接 socket 的应用层缓冲区(即 inbound_data_)的字节流中按照 OMClientRequest 的字节数提取出完整的请求消息,经过 reinterprec_cast 还原成 OMClientRequest 结构体,经过简单校验后加入 fifo_sequencer_

auto OrderServer::recvCallback(TCPSocket *socket, Nanos rx_time) noexcept {
    logger_.log("%:% %() % Received socket:% len:% rx:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket->socket_fd_, socket->next_rcv_valid_index_, rx_time);

    if (socket->next_rcv_valid_index_ >= sizeof(OMClientRequest){
        // 这里前者是 rcv_buffer 中当前已有的有效字节数;后者是一条完整请求的字节数
        // 这里是想要判断接收缓冲区中是否至少有一条完整消息
        // 如果不足,则等到下一次再读取解析

        size_t i = 0; // i 表示当前要解析的消息在缓冲区中的起始位置
        for (; 
            i + sizeof(OMClientRequest) <= socket->next_rcv_valid_index_; 
            i += sizeof(OMClientRequest)) {
                // 每次解析一条消息,因此 i 的步进是 sizeof(OMClientRequest)
                
                auto request = reinterpre_cast<const OMClientRequest*>(socket->inbound_data_.data() + i);

                logger_.log("%:% %() % Received %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), request->toString());

                if (UNLIKELY(cid_tcp_socket_[request->me_client_request_.client_id_] == nullptr)) { // 这是当前客户的第一条消息,当前 socket* array 还未有这个客户端的 socket
                    cid_tcp_socket_[request->me_client_request_.client_id_] = socket;

                    if (cid_tcp_socket_[request->me_client_request_.client_id_] != socket) {   // 当前这个客户端的 socket 与当前 socket 不符,应该拒绝该请求
                        
                        logger_.log("%:% %() % Received ClientRequest from ClientId:% on different socket:% expected:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), request->me_client_request_.client_id_, socket->socket_fd_, cid_tcp_socket_[request->me_client_request_.client_id_]->socket_fd_);

                        continue;
                    } 

                    // 这个客户下一个应该的 seq_num
                    auto &next_exp_seq_num = cid_next_exp_seq_num_[request->me_client_request_.client_id_];

                    if (request->seq_num_ != next_exp_seq_num) { // seq_num 不符

                        logger_.log("%:% %() % Incorrect sequence number. ClientId:% SeqNum expected:% received:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), request->me_client_request_.client_id_, next_exp_seq_num, request->seq_num_);
                        
                        continue;
                    }

                    ++next_exp_seq_num;

                    fifo_sequencer_.addClientRequest(rx_time, request->me_client_request);
                }

                // 此时已经解析完本轮 TCP 字节流当中完整的 ClientRequest 消息,i 已经指向了最后“半条”消息的起始点。
                // 换个角度理解,i 表示在本轮字节流范围内,目前已读的消息相比于字节流开头的 offset
                
                memcpy(socket->inbound_data.data(), socket->inbound_data_.data() + i, socket->next_rcv_valid_index_ - i)
                // 这里的 memcpy 就是要把那“半条消息”搬到缓冲区开头

                socket->next_rcv_valid_index_ -= i;
                // 然后把有效 index 减去本轮已读的 offset,使其反映缓冲区当中有效信息截止点
                
            }
    }) 

}

函数内的 for (; i + sizeof(OMClientRequest) <= socket->next_rcv_valid_index_; i += sizeof(OMClientRequest)) 这个循环遍历当前 socket 接收缓冲区中所有完整的 OMClientRequest

是有可能一次读到多条消息的,因为 TCP 字节流传输就是有可能一次 recvmsg() 读到半条消息,或者一条半消息。意味着有可能一次读到的多条消息共享一个 rx_time,因此无法知道每条消息到达内核的精确事件,因此只能做近似公平排序,而不是逐消息硬件时间排序。如果要改进的话,要在网卡和驱动上启动硬件时间戳之类的。

socket->inbound_data.data() 这个缓冲区是不会把其中的信息都清除的,只有每次对其中的信息进行覆写。控制读写位置的核心变量是 next_rcv_valid_index_。缓冲区只在构造时 resize(TCPBufferSize) 一次,之后从不 clear。它就是一块固定的内存,旧的字节会一直留着,只是变成无用的垃圾数据。有效区域是 [0, next_rcv_valid_index_],这之外的内容没有意义。

由于 recvCallback() 第一次读取这个缓冲区是从 0 开始,且每次读取的字节数都是一个 OMClientRequest 的字节数的整数倍,因此可以确保读入的每个 OMClientRequest 都是完整的 OMClientRequest

socket 收到数据之后,回调会逐层通知到 OrderServer,路径是:

  1. TCPServer::sendAndRecv()
  2. TCPSocket::sendAndRecv()
  3. TCPSocket::recv_callback_(this, kernel_time)
  4. OrderServer::recvCallback()

recvCallback() 与 recvFinishedCallback() 的分工:

  • recvCallback(socket, rx_time):每当一个客户端 socket 读到数据时调用。它解析并校验该客户端的请求,再连同接收事件交给 FIFOSequencer 暂存
  • recvFinishedCallback():这一轮所有可读 socket 都处理完后只调用一次。让 FIFOSequencer 将本轮收集的请求按事件排序,再发布给 Matching Engine

OrderServer::recvFinishedCallback()

负责把若干个 socket 的订单请求送入 FIFOSequencer。作为回调函数,由 TCPServer::sendAndRecv() 在本轮 epoll 的所有 receive_sockets_ 的读取事件处理完毕后调用。也即是说,在调用当前方法的时候,已经完成以下的事项:receive_sockets_ 当中的每一个 socket 已经把内核缓冲区的数据放到应用层缓冲区,并且调用注册好的回调函数 OrderServer::recvCallback()OMClientRequest 消息放入 fifo_sequencer_.pending_client_requests_ 当中。因此在方法只是调用 fifo_sequencer_.sequenceAndPublish()fifo_sequencer_.pending_client_requests_ 当中的请求按时间顺序排序。

auto OrderServer::recvFinishedCallback() noexcept {
    fifo_sequencer_.sequenceAndPublish();
}

FIFOSequencer

核心目的是保证公平

  1. OrderServer 收到请求时记录软件接收时间
  2. 请求先进入待等待数组(为什么需要有等待数组?)
  3. sequenceAndPublish() 按接收时间排序
  4. 按 FIFO 顺序写入 Matching Engine 的无锁队列

FIFOSequencer::addClientRequest()recvCallback() 针对 OrderServer 收到的每条 OMClientRequest 逐条调用,只负责将 (rx_time, request) 追加到内部暂存数组中,不排序、不发布。

FIFOSequencer::sequenceAndPublish()recvFinishedCallback() 针对整批 socket 调一次,负责使用 std::sort 按 recv_time_ 排序,然后逐条写入无锁队列,然后把 pending_size_ = 0,等待下一批。

FIFOSequencer 的成员对象

class FIFOSequencer {
private:
    clientRequestLFQueue* incoming_requests_ = nullptr; // FIFOSequencer 也要有与撮合引擎通信的无锁队列

    std::string time_str_;
    Logger* logger_ = nullptr;

    struct RecvTimeClientRequest { // 就是把 request 的结构体和到达时间再加一层封装
        Nanos recv_time_ = 0;
        MEClientRequest request_;

        auto operator<(const RecvTimeClientRequest& rhs) const {  // 这是小于的比较符号
            return (recv_time_ < rhs.recv_time_);
        }
    };

    // 每处理一个 OMClientRequest 那么大的 TCP 字节流,就写入这个 array 一次
    // 然后等到触发 sequenceAndPublish() 的时候再写入与撮合引擎的无锁队列
    std::array<RecvTimeClientRequest, ME_MAX_PENDING_REQUESTS> pending_client_requests_;
    size_t pending_size_ = 0; // 表示上述 array 的大小
};

FIFOSequencer::addClientRequests()

把给出的 MEClientRequest 对象加入 FIFOSequencer::pending_client_requests_

Q:为什么在这一步不保证加入是按照时间顺序的,还需要额外的 sequenceAndPublish() ? A:因为 FIFOSequencer::addClientRequests() 被调用的顺序由 TCPSocket::sendAndRecv() 决定,TCPSocket 在这个函数内调用回调函数,进而调用 FIFOSequencer::addClientRequests()。而 TCPSocket::sendAndRecv() 被调用的顺序由 TCPServer 遍历 receive_sockets_ 的顺序决定,这个顺序并不一定等于消息到达内核的顺序。为了公平,需要在 FIFOSequencer::sequenceAndPublish() 当中再按照消息到达内核的时间重新排序。

auto FIFOSequencer::addClientRequests(Nanos rx_time, const MEClientRequest& request) {
    if (pending_size >= ME_MAX_PENDING_REQUESTS) {
        FATAL("Too many pending requests");
    }

    pending_client_requests_.at(pending_size++) = RecvTimeClientRequest{rx_time, request}; 
}

把元素加入 pending_client_requests_ 本身不会让元素有序是吗?为什么?这里不是都已经按照加入的时间来加入了吗?

FIFOSequencer::sequenceAndPublish()

pending_client_requests_ 当中的 RecvTimeClientRequest 按照时序排列,然后加入到 OrderServer 与撮合引擎通信的无锁队列 incoming_requests_ 当中

auto FIFOSequencer::sequenceAndPublish() {
    // 当前 pending_client_requests_ 为空,直接返回
    if (UNLIKELY(!pending_size_)) { return; }

    logger_->log("%:% %() % Processing % requests.\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), pending_size_);

    // 对 pending_client_requests_ 当中的元素排序
    // 我们自定义了 FIFOSequencer::RecvTimeClientRequest::operator<,会按照 recv_time 排序
    std::ranges::sort(pending_client_requests_.begin(), pending_client_requests_.begin() + pending_size_);

    // 把客户请求加入与撮合引擎的无锁队列
    for (size_t i = 0; i < pending_size_; ++i) {
        // const 引用避免复制元素
        const auto &client_request = pending_client_requests_.at(i);

        logger_->log("%:% %() % Writing RX:% Req:% to FIFO.\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), client_request.recv_time_, client_request.request_.toString());

        auto next_to_write_to = incoming_requests_->getNextToWriteTo();

        next_to_write_to* = client_request.request_;

        incoming_requests_->updateWriteIndex();
    }

    // 把 pending_client_requests_ 的元素计数重新置零
    pending_size_ = 0;
}

MarketDataPublisher

SnapshotSynthesizer 负责间隔一段时间就向 snapshot stream 发送一次订单簿快照

这个订单簿快照其实不是一个真正的 “快照”,而是短时间内发送 CLEAR 和多个 ADD 事件,其中 CLEAR 告诉客户端,随后的 ADD 事件将从零开始重建订单簿。而订单簿上的每一个订单都对应后续的一个 ADD 事件。

SNAPSHOT_START 标记快照流开始;SNAPSHOT_END 标记快照流结束

SnapshotSynthesizer 运行在独立线程中。

MarketDataPublisher 的成员变量

class MarketDataPublisher {
private:
    // 市场更新序列号,从 1 开始
    size_t next_inc_seq_number_ = 1;
    // 撮合引擎与 MDP 之间的队列
    MEMarketUpdateLFQueue *outgoing_md_updates_ = nullptr;

    // MDP 与 snapshot synthesizer 之间的队列
    MDPMarketUpdateLFQueue snapshot_md_updates_;

    volatile bool run_ = false;  // 话说 volatile 关键词的作用究竟是啥,好像没系统学过

    std::string time_str_;
    Logger logger_;

    // 因为是 UDP multicast,所以这里只有唯一的 socket,不区分监听/连接
    Common::McastSocket incremental_socket_;

    SnapshotSynthesizer *snapshot_synthesizer_ = nullptr;
}

MarketDataPublisher::start()

auto MarketDataPublisher::start() {
    run_ = true;

    ASSERT(Common::createAndStartThread(-1, "Exchange/MarketDataPublisher", [this]() { run(); }) != nullptr, "Failed to start MarketData thread.");

    snapshot_synthesizer_->start();
}

MarketDataPublisher::run()

auto MarketDataPublisher::run() noexcept {
    logger_.log("%:% %() %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_));

    while (run_) {
        // 在工作循环内不断遍历与撮合引擎通信的无锁队列,
        for (auto market_update = outgoing_md_updates_->getNextToRead(); outgoing_md_updates_->size() && market_update; market_update = outgoing_md_updates_ ->getNextToRead()) {

            logger_.log("%:% %() % Sending seq:% %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), next_inc_seq_num_, market_update->toString().c_str());

            // incremental stream 的消息是不用经过 SnapshotSynthesizer 的
            // Q: 为什么这里要分两次来 send()?
            // A: 分别发送 MDPMarketUpdate 结构体的 seq_num_ 和 MEMarketUpdate 两部分
            incremental_socket_.send(&next_inc_seq_num_, sizeof(next_inc_seq_num_));
            incremental_socket_.send(market_update, sizeof(MEMarketUpdate));
            outgoing_md_updates_->updateReadIndex();

            // 与 SnapshotSynthesizer 通信
            auto next_write = snapshot_md_updates_.getNextToWriteTo();
            // Q:为什么这里不是直接解引用写入?
            // A:因为队列存的是 MDPMarketUpdate 结构体,其2个字段要分别写入
            next_write->seq_num_ = next_inc_seq_num_;
            next_write->me_market_update_ = *market_update;
            snapshot_md_updates_.updateWriteIndex();

            // 市场更新的序列号
            ++next_inc_seq_num_;
        }

        // 是把当前无锁队列当中所有市场更新消息写入应用层缓冲区后,一并写入内核缓冲区
        incremental_socket_.sendAndRecv();
    }

}

SnapshotSynthesizer

SnapshotSynthesizer 成员变量

class SnapshotSynthesizer {
private:
    // MDP 与 SnapshotSynthesizer 之间的通信无锁队列
    MDPMarketUpdateLFQueue *snapshot_md_updates_ = nullptr;

    Logger logger_;

    volatile bool run_ = false;

    std::string time_str_;

    // snapshot stream 的 UDP socket
    McastSocket snapshot_socket_;

    // ??:为什么这里要嵌套 array
    std::array<std::array<MEMarketUpdate*, ME_MAX_ORDER_IDS>, ME_MAX_TICKERS> ticker_orders_;

    size_t last_inc_seq_num_ = 0;  // 标记当前 snapshot 截止到哪一条增量更新序号
    Nanos last_snapshot_time_ = 0;

    MemPool<MEMarketUpdate> order_pool_;
};

std::array<std::array<MEMarketUpdate*, ME_MAX_ORDER_IDS>, ME_MAX_TICKERS> ticker_orders_; 这是一个 array 套 array 得到的二维数组,用 ticker_id_order_id_ 确认一张订单。意味着 SnapshotSynthesizer 要维护一个撮合引擎当中订单簿的副本快照。后续用于在 snapshot stream 上发布订单簿快照。

Q:既然 SnapshotSynthesizer 有用二维数组存储的订单簿快照,为什么 snapshot stream 仍然要靠 CLEAR 和一条条 ADD 消息来发布快照,为什么不直接把整个二维数组发给客户端? A:因为二维数组内部存储的是订单指针,指针直接发给客户端无意义。且二维数组内部有大量空槽位,整个数组发送不高效。在 snapshot stream 当中发送 CLEARADD 消息可以复用已有数据格式。

// TODO:当然可以设计一套更高效的专门的批量快照协议。要额外花心思,不能复用已有协议。

SnapshotSynthesizer::addToSnapshot()

void SnapshotSynthesizer::addToSnapshot(const MDPMarketUpdate *market_update) {
    // 抽取出 MEMarketUpdate 结构体
    const auto &me_market_update = market_update->me_market_update_;

    // 获取当前合约的所有订单
    auto* orders = &ticker_orders_.at(me_market_update.ticker_id_);

    // 这里用一个 switch statement 来处理各种市场更新
    switch (me_market_update.type_) {
        case MarketUpdateType::ADD: {
            auto order = orders->at(me_market_update.order_id_);

            // 确保该订单号的订单不存在
            ASSERT(order == nullptr, "Received:" + me_market_update.toString() + " but order already exists:" + (order ? order->toString() : ""));
            // 在内存池上分配空间,其实只是把 new 换成 order_pool_.allocate() 而已
            orders->at(me_market_update.order_id_) = order_pool_.allocate(me_market_update);
            break;
        }

        case MarketUpdateType::MODIFY {
            auto order = orders->at(me_market_update.order_id_);

            // 确保该订单号的订单已经存在
            ASSERT(order != nullptr,  "Received:" + me_market_update.toString() + " but order does not exist.");

            // 检查不变量,防止数组下标当中的 order_id_ 与订单实际 order_id 不一致
            ASSERT(order->order_id_ == me_market_update.order_id_, "Expecting existing order to match new one.");
            // 确保 Side 一致
            ASSERT(order->side_ == me_market_update.side_, "Expecting existing order to match new one.");

            // 更新订单的量价
            order->price = me_market_update.price_;
            order->qty = me_market_update.qty_;

            break;
        }

        case MarketUpdateType::CANCEL {
            auto* order = orders.at(me_market_update.order_id_);

            ASSERT(order != nullptr,  "Received:" + me_market_update.toString() + " but order does not exist.");
            ASSERT(order->order_id_ == me_market_update.order_id_, "Expecting existing order to match new one.");
            ASSERT(order->side_ == me_market_update.side_, "Expecting existing order to match new one.");

            order_pool_.deallocate(order);
            orders->at(me_market_update.order_id_) = nullptr;
            break;
        }

        // 对订单簿没有影响的市场更新
        case MarketUpdateType::SNAPSHOT_START:
        case MarketUpdateType::CLEAR:
        case MarketUpdateType::SNAPSHOT_END:
        case MarketUpdateType::TRADE:
        case MarketUpdateType::INVALID:
        break;
    }

    // 确保撮合引擎传来的每条 MDPMarketUpdate 都处理了,更新了 SnapshotSynthesizer 的订单簿快照
    ASSERT(market_update->seq_num == last_inc_seq_num_ + 1, "Expected incremental seq_nums to increase.")
    last_inc_seq_num_ = market_update->seq_num_;
}

SnapshotSynthesizer::publishSnapshot()

snapshot stream 实际是通过一系列消息(而不是单一条包含订单簿所有信息的消息)来帮助客户端重建订单簿的。涉及到的消息类型如下:

  1. MarketUpdateType::SNAPSHOT_START:表示随后的消息用于重建订单簿
  2. MarketUpdateType::CLEAR:表示客户端可以清楚其订单簿上的所有订单,等待随后添加
  3. MarketUpdateType::ADD:订单簿快照中每一个订单对应一条这个类型的消息
  4. MarketUpdateType::SNAPSHOT_END:表示订单簿快照消息结束。
auto SnapshotSynthesizer::publishSnapshot() {
    size_t snapshot_size = 0;

    // MarketUpdateType::SNAPSHOT_START 消息标识订单簿快照消息开始
    const MDPMarketUpdate start_market_update{snapshot_size++, {MarketUpdateType::SNAPSHOT_START, last_inc_seq_num_}}; //??:为什么用到了 snapshot_size++ 变量?

    logger_.log("%:% %() % %\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_), start_market_update.toString());

    snapshot_socket_.send(&start_market_update, sizeof(MDPMarketUpdate));

    // 遍历所有合约
    for (size_t ticker_id = 0; ticker_id < ticker_orders_.size(); ++ticker_id) {
        // 这个合约下的所有订单
        const auto& orders = ticker_orders_.at(ticker_id);

        MEMarketUpdate me_market_update;  //??:从空的结构体开始构造吗?为什么不直接构造有数据的结构体?
        me_market_update.type_ = MarketUpdateType::CLEAR;
        me_market_update.ticker_id_ = ticker_id;

        const MDPMarketUpdate clear_market_update{snapshot_size++, me_market_update};
        logger_.log("%:% %() % %\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_), clear_market_update.toString());
        
        snapshot_socket_.send(&clear_market_update, sizeof(MDPMarketUpdate));

        // 遍历该合约下所有订单
        for (const auto order: orders) {
            if (order) { // ??:加这层判断意义何在?
                // ??:order 是 MEMarketUpdate* 类型的,那么为何这里还要 *order,那不就是指针的指针了吗?
                const MDPMarketUpdate market_update{snapshot_size++, *order};

                logger_.log("%:% %() % %\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_), market_update.toString());

                snapshot_socket_.send(&market_update, sizeof(MDPMarketUpdate));
                snapshot_socket_.sendAndRecv();
            }
        }
    } // 遍历所有合约结束

    // MarketUpdateType::SNAPSHOT_START 消息标识订单簿快照消息结束
    const MDPMarketUpdate end_market_update{snapshot_size++, {MarketUpdateType::SNAPSHOT_END, last_inc_seq_num_}};

    logger_.log("%:% %() % %\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_), end_market_update.toString());

    snapshot_socket_.send(&end_market_update, sizeof(MDPMarketUpdate));
    snapshot_socket_.sendAndRecv();

    logger_.log("%:% %() % Published snapshot of % orders.\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_), snapshot_size - 1);
}

SnapshotSynthesizer::run()

void SnapshotSynthesizer::run() {
    logger_.log("%:% %() %\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_));

    while(run_) {
        // 遍历 MDP 与 SnapshotSynthesizer 之间的无锁队列
        for (auto market_update = snapshot_md_updates_->getNextToRead(); snapshot_md_updates_->size() && market_update; market_update = snapshot_md_updates->getNextToRead()) {

            logger_.log("%:% %() % Processing %\n", __FILE__, __LINE__, __FUNCTION__, getCurrentTimeStr(&time_str_), market_update->toString().c_str())

            // 让市场更新作用于 SnapshotSynthesizer 维护的订单簿快照
            addToSnapshot(market_update);

            snapshot_md_updates_->updateReadIndex();
        } // 循环结束,订单簿快照得到了一轮更新

        // 若距离上次发送快照的时间间隔已经足够长,则发送快照
        if (getCurrentNanos() - last_snapshot_time_ > 60 * NANOS_TO_SECS) {
            last_snapshot_time = getCurrentNanos();
            publishSnapshot();
        }
    }
}

交易所主程序

我们上一章学习的交易所主程序时还未学习 OrderServer 和 MarketDataPublisher。此处在交易所主程序 exchange_main 当中加入这两个组件。

Exchange::MatchingEngine *matching_engine = nullptr;
Exchange::MarketDataPublisher *market_data_publisher = nullptr;
Exchange::OrderServer *order_server = nullptr;

int main(int, char **) {
    logger = new Common::Logger("exchange_main.log");

    std::signal(SIGINT, signal_handler);

    const int sleep_time = 100 * 1000;

    // 创建三条无锁队列
    Exchange::ClientRequestLFQueue client_requests(ME_MAX_CLIENT_UPDATES);
    Exchange::ClientResponseLFQueue client_responses(ME_MAX_CLIENT_UPDATES);
    Exchange::MEMarketUpdateLFQueue market_updates(ME_MAX_MARKET_UPDATES);

    std::string time_str;

    // 创建并启动撮合引擎
    matching_engine = new Exchange::MatchingEngine(&client_requests, &client_responses, &market_updates);
    matching_engine->start();

    // 配置 MDP 的接口、IP和端口,然后创建并启动 MDP
    const std::string mkt_pub_iface = "lo"; // 网络接口。其实就是走哪个网卡是吧?

    const std::string inc_pub_ip = "123.456.78.2";
    const int inc_inc_port = 20001;

    const std::string snap_pub_ip = "123.456.78.1";
    const int snap_pub_port = 20000;

    market_data_publisher = new Exchange::MarketDataPublisher(&market_updates, mkt_pub_iface, snap_pub_ip, snap_pub_port, inc_pub_ip, inc_pub_port);
    market_data_publisher->start();

    // 配置 OrderServer 的接口、端口
    const std::string order_gw_iface = "lo";  // 这里又叫 gw,有病吧,gw 应该是客户端的吧。TODO:重命名
    const int order_gw_port = 12345;

    order_server = new Exchange::OrderServer(&client_requests, &client_responses, order_gw_iface, order_gw_port);
    order_server->start();

    while (true) {
        usleep(sleep_time * 1000);  // 啥也不用做了,睡觉就好。
    }
}