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

交易系统: 无锁队列 交易系统: 无锁队列 交易系统: 无锁队列

在交易系统当中,同机器上的不同线程之间通过放在共享内存的无锁队列来互相通信,且通信的情况限制在单方向、一对一的通信,不存在多个生产者、多个消费者的情况。可以使用 Ring Buffer 来高效且优雅实现这种 SPSC (Single Producer, Single Consumer) 通信。

Ring Buffer

先看成员变量

class LFQueue final {    
private:
    std::vector<T> store_;
    std::atomic<size_t> next_write_index_{0};
    std::atomic<size_t> next_read_index_{0};
    std::atomic<size_t> num_elements_{0};
}

有两个 index:

  • next_write_index_ 表示下一次写入的位置
  • next_read_index_ 表示下一次读取的位置 由于目前我们只讨论 SPSC 的情况,因此每个 index 都只由一个线程推进(producer 或者 consumer),因此不需要多个生产者中间竞争写入位置,不需要多个消费者之间竞争读位置,实现可以显著简单,用两个 index 就搞定。
  • num_elements 表示队列当中尚未被消费者读取的数据

在正常运转当中,next_write_index_ 一定要快过 next_read_index_ ,且最多快一整圈。两者相等时意味着队列为空;num_elements_ == store_.size() 的时候意味着队列已满

alt text

无锁队列初始化

template<typename T>
class LFQueue final {
    LFQueue(std::size_t num_elems):
        store_(num_elems, T()) {}
};

操作方法主要是 getter 和 setter

  • getNextToRead()
  • updateReadIndex()
  • getNextToWrite()
  • updateWriteIndex()

getNextToRead()

// 注意,这里返回的不是 next_write_index_ 这个 index,而是指向 store_ 当中元素的指针
const T* getNextToRead() const noexcept {
    if (next_read_index_ == next_write_index_) {  // 判断队列是否为空
        return nullptr;
    }
    return &store_[next_read_index_];
}

updateReadIndex()

每一个索引都是这样算的:(index + 1) % store_.size()

void updateReadIndex() noexcept {
    next_read_index_ = (next_read_index_ + 1) % store_.size();
    num_elements--;
}

调用方(consumer)需要遵守以下流程:

if (const auto* slot = lfq.getNextToRead()) {  // 不是 nullptr,证明队列非空
    process(*slot);
    lfq.updateReadIndex();
}

getNextToWriteTo()

auto getNextToWriteTo() noexcept {
    if (num_elements_.load() == store_.size()) { // 队列已满
        return nullptr;
    }
    return &store_[next_write_index_];
}

updateWriteIndex()

void updateWriteIndex() {
    next_write_index_ = (next_write_index_ + 1) % store_.size();
    num_elements++;
}

调用方遵守以下流程:

if (auto* slot = lfq.getNextToWriteTo()) {
    *slot = message;
    lfq.updateWriteIndex();
}

注意:这里在生产者 updateWriteIndex() 和消费者 updateReadIndex() 之后,没有一个机制唤醒对方,所以二者其实是异步的。另外由于这里的 LFQueue 使用了模板编程,队列内部可以存储任意自定义对象(也就是自定义的消息类型),我们希望达到的效果就是每次生产者只 push 一条消息;每次消费者也只读取一条消息。