如何确保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
相关产品推荐
相关产品推荐

