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

如何确保QuickFIX C++库发送的异步消息均被接收方确认?

问题描述

使用QuickFIX C++库通过sendToTarget方法异步发送消息流时,遇到一个难题:在停止Initiator前,如何确定所有消息已被接收方(Acceptor)接收并确认?

由于sendToTarget是异步调用,且接收方不会回复业务消息,无法直接通过业务回执判断消息确认状态,导致难以确定何时可以安全停止Initiator。

简化代码如下:

// Code snippet to send messages using QuickFIX C++ library

#include <quickfix/Application.h>
#include <quickfix/Session.h>

class MyApplication : public FIX::Application {
public:
    void sendMessages() {
        // Loop to send multiple messages
        for (int i = 0; i < numMessages; ++i) {
            // Construct and send message
            FIX::Message message;
            // Populate message fields
            // ...

            // Send message asynchronously
            FIX::Session::sendToTarget(message);
        }
    }

    // Override other necessary methods of FIX::Application class
};

int main() {
    MyApplication app(...);
    FIX::SocketInitiator initiator(app, ...);
    initiator.start();
    app.sendMessages();
    initiator.stop();
}
解决方案

要解决这个问题,需要利用FIX协议的底层会话确认机制(而非业务消息回执),结合QuickFIX的回调方法跟踪接收方的消息接收状态:

1. 跟踪已发送消息的最大序列号

每条FIX消息都带有MsgSeqNum(消息序列号),发送方会话会自动递增这个值。我们需要记录所有待确认消息的最大序列号,以此作为确认基准。

2. 监听接收方的会话确认信号

接收方会通过**心跳(Heartbeat)或测试请求(TestRequest)**回复NextExpectedMsgSeqNum字段,这个字段代表接收方期望收到的下一条消息的序列号。当NextExpectedMsgSeqNum大于等于我们发送的最大序列号时,说明接收方已经收到了所有已发送的消息。

3. 实现同步等待逻辑

使用条件变量+互斥锁的组合,在发送完所有消息后,阻塞等待直到确认所有消息被接收方收到,再停止Initiator。如果心跳间隔过长,可以主动发送TestRequest触发接收方立即回复,加快确认速度。

修改后的代码示例
#include <quickfix/Application.h>
#include <quickfix/Session.h>
#include <quickfix/MessageCracker.h>
#include <quickfix/Heartbeat.h>
#include <quickfix/TestRequest.h>
#include <mutex>
#include <condition_variable>
#include <atomic>

class MyApplication : public FIX::Application, public FIX::MessageCracker {
private:
    std::atomic<int> m_sentMaxSeqNum{0};
    std::mutex m_mutex;
    std::condition_variable m_cv;
    bool m_allConfirmed{false};
    FIX::SessionID m_targetSession; // 目标会话ID,需根据实际配置初始化

public:
    void sendMessages() {
        // 发送消息前先获取当前会话的初始序列号
        if (auto session = FIX::Session::lookupSession(m_targetSession)) {
            m_sentMaxSeqNum = session->getSenderMsgSeqNum();
        }

        // 循环发送消息
        for (int i = 0; i < numMessages; ++i) {
            FIX::Message message;
            // 填充消息字段(注意设置正确的MsgType等)
            // ...

            FIX::Session::sendToTarget(message, m_targetSession);
            
            // 更新已发送的最大序列号
            if (auto session = FIX::Session::lookupSession(m_targetSession)) {
                m_sentMaxSeqNum = session->getSenderMsgSeqNum() - 1; // send后序列号已递增,所以减1得到当前发送的消息序列号
            }
        }

        // 主动发送TestRequest触发接收方立即回复,避免等待心跳间隔
        sendTestRequest();
    }

    void sendTestRequest() {
        FIX::TestRequest testReq;
        testReq.set(FIX::TestReqID("WAIT_FOR_CONFIRM"));
        FIX::Session::sendToTarget(testReq, m_targetSession);
    }

    // 重写onMessage方法,处理接收方的心跳
    void onMessage(const FIX::Heartbeat& msg, const FIX::SessionID& sessionID) override {
        checkConfirmation(msg.getHeader(), sessionID);
    }

    // 检查接收方的确认状态
    void checkConfirmation(const FIX::Header& header, const FIX::SessionID& sessionID) {
        if (sessionID != m_targetSession) return;

        int nextExpectedSeq = header.getField(FIX::NextExpectedMsgSeqNum());
        if (nextExpectedSeq > m_sentMaxSeqNum) {
            std::lock_guard<std::mutex> lock(m_mutex);
            m_allConfirmed = true;
            m_cv.notify_all();
        }
    }

    // 等待所有消息被确认
    void waitForAllConfirmations() {
        std::unique_lock<std::mutex> lock(m_mutex);
        m_cv.wait(lock, [this](){ return m_allConfirmed; });
    }

    // 重写其他必要的Application方法
    void onCreate(const FIX::SessionID&) override {}
    void onLogon(const FIX::SessionID&) override {}
    void onLogout(const FIX::SessionID&) override {}
    void toAdmin(FIX::Message&, const FIX::SessionID&) override {}
    void toApp(FIX::Message&, const FIX::SessionID&) throw(FIX::DoNotSend) override {}
    void fromAdmin(const FIX::Message&, const FIX::SessionID&) throw(FIX::FieldNotFound, FIX::IncorrectDataFormat, FIX::IncorrectTagValue, FIX::RejectLogon) override {}
    void fromApp(const FIX::Message&, const FIX::SessionID&) throw(FIX::FieldNotFound, FIX::IncorrectDataFormat, FIX::IncorrectTagValue, FIX::UnsupportedMessageType) override {}
};

int main() {
    MyApplication app;
    // 初始化SessionID、SocketInitiator配置等
    // ...

    FIX::SocketInitiator initiator(app, ...);
    initiator.start();
    
    app.sendMessages();
    app.waitForAllConfirmations(); // 等待所有消息被接收方确认
    
    initiator.stop();
    return 0;
}
关键说明
  • m_sentMaxSeqNum记录已发送消息的最大序列号,每次发送后更新。
  • 通过onMessage监听接收方的心跳消息,解析NextExpectedMsgSeqNum字段判断是否所有消息已被接收。
  • 主动发送TestRequest可以避免等待默认心跳间隔,加快确认速度。
  • 使用条件变量m_cv实现同步等待,确保只有在所有消息确认后才停止Initiator。

内容的提问来源于stack exchange,提问作者Peter

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 00:43:15