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

交易系统:UDP McastSocket 类 交易系统:UDP McastSocket 类 交易系统:UDP McastSocket 类

成员对象

struct McastSocket {
    init socket_fd_ = -1;

    // 应用层缓冲区
    std::vector<char> outbound_data_;
    size_t next_send_valid_index_ = 0;
    std::vector<char> inbound_data_;
    size_t next_rcv_valid_index_ = 0;

    // 只有 MDC 需要用到这个回调函数吗?
    std::function<void(McastSocket* s)> recv_callback_ = nullptr;

    std::string time_str_;
    Logger &logger; 
}

在 McstSocket 当中的 logger 对象是引用类型的,而在一些其它对象当中的 logger 是非引用类型。McastSocket 借用所属组件的 logger,把日志写入该组件的日志文件。外部 logger 必须活得比 McastSocket 对象更久。

McastSocket::init()

负责创建并配置 multicast socket。我发现 UDP multicast 的 socket 的配置比 TCP socket 的配置简单很多,是不是?

auto McastSocket::init(const std::string &ip, const std::string &iface, int port, bool is_listening) {
    const Common::SocketCfg socket_cfg{ip, iface, port, true, is_listening, false};
    socket_fd_ = Common::createSocket(logger_, socket_cfg);
    return socket_fd_;
}

join() 和 leave()

发送者向组播 IP 发送,加入该组的接收者接收数据。

McastSocket::join() 用于接收端加入组播组。而发送方是无需加入就可发送的。

由于 McastSocket::leave() 的实现实质上是关闭 socket,因此发送方和接收方都可以调用

bool McastSocket::join(const std::string& ip) {
    return Common::join(socket_fd_, ip);
}

auto McastSocket::leave(const std::string&, int) {
    // 这个方法的参数没写错,是合法的“省略参数名”语法,表示函数接收这些参数,但实现不使用它们
    // 这里的 leave() 是直接关闭 socket,发送端也可调用。leave 的方法名其实不够准确
    // TODO:重构代码,把 leave() 写成单纯从组播组中退出;为关闭 socket 写一个新方法
    close(socket_fd_);
    socket_fd_ = -1;
}

McastSocket::sendAndRecv()

负责把 socket 内核缓冲区当中当中的字节流复制到应用层接收缓冲区,并且把应用层发送缓冲区的字节流复制到 socket 内核缓冲区。

auto McastSocket::sendAndRecv() noexcept {
    // 把内核缓冲区的信息拷贝到应用层缓冲区,使用 recv() 系统调用
    const ssize_t n_rcv = recv(socket_fd_, inbound_data_.data() + next_rcv_valid_index_, McastBufferSize - next_rcv_valid_index_, MSG_DONTWAIT);

    if (n_rcv > 0) {
        next_rcv_valid_index_ += n_rcv;

        logger_.log("%:% %() % read socket:% len:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket_fd_,  next_rcv_valid_index_);

        // 让回调函数来处理应用层缓冲区的数据
        recv_callback_(this);
    }

    // 把应用层缓冲区的字节移动到内核缓冲区
    if (next_send_valid_index_ > 0) {
        ssize_t n = ::send(socket_fd_, outbound_data_.data(), next_send_valid_index_, MSG_DONTWAIT | MSG_NOSIGNAL);

        logger_.log("%:% %() % send socket:% len:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket_fd_, n);

        next_send_valid_index_ = 0;

        return (n_rcv > 0);
    }
}

这里使用的 recv() 也是系统调用。recvmsg()recv() 的区别在于:

  • recv() 直接传入缓冲区地址、长度和标志
  • recvmsg() 通过 msghdr 描述多个缓冲区,因此可以接收实践戳等辅助信息

在 OrderServer 当中,FIFOSequencer 依赖每个订单到达内核的时间,因此需要额外的缓冲区来存储时间戳这一辅助信息,因此 TCPSocket 使用 recvmsg() 系统调用。但是在 MarketDataPublisher 这边,推送市场更新的 McastSocket 不需要记录辅助信息,因此使用 recv() 系统调用即可。

McastSocket::send()

auto McastSocket::send(const void* data, size_t len) noexcept {
    // const void* 表示指向未知类型数据的指针,不能通过它修改数据
    // 因此 send() 可以接收不同结构体或字节数组的地址,配合 len 确定复制多少字节
    
    // memcpy 把待发送数据追加到应用层缓冲区
    // 参数:memcpy(目标地址, 源地址, 字节数); memcpy 不是系统调用,完全运行在用户态
    memcpy(outbound_data_.data() + next_send_valid_index_, data, len);
    next_send_valid_index_ += len;
    ASSERT(next_send_valid_index_ < McastBufferSize, "Mcast socket buffer filled up and sendAndRecv() not called.");
}