You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

QuickFix C++ 如何在其他函数运行时持续读取market data

问题描述

目标

持续打印来自服务端的market data消息,同时运行另一个包含while循环的函数,监听来自非行情类信息源的输入。核心要求为:无论其他函数是否运行,监听market data的逻辑需始终保持运行状态。

详情说明

下方示例程序可成功向服务端发送market data request消息,后续可按预期持续打印包含bids、asks等交易信息的Market data消息。之后程序会调用另一函数,通过while循环监听其他信息源的输入。
负责检测新market data消息的函数为Application::onMessage(const FIX44::MarketDataSnapshotFullRefresh& mdMessage, const FIX::SessionID& session)。

问题现象

当负责监听其他输入的第二个函数被调用后,终端输出的market data消息流会停止打印,直到收到对应输入后才会恢复,存在明确的blocking阻塞行为。

简化示例代码

#include "Application.h"
#include "quickfix/Session.h"

#include <iostream>

static bool listeningToMarketData = false;
static bool requestMadeAlready = false;

// 登录成功输出
void Application::onLogon( const FIX::SessionID& sessionID ) {
    std::cout << std::endl << "Logon - " << sessionID << std::endl;
}

// 登出输出
void Application::onLogout( const FIX::SessionID& sessionID ) {
    std::cout << std::endl << "Logout - " << sessionID << std::endl;
}

// 客户端发送的管理类消息处理
void Application::toAdmin( FIX::Message& message, const FIX::SessionID& sessionID) {

    // 登录消息设置用户名密码
    if (FIX::MsgType_Logon == message.getHeader().getField(FIX::FIELD::MsgType))
    {
        message.getHeader().setField(FIX::Username("XXXXX"));
        message.getHeader().setField(FIX::Password("XXXXX"));
    }
}

// 服务端推送消息接收处理
void Application::fromApp( const FIX::Message& message, const FIX::SessionID& sessionID )
throw( FIX::FieldNotFound, FIX::IncorrectDataFormat, FIX::IncorrectTagValue, FIX::UnsupportedMessageType ) {
    crack( message, sessionID );
    std::cout << std::endl << "INCOMING: " << message << std::endl;
}

// 客户端发送应用类消息处理
void Application::toApp( FIX::Message& message, const FIX::SessionID& sessionID )
throw( FIX::DoNotSend ) {
    try {
        FIX::PossDupFlag possDupFlag;
        message.getHeader().getField( possDupFlag );
        if ( possDupFlag ) throw FIX::DoNotSend();
    }
    catch ( FIX::FieldNotFound& ) {}

    std::cout << std::endl << "OUTGOING: " << message << std::endl;
}

// 行情消息处理函数
void Application::onMessage(const FIX44::MarketDataSnapshotFullRefresh& mdMessage, const FIX::SessionID& session) {

    static bool isFirstMdMessage = true;
    listeningToMarketData = true;
}

void listenToMarketData() {
        FIX44::MarketDataRequest marketDataRequest(
                              FIX::MDReqID("ABC1"),
                              FIX::SubscriptionRequestType('1'),
                              FIX::MarketDepth(1));

        marketDataRequest.set( FIX::MDUpdateType(0) );
        marketDataRequest.set( FIX::NoMDEntryTypes(2) );
        marketDataRequest.set( FIX::NoRelatedSym(1) );

        FIX44::MarketDataRequest::NoRelatedSym noRelatedSym;
        FIX44::MarketDataRequest::NoMDEntryTypes noMDEntryTypes1;
        FIX44::MarketDataRequest::NoMDEntryTypes noMDEntryTypes2;

        noRelatedSym.set(FIX::SecurityIDSource("8"));
        noRelatedSym.set(FIX::SecurityID("100800"));
        noMDEntryTypes1.set( FIX::MDEntryType('0') ); // 买盘
        noMDEntryTypes2.set( FIX::MDEntryType('1') ); // 卖盘

        marketDataRequest.addGroup( noRelatedSym );
        marketDataRequest.addGroup( noMDEntryTypes1 );
        marketDataRequest.addGroup( noMDEntryTypes2 );

        FIX::Header& mdHeader = marketDataRequest.getHeader();
        mdHeader.setField(FIX::TargetCompID("TARGET"));
        mdHeader.setField(FIX::SenderCompID("SENDER"));

        FIX::Session::sendToTarget(marketDataRequest);
}

void request() {
    requestMadeAlready = true;

    while (true) {
        std::cout << "listening" << std::endl;
        usleep(1000000);
    }
}

// 客户端启动入口
void Application::run() {
    initialiseSystem();
}

// 初始化行情订阅和其他监听逻辑
void Application::initialiseSystem() {

    while (true) {
        try {
            if (!listeningToMarketData) {

                listenToMarketData();
            }

            usleep(5000000);

            if (!requestMadeAlready) {
                request();
            }
        }
        catch ( std::exception & e )
        {
            std::cout << "Message Not Sent: " << e.what() << std::endl;
        }
    }
}

咨询问题

Q1. 如何让程序在并行运行其他输入监听函数的场景下,持续正常打印接收到的market data消息?


问题根因

现有代码是单线程执行模型,request()内部写了永久while循环,一旦调用就会完全占住当前线程的执行权。QuickFIX引擎本身的socket消息接收、回调触发逻辑都运行在这个线程中,被阻塞后自然无法处理服务端推送的行情数据,直到阻塞逻辑释放才会恢复。

解决方案

把阻塞的非行情监听逻辑放到独立线程执行,不要占用QuickFIX的工作线程。

  1. 引入C++标准库的线程支持头文件,不需要依赖第三方组件。
  2. 调整初始化逻辑,不要在主循环里直接调用阻塞的request()函数,第一次触发条件时启动独立线程运行非行情监听,启动后立刻返回主循环,不等待监听线程结束。
  3. 多线程共享的标志位替换为原子类型,避免数据竞争导致的未定义行为。

核心修改代码示例

// 新增头文件引用
#include <atomic>
#include <thread>

// 普通bool替换为原子bool,解决多线程读写竞争问题
static std::atomic<bool> listeningToMarketData{false};
static std::atomic<bool> requestMadeAlready{false};

// 其余原有业务函数保持不变,仅修改initialiseSystem逻辑
void Application::initialiseSystem() {

    while (true) {
        try {
            if (!listeningToMarketData) {
                listenToMarketData();
            }

            usleep(5000000);

            if (!requestMadeAlready) {
                requestMadeAlready = true;
                // 启动独立线程运行非行情监听,不阻塞当前主循环
                std::thread listenerThread(request);
                // 线程分离,自行管理生命周期,不需要主线程等待结束
                listenerThread.detach();
            }
        }
        catch ( std::exception & e )
        {
            std::cout << "Message Not Sent: " << e.what() << std::endl;
        }
    }
}

注意事项

  • 所有QuickFIX的回调函数(包括onMessage、fromApp、toApp等)内部不要写任何阻塞逻辑,否则同样会卡断整个消息处理流程。
  • QuickFIX的sendToTarget接口本身是线程安全的,在子线程里直接调用发送消息不需要额外加锁。
  • 如果多个线程需要同时调用std::cout打印日志,建议加一个全局std::mutex做输出保护,避免不同线程的打印内容交错乱码。

内容的提问来源于stack exchange,提问作者p.luck

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 20:15:43