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

交易系统:撮合引擎运行态及通信机制 交易系统:撮合引擎运行态及通信机制 交易系统:撮合引擎运行态及通信机制

交易所主程序 exchange_main

exchange_main 负责创建交易所组件和它们之间的队列,启动 Matching Engine,并在收到终止信号后协调清理退出。

全局变量和 main() 函数

// 此处先声明全局变量,不初始化
Common::Logger* logger = nullptr;
Exchange::MatchingEngine* matching_engine = nullptr;

int main(int, char **) {
    // 日志记录器
    logger = new Common::Logger("exchange_main.log");

    // 注册一个信号处理函数,受到 `SIGINT` 信号的时候,调用 `signal_handler` 处理函数
    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::MEMatketUpdateLFQueue market_updates(ME_MAX_MARKET_UPDATES);

    // 初始化 matching_engine
    std::string time_str;
    logger->log("%:% &() % Starting Matching Engine...\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str));
    matching_engine = new Exchange::MatchingEngine(&client_requests, &client_responses, &market_updates);

    // 启动撮合引擎并启动无限循环
    matching_engine->start();
    while (true) {
        logger->log("%:% %() % Sleeping for a few milliseconds..\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str));
        usleep(sleep_time * 1000);
    }

}

这里其实可以总结一下在本书当中出现的 3 个 namespace:

  1. Common:包含通用基础设施,不属于交易所或策略业务
  • 比如 LFQueueLogger、线程工具、时间工具、宏、基础类型和公共常量
  1. Exchange:交易所侧组件
  • 例如 MatchingEngineMEOrderBookOrderServerMarketDataPublisher
  • 以及交易所内部的请求、响应和市场更新类型
  1. Trading:客户端侧交易系统
  • 例如 MarketDataConsumer、客户端订单网关、TradeEngine、订单管理、风险管理、持仓和策略

这样至少有 2 个好处:

  1. OrderId、Logger、MatchingEngine 等名字很通用。用 namespace 后可以区分
  2. Exchange 和 Trading 可以依赖 Common,但是不应该存在反向依赖

main() 当中的 Exchange::TradingEngine 类的成员和方法在下面细说。

signal_handler 负责处理 SIGINT 信号

signal_handler() 函数用来处理 SIGINT 信号。这个信号一般由 ctrl + C 触发,用来表示停止交易所程序。

这里直接使用 sleep_for,显然是粗糙的等待机制,不是生产级做法。正确的顺序应该是 “收到信号 -> 设置停止标志 -> stop() 各个组件 -> join() 工作线程 -> 释放对象 -> 正常退出”

这里析构 loggermatching_engine 的做法是:先调用 delete 释放 loggermatching_engine 所指向的对象的内存,然后将这两个指针置为 nullptr。

void signal_handler(int) {
    using namespace std::literals::chrono_literals;
    std::this_thread::sleep_for(10s);

    delete logger;
    logger = nullptr; 
    delete matching_engine;
    matching_engine = nullptr;

    std::this_thread::sleep_for(10s);
    exit(EXIT_SUCCESS);
}

撮合引擎的成员和方法

撮合引擎的成员对象

class MatchingEngine final {
private:
    OrderBookHashMap ticker_order_book_;
    ClientRequestLFQueue *incoming_requests_ = nullptr;
    ClientResponseLFQueue *outgoing_ogw_responses_ = nullptr;
    MEMarketUpdateLFQueue *outgoing_md_updates_ = nullptr;
    volatile bool run_ = false;
    std::string time_str_;
    Logger logger_;
}

一个 Matching Engine 线程可以管理多个 ticker 的订单簿。书中的方案是一个 Matching Engine 线程管理所有 ticker 的订单簿。ticker_order_book 是一张查找表,可以通过 ticker_id_ 找到对应 ticker 的 MEOrderBook。然而在实际生产中,如果使用多个 Matching Engine 线程分别管理一部分 ticker 的订单簿,则可能引入 MPSC 或者 SPMC 问题,超出了此处的讨论范围。

Order Server 通过 epoll 与若干个 Client 维持 TCP 连接。Client 发过来的订单请求进入 Order Server 的 FIFO Sequencer,然后按照时间顺序进入 ClientRequestLFQueue。然后 Matching Engine 线程根据 MEClientRequestticker_id_ 字段,更新对应 ticker 的订单簿。随后再分别向 ClientResponseLFQueueMEMarketUpdateLFQueue 添加订单更新信息。

run_ 布尔对象用来标识线程循环是否继续。使用了 volatile 关键词是因为这个变量可能被多个线程访问。

time_str_ 用于暂存格式化后的当前时间字符串,用于日志记录。

构造函数 MatchingEngine()

主要工作是建立 ticker_order_book 查找表当中每个 ticker 的订单簿

MatchingEngine::MatchingEngine(
    ClientRequestLFQueue* client_requests, 
    ClientResponseLFQueue* client_responses,
    MEMarketUpdateLFQueue* market_updates
    ) : incoming_requests_(client_requests),
    outgoing_ogw_responses_(client_requests),
    outgoing_md_updates_(market_updates),
    logger_("exchange_matching_engine.log") {
        for(size_t i = 0; i < ticker_order_book.size(), ++i) {
            ticker_order_book[i] = new MEOrderBook(i, &logger_, this);
        }
    }

记得把默认构造函数、拷贝/移动构造函数和赋值运算符都 =delete。

析构函数 ~MatchingEngine()

这里要注意,MatchingEngine 只拥有订单簿,不拥有三条 LFQueue,因此对待它们的析构方式不同。

MatchingEngine::~MatchingEngine() {
    run_ = false;
    using namespace std::literals::chrono_literals;
    std::this_thread::sleep_for(1s);  // 粗糙实现

    incoming_requests = nullptr;
    outgoing_ogw_responses = nullptr;
    outgoing_md_responses = nullptr;

    for (auto order_book: ticker_order_book) {
        delete order_book;
        order_book = nullptr;
    }
}

这里因为 MatchingEngine 类并不拥有那三条 LFQueue,只拥有指向它们的指针,因此在析构函数这里针对这些指针只要置为 nullptr 即可。但是 MatchingEngine 持有订单簿查找表,因此要针对查找表当中的每一个 MEOrderBook 调用 delete 释放内存并且将指针置为 nullptr。

start() 创建并启动撮合引擎线程

auto MatchingEngine::start() {
    run_ = true;
    if (Common::createAndStartThread(-1, "Exchange/MatchingEngine", [this]() { run(); }) != nullptr) {
        std::cerr << "Failed to start MarchingEngine thread."; // 实际上这里用 cerr 是不对的
    }
}

stop() 将 run_ 置为 false

void MatchingEngine::stop() {
    run_ = false;
}

线程会运行的 run() 函数

为什么要区别 start() 和 run()?start() 管理线程生命周期,创建并启动线程;run() 定义工作线程实际执行的循环。如果把两者合并,那么调用 start() 的线程会被阻塞,无法继续创建或启动 OrderServerMarketDataPublisher 等组件。

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

    while(run_) {
        // 从无锁队列当中获取指向 MEClientRequest 的指针
        const auto me_client_request = incoming_requests_->getNextToRead();

        if (LIKELY(me_client_request)) {
            // 打印日志
            logger_.log("%:% %() % Processing %\n", __FILE__, __LINE__,  __FUNCTION__, Common::getCurrentTimeStr(&time_str_),
            me_client_request->toString());

            // 由专门的处理函数处理拿到的 MEClientRequest
            processClientRequest(me_client_request);
            incoming_requests_->updateReadIndex();
        }
    }
}

这里之所以要加上 noexcept 是为了表明,如果有一个异常在这个函数内没有被 catch,并且尝试离开这个函数(到它的调用者处继续追寻异常源头),那么程序会直接中止。把 noexcept 用在热路径代码上,实质上是想要说明撮合热路径不使用异常恢复机制,出现异常就认为是不可恢复的程序错误,直接停机。

processClientRequest()

这里用了一个 switch statement 根据客户请求类型调用 order_book 的相应方法完成工作。

auto processClientRequest(const MEClientRequest* client_request) noexcept {
    auto order_book = ticker_order_book[client_request->ticker_id_];

    switch(client_request_->type_) {
        case ClientRequestType::NEW: {
            order_book.add(client_request->client_id_, client_request->order_id_, client_request->ticker_id_, client_request->side_, client_request->price_, client_request->qty_);
        }
        break;
        case ClientRequestType::CANCEL: {
            order_book.cancel(client_request->client_id_, client_request->client_order_id_, client_request->ticker_id_);
        }
        break;
        default: {
            FATAL("Received invalid client request type:" + clientRequestTypeToString(client_request->type_));
        }
        break;
    }
}

sendClientResponse()

auto MatchingMachine::sendClientResponse(const MEClientResponse* client_responses) noexcept {
    logger_.log("%:% %() % Sending %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), client_response->toString());

    auto next_write = outgoing_ogw_responses_->getNextToWriteTo();
    next_write* = client_responses;
    outgoing_ogw_response_->updateWriteIndex();
}

这里 MEClientResponse 对象内没有动态资源,都是标量字段,所以哪怕把 next_write* = client_responses 换成 next_write* = std::move(client_responses) 也不会增加收益。即复制与移动的开销是一致的。

sendMarketUpdate()

同理,向 Market Data Publisher 发送数据的函数:

auto sendMarketUpdate(const MEMarketUpdate *market_update) noexcept {
    ogger_.log("%:% %() % Sending %\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), market_update->toString());

    auto next_write = outgoing_md_update_->getNextToWriteTo();
    *next_write = *market_update;
    outgoing_md_updates->updateWriteIndex();
}