文章目录
交易系统:撮合引擎运行态及通信机制 交易系统:撮合引擎运行态及通信机制 交易系统:撮合引擎运行态及通信机制
交易所主程序 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:
- Common:包含通用基础设施,不属于交易所或策略业务
- 比如
LFQueue、Logger、线程工具、时间工具、宏、基础类型和公共常量
- Exchange:交易所侧组件
- 例如
MatchingEngine、MEOrderBook、OrderServer、MarketDataPublisher - 以及交易所内部的请求、响应和市场更新类型
- Trading:客户端侧交易系统
- 例如 MarketDataConsumer、客户端订单网关、TradeEngine、订单管理、风险管理、持仓和策略
这样至少有 2 个好处:
- OrderId、Logger、MatchingEngine 等名字很通用。用 namespace 后可以区分
- Exchange 和 Trading 可以依赖 Common,但是不应该存在反向依赖
main() 当中的 Exchange::TradingEngine 类的成员和方法在下面细说。
signal_handler 负责处理 SIGINT 信号
signal_handler() 函数用来处理 SIGINT 信号。这个信号一般由 ctrl + C 触发,用来表示停止交易所程序。
这里直接使用 sleep_for,显然是粗糙的等待机制,不是生产级做法。正确的顺序应该是 “收到信号 -> 设置停止标志 -> stop() 各个组件 -> join() 工作线程 -> 释放对象 -> 正常退出”
这里析构 logger 和 matching_engine 的做法是:先调用 delete 释放 logger 和 matching_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 线程根据 MEClientRequest 的 ticker_id_ 字段,更新对应 ticker 的订单簿。随后再分别向 ClientResponseLFQueue 和 MEMarketUpdateLFQueue 添加订单更新信息。
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() 的线程会被阻塞,无法继续创建或启动 OrderServer 、 MarketDataPublisher 等组件。
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();
}